Hyperlinkv0.8.0-beta.28

PubSub

PubSub.takeconsteffect/PubSub.ts:1160
<A>(self: Subscription<A>): Effect.Effect<A>

Takes a single message from the subscription. If no messages are available, this will suspend until a message becomes available.

Example (Taking a message)

import { Effect, Fiber, 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)

    // Start a fiber to take a message (will suspend)
    const takeFiber = yield* Effect.forkChild(
      PubSub.take(subscription)
    )

    // Publish a message
    yield* PubSub.publish(pubsub, "Hello")

    // The take will now complete
    const message = yield* Fiber.join(takeFiber)
    console.log("Received:", message) // "Hello"
  }))
})
subscriptions
Source effect/PubSub.ts:116019 lines
export const take = <A>(self: Subscription<A>): Effect.Effect<A> =>
  Effect.suspend(() => {
    if (self.shutdownFlag.current) {
      return Effect.interrupt
    }
    if (self.replayWindow.remaining > 0) {
      const message = self.replayWindow.take()!
      return Effect.succeed(message)
    }
    const message = self.pollers.length === 0
      ? self.subscription.poll()
      : MutableList.Empty
    if (message === MutableList.Empty) {
      return pollForItem(self)
    } else {
      self.strategy.onPubSubEmptySpaceUnsafe(self.pubsub, self.subscribers)
      return Effect.succeed(message)
    }
  })
Referenced by 3 symbols