MCPcopy Create free account
hub / github.com/Effect-TS/effect / fromQueue

Function fromQueue

packages/effect/src/Sink.ts:514–526  ·  view source on GitHub ↗
(
  queue: Queue.Queue<A, Cause.Done>
)

Source from the content-addressed store, hash-verified

512 * @since 2.0.0
513 */
514export const fromQueue = <A>(
515 queue: Queue.Queue<A, Cause.Done>
516): Sink<void, A> =>
517 fromTransform((upstream) =>
518 upstream.pipe(
519 Effect.flatMap((arr) => Queue.offerAll(queue, arr)),
520 Effect.forever({ disableYield: true }),
521 Pull.catchDone((_) => {
522 Queue.endUnsafe(queue)
523 return endVoid
524 })
525 )
526 )
527
528/**
529 * Creates a sink that publishes every consumed input element to a `PubSub`.

Callers

nothing calls this directly

Calls 3

offerAllMethod · 0.80
fromTransformFunction · 0.70
pipeMethod · 0.65

Tested by

no test coverage detected