<A>(self: PubSub<A>): Effect.Effect<boolean>Returns true if the Pubsub contains zero elements, false otherwise.
Example (Checking whether a PubSub is empty)
import { Effect, PubSub } from "effect"
const program = Effect.gen(function*() {
const pubsub = yield* PubSub.bounded<string>(10)
// Initially empty
const initiallyEmpty = yield* PubSub.isEmpty(pubsub)
console.log("Initially empty:", initiallyEmpty) // true
yield* Effect.scoped(Effect.gen(function*() {
const subscription = yield* PubSub.subscribe(pubsub)
// Publish a message for the active subscription
yield* PubSub.publish(pubsub, "Hello")
const nowEmpty = yield* PubSub.isEmpty(pubsub)
console.log("Now empty:", nowEmpty) // false
yield* PubSub.take(subscription)
}))
})export const const isEmpty: <A>(
self: PubSub<A>
) => Effect.Effect<boolean>
Returns true if the Pubsub contains zero elements, false otherwise.
Example (Checking whether a PubSub is empty)
import { Effect, PubSub } from "effect"
const program = Effect.gen(function*() {
const pubsub = yield* PubSub.bounded<string>(10)
// Initially empty
const initiallyEmpty = yield* PubSub.isEmpty(pubsub)
console.log("Initially empty:", initiallyEmpty) // true
yield* Effect.scoped(Effect.gen(function*() {
const subscription = yield* PubSub.subscribe(pubsub)
// Publish a message for the active subscription
yield* PubSub.publish(pubsub, "Hello")
const nowEmpty = yield* PubSub.isEmpty(pubsub)
console.log("Now empty:", nowEmpty) // false
yield* PubSub.take(subscription)
}))
})
isEmpty = <function (type parameter) A in <A>(self: PubSub<A>): Effect.Effect<boolean>A>(self: PubSub<A>(parameter) self: {
pubsub: PubSub.Atomic<A>;
subscribers: PubSub.Subscribers<A>;
scope: Scope.Closeable;
shutdownHook: Latch.Latch;
shutdownFlag: MutableRef.MutableRef<boolean>;
strategy: PubSub.Strategy<A>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
self: interface PubSub<in out A>A PubSub<A> is an asynchronous message hub into which publishers can publish
messages of type A and subscribers can subscribe to take messages of type
A.
Example (Publishing and subscribing to messages)
import { Effect, PubSub } from "effect"
const program = Effect.gen(function*() {
// Create a bounded PubSub with capacity 10
const pubsub = yield* PubSub.bounded<string>(10)
// Subscribe and consume messages
yield* Effect.scoped(Effect.gen(function*() {
const subscription = yield* PubSub.subscribe(pubsub)
// Publish messages
yield* PubSub.publish(pubsub, "Hello")
yield* PubSub.publish(pubsub, "World")
const message1 = yield* PubSub.take(subscription)
const message2 = yield* PubSub.take(subscription)
console.log(message1, message2) // "Hello", "World"
}))
})
Companion namespace containing the low-level building blocks used by
PubSub, including atomic implementations, backing subscriptions, replay
windows, and delivery strategies.
PubSub<function (type parameter) A in <A>(self: PubSub<A>): Effect.Effect<boolean>A>): import EffectEffect.type Effect.Effect = /*unresolved*/ anyEffect<boolean> => import EffectEffect.map(const size: <A>(
self: PubSub<A>
) => Effect.Effect<number>
Returns the current number of messages retained by the PubSub for active
subscribers.
Details
If the PubSub has been shut down, the returned effect succeeds with 0.
The size is not a count of waiting subscribers or suspended publishers.
Example (Getting PubSub size)
import { Effect, PubSub } from "effect"
const program = Effect.gen(function*() {
const pubsub = yield* PubSub.bounded<string>(10)
// Initially empty
const initialSize = yield* PubSub.size(pubsub)
console.log("Initial size:", initialSize) // 0
yield* Effect.scoped(Effect.gen(function*() {
const subscription = yield* PubSub.subscribe(pubsub)
// Publish some messages for the active subscription
yield* PubSub.publish(pubsub, "msg1")
yield* PubSub.publish(pubsub, "msg2")
const afterPublish = yield* PubSub.size(pubsub)
console.log("After publishing:", afterPublish) // 2
yield* PubSub.takeAll(subscription)
}))
})
size(self: PubSub<A>(parameter) self: {
pubsub: PubSub.Atomic<A>;
subscribers: PubSub.Subscribers<A>;
scope: Scope.Closeable;
shutdownHook: Latch.Latch;
shutdownFlag: MutableRef.MutableRef<boolean>;
strategy: PubSub.Strategy<A>;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
self), (size: anysize) => size: anysize === 0)