Hyperlinkv0.8.0-beta.28

PubSub

PubSub.takeAllconsteffect/PubSub.ts:1208
<A>(self: Subscription<A>): Effect.Effect<Arr.NonEmptyArray<A>>

Takes all available messages from the subscription, suspending if no items are available.

Example (Taking all available messages)

import { Effect, PubSub } from "effect"

const program = Effect.gen(function*() {
  const pubsub = yield* PubSub.bounded<string>(10)

  yield* Effect.scoped(Effect.gen(function*() {
    const subscription = yield* PubSub.subscribe(pubsub)

    // Publish multiple messages
    yield* PubSub.publishAll(pubsub, ["msg1", "msg2", "msg3"])

    // Take all available messages at once
    const allMessages = yield* PubSub.takeAll(subscription)
    console.log("All messages:", allMessages) // ["msg1", "msg2", "msg3"]
  }))
})
subscriptions
Source effect/PubSub.ts:120819 lines
export const takeAll = <A>(self: Subscription<A>): Effect.Effect<Arr.NonEmptyArray<A>> =>
  Effect.suspend(function loop(value?: [A]): Effect.Effect<Arr.NonEmptyArray<A>> {
    if (self.shutdownFlag.current) {
      return Effect.interrupt
    }
    let as = self.pollers.length === 0
      ? self.subscription.pollUpTo(Number.POSITIVE_INFINITY)
      : []
    if (value) {
      as = value.concat(as)
    }
    self.strategy.onPubSubEmptySpaceUnsafe(self.pubsub, self.subscribers)
    if (self.replayWindow.remaining > 0) {
      return Effect.succeed(self.replayWindow.takeAll().concat(as) as Arr.NonEmptyArray<A>)
    } else if (!Arr.isArrayNonEmpty(as)) {
      return Effect.flatMap(pollForItem(self), (item) => loop([item]))
    }
    return Effect.succeed(as)
  })
Referenced by 1 symbols