HubModule
Hub
Hub<'a> delivers every published value to every current subscription. Each subscription chooses its own
QueueStrategy, so one publisher can feed a lossless consumer and latest-value consumers at the same time, and every
subscriber sees values in the same order.
subscribe returns the subscription as a Dequeue owned by
the current scope; to subscribe for the life of a stream, use
FlowStream.fromHub. publish
waits while a BackPressure subscription is full, and no later value reaches any subscriber meanwhile;
tryPublish publishes only when no waiting is needed, to every subscription or to none.
tryPublishNow provides that immediate attempt to a synchronous callback and reports
whether it published, found another publisher busy, found a full subscription, or found the hub shut down.
shutdown lets every subscriber drain its backlog.
A subscriber that joins late sees only later values; for a current value followed by every change, use a
SubscriptionRef. Read the Hub guide,
and the hub torture test for every guarantee checked together.
Summary
| Name | Signature | Synopsis |
|---|---|---|
| Operations | ||
| make | Hub.make () | Creates a hub with no subscriptions. |
| makeScoped | Hub.makeScoped () | Creates a hub that is shut down when the current scope closes. |
| publish | Hub.publish value hub | Delivers a value to every current subscription. |
| publishAll | Hub.publishAll values hub | Publishes values in order, with no other publish interleaved, and sums their results. |
| tryPublish | Hub.tryPublish value hub | Publishes a value only if that needs no waiting; returns None without publishing otherwise. |
| tryPublishNow | Hub.tryPublishNow value hub | Tries to publish immediately from a synchronous callback. |
| subscribe | Hub.subscribe strategy hub | Subscribes to the hub with a buffering strategy. |
| subscriberCount | Hub.subscriberCount hub | Returns the number of current subscriptions. |
| shutdown | Hub.shutdown hub | Shuts the hub and all of its subscriptions down. |
| isShutdown | Hub.isShutdown hub | Returns whether the hub has been shut down. |
| awaitShutdown | Hub.awaitShutdown hub | Suspends until the hub is shut down. |
Operations
Parameters
| Name | Type | Description |
|---|---|---|
| strategy | QueueStrategy | |
| hub | Hub<'a> |
Returns
Verification Examples
flow {
let! (readings: Hub<float>) = Hub.make ()
let! latest = readings |> Hub.subscribe (QueueStrategy.Sliding 1)
do! readings |> Hub.publishAll [ 20.0; 21.5 ] |> Flow.ignore
return! Dequeue.takeAll latest
}
|> Flow.scoped
