Hyperlinkv0.8.0-beta.28

Hyperlink

Hyperlink.runForEachTagconstsrc/Hyperlink.ts:5801
<A extends TaggedEvent, Cases extends TagHandlers<A, unknown, unknown>>(
  handlers: Cases
): (
  self: Stream.Stream<A>
) => Effect.Effect<void, HandlersError<Cases>, HandlersContext<Cases>>
<A extends TaggedEvent, const K extends A["_tag"], E, R>(
  tag: K,
  f: (
    event: Extract<A, { readonly _tag: K }>
  ) => Effect.Effect<void, E, R>
): (self: Stream.Stream<A>) => Effect.Effect<void, E, R>
<A extends TaggedEvent, const K extends A["_tag"], E, R>(
  self: Stream.Stream<A>,
  tag: K,
  f: (
    event: Extract<A, { readonly _tag: K }>
  ) => Effect.Effect<void, E, R>
): Effect.Effect<void, E, R>
<A extends TaggedEvent, Cases extends TagHandlers<A, unknown, unknown>>(
  self: Stream.Stream<A>,
  handlers: Cases
): Effect.Effect<void, HandlersError<Cases>, HandlersContext<Cases>>

Consume a tagged-event Stream, dispatching each element to a handler by its _tag — the stream-native replacement for lifecycle callbacks (one off-fiber consumer, not a fiber per item). Pass a single tag + handler or a handler map; data-first or pipeable. Built on Match, so handlers are fully typed with no casts; unhandled tags are ignored.

yield* jobs.events.pipe(Hyperlink.runForEachTag({
  Failed:  ({ entry, cause }) => Effect.logError(`failed ${entry.entryId}`, cause),
  Drained: ({ completed })    => Effect.log(`drained @ ${completed}`),
}))
yield* Hyperlink.runForEachTag(jobs.events, "Failed", (e) =>
  Effect.failCause(e.cause).pipe(Effect.catchTags({ Timeout: …, Rejected: … })),
)
reactivityStream
Source src/Hyperlink.ts:580149 lines
export const runForEachTag: {
  // ── data-last (pipeable) ──
  // `Cases` is inferred from the literal; E/R are EXTRACTED from the handlers via `infer`
  // (not inferrable params), so `A` can unify at the pipe site without dragging R to `unknown`.
  <
    A extends TaggedEvent,
    Cases extends TagHandlers<A, unknown, unknown>,
  >(
    handlers: Cases,
  ): (
    self: Stream.Stream<A>,
  ) => Effect.Effect<void, HandlersError<Cases>, HandlersContext<Cases>>;
  <A extends TaggedEvent, const K extends A["_tag"], E, R>(
    tag: K,
    f: (event: Extract<A, { readonly _tag: K }>) => Effect.Effect<void, E, R>,
  ): (self: Stream.Stream<A>) => Effect.Effect<void, E, R>;
  // ── data-first ──
  <A extends TaggedEvent, const K extends A["_tag"], E, R>(
    self: Stream.Stream<A>,
    tag: K,
    f: (event: Extract<A, { readonly _tag: K }>) => Effect.Effect<void, E, R>,
  ): Effect.Effect<void, E, R>;
  <A extends TaggedEvent, Cases extends TagHandlers<A, unknown, unknown>>(
    self: Stream.Stream<A>,
    handlers: Cases,
  ): Effect.Effect<void, HandlersError<Cases>, HandlersContext<Cases>>;
} = Fn.dual(
  (args) => Stream.isStream(args[0]),
  // Impl is typed over the concrete `TaggedEvent` base so `Match` (which needs a concrete
  // union, not a generic) type-checks with no casts; the overload signatures above carry the
  // precise per-tag types to callers.
  <E, R>(
    self: Stream.Stream<TaggedEvent>,
    tagOrHandlers: string | TagHandlers<TaggedEvent, E, R>,
    f?: (event: TaggedEvent) => Effect.Effect<void, E, R>,
  ): Effect.Effect<void, E, R> => {
    const matcher =
      typeof tagOrHandlers === "string"
        ? Match.type<TaggedEvent>().pipe(
            Match.tag(tagOrHandlers, f ?? (() => Effect.void)),
            Match.orElse(() => Effect.void),
          )
        : Match.type<TaggedEvent>().pipe(
            Match.tags(tagOrHandlers),
            Match.orElse(() => Effect.void),
          );
    return Stream.runForEach(self, matcher);
  },
);
Referenced by 1 symbols