Hyperlinkv0.9.0-beta.0

WorkPool

WorkPool.priorityconstsrc/WorkPool.ts:1785
<Self>(): {
  <F extends Schema.Struct.Fields, Success extends Schema.Top, HSelf>(
    key: string,
    config: PriorityTagConfig<F, Success> & {
      readonly node: NodeKey<HSelf>
    }
  ): NodeBoundTag<Self, PriorityInstanceSpec<F>, HSelf>
  <
    F extends Schema.Struct.Fields,
    Success extends Schema.Top = Schema.Void
  >(
    key: string,
    config: PriorityTagConfig<F, Success>
  ): HyperlinkTag<Self, PriorityInstanceSpec<F>>
}

Define an N-level managed queue as a named service Tag (also exported as priority): class Jobs extends WorkPool.priority<Jobs>()("@app/Jobs", { … }) {}. The priority (N-level lane) peer of Tag — same WorkPool, with laneCount / namedLanes priority lanes and add(item, lane?). class Jobs extends WorkPool.priority<Jobs>()("@app/Jobs", { payload, laneCount: 2 }) {}. The class is the Tag — yield* Jobs resolves the handle, layer provides it and serve exposes it over RPC (both dispatch to the leveled engine for a priority tag). payload is the item schema; optional success / error add the worker wire schemas.

constructorsTagprioritylayerserve
Source src/WorkPool.ts:178553 lines
export const priority = <Self>() => {
  function build<
    F extends Schema.Struct.Fields,
    Success extends Schema.Top,
    HSelf,
  >(
    key: string,
    config: PriorityTagConfig<F, Success> & { readonly node: NodeKey<HSelf> },
  ): NodeBoundTag<Self, PriorityInstanceSpec<F>, HSelf>;
  function build<
    F extends Schema.Struct.Fields,
    Success extends Schema.Top = typeof Schema.Void,
  >(
    key: string,
    config: PriorityTagConfig<F, Success>,
  ): HyperlinkTag<Self, PriorityInstanceSpec<F>>;
  function build<F extends Schema.Struct.Fields, Success extends Schema.Top>(
    key: string,
    config: PriorityTagConfig<F, Success>,
  ): HyperlinkTag<Self, PriorityInstanceSpec<F>> {
    const laneConfig: PriorityTagLaneConfig = {
      laneCount: config.laneCount,
      namedLanes: config.namedLanes ?? {},
    };
    const wire = { success: config.success, error: config.error };
    const spec = assertPriorityInstanceSpec<F>(
      prioritySpec(config.payload, laneConfig, wire),
      prioritySpec(config.payload, laneConfig),
      wire,
    );
    const base =
      config.node === undefined
        ? Hyperlink.Tag<Self>()(key, spec, { description: config.description, kind: priorityKind })
        : Hyperlink.Tag<Self>()(key, spec, {
            description: config.description,
            kind: priorityKind,
            node: config.node,
          });
    const ready = Hyperlink.withReadiness(base, (svc) =>
      Effect.map(svc.status.get, (status) => ({
        ready: status.phase === "running",
        ...(status.phase === "running"
          ? {}
          : { detail: `phase: ${status.phase}` }),
      })),
    );
    return stampQueueWireSchemas(ready, {
      success: config.success,
      error: config.error,
    });
  }
  return build;
};