Hyperlinkv0.8.0-beta.28

CustomQueueHyperlink

CustomQueueHyperlink.Tagconstsrc/CustomQueueHyperlink.ts:365
<Self>(): {
  <F extends Schema.Struct.Fields, Success extends Schema.Top, HSelf>(
    key: string,
    config: CustomQueueTagConfig<F, Success> & {
      readonly node: NodeKey<HSelf>
    }
  ): NodeBoundTag<Self, CustomQueueInstanceSpec<F>, HSelf>
  <
    F extends Schema.Struct.Fields,
    Success extends Schema.Top = Schema.Void
  >(
    key: string,
    config: CustomQueueTagConfig<F, Success>
  ): HyperlinkTag<Self, CustomQueueInstanceSpec<F>>
}

Define an N-level managed queue as a named service Tag (also exported as customQueueTag): class Jobs extends CustomQueueHyperlink.Tag<Jobs>()("@app/Jobs", { … }) {}. The class is the Tag — yield* Jobs resolves the queue handle, layer provides it and serve exposes it over RPC. payload is the item schema; levelCount / namedLevels declare the priority levels; optional success / error add the worker wire schemas.

export const customQueueTag = <Self>() => {
  function build<
    F extends Schema.Struct.Fields,
    Success extends Schema.Top,
    HSelf,
  >(
    key: string,
    config: CustomQueueTagConfig<F, Success> & { readonly node: NodeKey<HSelf> },
  ): NodeBoundTag<Self, CustomQueueInstanceSpec<F>, HSelf>;
  function build<
    F extends Schema.Struct.Fields,
    Success extends Schema.Top = typeof Schema.Void,
  >(
    key: string,
    config: CustomQueueTagConfig<F, Success>,
  ): HyperlinkTag<Self, CustomQueueInstanceSpec<F>>;
  function build<F extends Schema.Struct.Fields, Success extends Schema.Top>(
    key: string,
    config: CustomQueueTagConfig<F, Success>,
  ): HyperlinkTag<Self, CustomQueueInstanceSpec<F>> {
    const levelConfig: CustomQueueTagLevelConfig = {
      levelCount: config.levelCount,
      namedLevels: config.namedLevels ?? {},
    };
    const wire = { success: config.success, error: config.error };
    const spec = assertCustomQueueInstanceSpec<F>(
      customQueueSpec(config.payload, levelConfig, wire),
      customQueueSpec(config.payload, levelConfig),
      wire,
    );
    const base =
      config.node === undefined
        ? Hyperlink.Tag<Self>()(key, spec, { description: config.description, kind })
        : Hyperlink.Tag<Self>()(key, spec, {
            description: config.description,
            kind,
            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;
};