Hyperlinkv0.8.0-beta.28

Channel

Channel.runIntoPubSubconsteffect/Channel.ts:8393
<OutElem>(
  pubsub: PubSub.PubSub<OutElem>,
  options?: { readonly shutdownOnEnd?: boolean | undefined } | undefined
): <OutErr, OutDone, Env>(
  self: Channel<OutElem, OutErr, OutDone, unknown, unknown, unknown, Env>
) => Effect.Effect<void, never, Env>
<OutElem, OutErr, OutDone, Env>(
  self: Channel<OutElem, OutErr, OutDone, unknown, unknown, unknown, Env>,
  pubsub: PubSub.PubSub<OutElem>,
  options?: { readonly shutdownOnEnd?: boolean | undefined } | undefined
): Effect.Effect<void, never, Env>

Runs a channel and publishes each output element to a PubSub.

Details

The channel's output values are published as individual PubSub messages. Use options.shutdownOnEnd to shut down the PubSub when channel execution ends.

destructors
Source effect/Channel.ts:839329 lines
export const runIntoPubSub: {
  <OutElem>(
    pubsub: PubSub.PubSub<OutElem>,
    options?: {
      readonly shutdownOnEnd?: boolean | undefined
    } | undefined
  ): <OutErr, OutDone, Env>(
    self: Channel<OutElem, OutErr, OutDone, unknown, unknown, unknown, Env>
  ) => Effect.Effect<void, never, Env>
  <OutElem, OutErr, OutDone, Env>(
    self: Channel<OutElem, OutErr, OutDone, unknown, unknown, unknown, Env>,
    pubsub: PubSub.PubSub<OutElem>,
    options?: {
      readonly shutdownOnEnd?: boolean | undefined
    } | undefined
  ): Effect.Effect<void, never, Env>
} = dual(
  (args) => isChannel(args[0]),
  <OutElem, OutErr, OutDone, Env>(
    self: Channel<OutElem, OutErr, OutDone, unknown, unknown, unknown, Env>,
    pubsub: PubSub.PubSub<OutElem>,
    options?: {
      readonly shutdownOnEnd?: boolean | undefined
    } | undefined
  ) =>
    runForEach(self, (value) => PubSub.publish(pubsub, value)).pipe(
      options?.shutdownOnEnd === true ? Effect.ensuring(PubSub.shutdown(pubsub)) : identity_
    )
)
Referenced by 1 symbols