Hyperlinkv0.8.0-beta.28

Channel

Channel.toPubSubconsteffect/Channel.ts:8331
(
  options:
    | {
        readonly capacity: "unbounded"
        readonly replay?: number | undefined
        readonly shutdownOnEnd?: boolean | undefined
      }
    | {
        readonly capacity: number
        readonly strategy?: "dropping" | "sliding" | "suspend" | undefined
        readonly replay?: number | undefined
        readonly shutdownOnEnd?: boolean | undefined
      }
): <OutElem, OutErr, OutDone, Env>(
  self: Channel<OutElem, OutErr, OutDone, unknown, unknown, unknown, Env>
) => Effect.Effect<PubSub.PubSub<OutElem>, never, Env | Scope.Scope>
<OutElem, OutErr, OutDone, Env>(
  self: Channel<OutElem, OutErr, OutDone, unknown, unknown, unknown, Env>,
  options:
    | {
        readonly capacity: "unbounded"
        readonly replay?: number | undefined
        readonly shutdownOnEnd?: boolean | undefined
      }
    | {
        readonly capacity: number
        readonly strategy?: "dropping" | "sliding" | "suspend" | undefined
        readonly replay?: number | undefined
        readonly shutdownOnEnd?: boolean | undefined
      }
): Effect.Effect<PubSub.PubSub<OutElem>, never, Env | Scope.Scope>

Converts a channel to a PubSub for concurrent consumption.

Details

shutdownOnEnd indicates whether the PubSub should be shut down when the channel ends. By default this is true.

destructors
Source effect/Channel.ts:833150 lines
export const toPubSub: {
  (
    options: {
      readonly capacity: "unbounded"
      readonly replay?: number | undefined
      readonly shutdownOnEnd?: boolean | undefined
    } | {
      readonly capacity: number
      readonly strategy?: "dropping" | "sliding" | "suspend" | undefined
      readonly replay?: number | undefined
      readonly shutdownOnEnd?: boolean | undefined
    }
  ): <OutElem, OutErr, OutDone, Env>(
    self: Channel<OutElem, OutErr, OutDone, unknown, unknown, unknown, Env>
  ) => Effect.Effect<PubSub.PubSub<OutElem>, never, Env | Scope.Scope>
  <OutElem, OutErr, OutDone, Env>(
    self: Channel<OutElem, OutErr, OutDone, unknown, unknown, unknown, Env>,
    options: {
      readonly capacity: "unbounded"
      readonly replay?: number | undefined
      readonly shutdownOnEnd?: boolean | undefined
    } | {
      readonly capacity: number
      readonly strategy?: "dropping" | "sliding" | "suspend" | undefined
      readonly replay?: number | undefined
      readonly shutdownOnEnd?: boolean | undefined
    }
  ): Effect.Effect<PubSub.PubSub<OutElem>, never, Env | Scope.Scope>
} = dual(
  2,
  Effect.fnUntraced(function*<OutElem, OutErr, OutDone, Env>(
    self: Channel<OutElem, OutErr, OutDone, unknown, unknown, unknown, Env>,
    options: {
      readonly capacity: "unbounded"
      readonly replay?: number | undefined
      readonly shutdownOnEnd?: boolean | undefined
    } | {
      readonly capacity: number
      readonly strategy?: "dropping" | "sliding" | "suspend" | undefined
      readonly replay?: number | undefined
      readonly shutdownOnEnd?: boolean | undefined
    }
  ) {
    const pubsub = yield* makePubSub<OutElem>(options)
    yield* Effect.forkScoped(runIntoPubSub(self, pubsub, {
      shutdownOnEnd: options.shutdownOnEnd !== false
    }))
    return pubsub
  })
)