Hyperlinkv0.8.0-beta.28

PubSub

PubSub.DroppingStrategyclasseffect/PubSub.ts:2517
DroppingStrategy<A>

Represents the dropping strategy for bounded PubSub values.

When to use

Use to keep publishers fast by dropping new messages when the PubSub is at capacity.

Details

A publish that arrives while the PubSub is full is dropped instead of waiting for capacity.

Gotchas

Subscribers may miss messages published while they are subscribed.

Example (Applying a dropping strategy)

import { Effect, PubSub } from "effect"

const program = Effect.gen(function*() {
  // Create PubSub with dropping strategy
  const pubsub = yield* PubSub.dropping<string>(2)

  // Or explicitly create with dropping strategy
  const customPubsub = yield* PubSub.make<string>({
    atomicPubSub: () => PubSub.makeAtomicBounded(2),
    strategy: () => new PubSub.DroppingStrategy()
  })

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

    // Fill the PubSub
    const pub1 = yield* PubSub.publish(pubsub, "msg1") // true
    const pub2 = yield* PubSub.publish(pubsub, "msg2") // true
    const pub3 = yield* PubSub.publish(pubsub, "msg3") // false (dropped)

    console.log("Publication results:", [pub1, pub2, pub3]) // [true, true, false]

    // Subscribers will only see the first two messages
    const messages = yield* PubSub.takeAll(subscription)
    console.log("Received messages:", messages) // ["msg1", "msg2"]
  }))
})
models
Source effect/PubSub.ts:251734 lines
export class DroppingStrategy<in out A> implements PubSub.Strategy<A> {
  get shutdown(): Effect.Effect<void> {
    return Effect.void
  }

  handleSurplus(
    _pubsub: PubSub.Atomic<A>,
    _subscribers: PubSub.Subscribers<A>,
    _elements: Iterable<A>,
    _isShutdown: MutableRef.MutableRef<boolean>
  ): Effect.Effect<boolean> {
    return Effect.succeed(false)
  }

  onPubSubEmptySpaceUnsafe(
    _pubsub: PubSub.Atomic<A>,
    _subscribers: PubSub.Subscribers<A>
  ): void {
    //
  }

  completePollersUnsafe(
    pubsub: PubSub.Atomic<A>,
    subscribers: PubSub.Subscribers<A>,
    subscription: PubSub.BackingSubscription<A>,
    pollers: MutableList.MutableList<Deferred.Deferred<A>>
  ): void {
    return strategyCompletePollersUnsafe(this, pubsub, subscribers, subscription, pollers)
  }

  completeSubscribersUnsafe(pubsub: PubSub.Atomic<A>, subscribers: PubSub.Subscribers<A>): void {
    return strategyCompleteSubscribersUnsafe(this, pubsub, subscribers)
  }
}
Referenced by 2 symbols