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

NameSignatureSynopsis
Operations
makeHub.make ()Creates a hub with no subscriptions.
makeScopedHub.makeScoped ()Creates a hub that is shut down when the current scope closes.
publishHub.publish value hubDelivers a value to every current subscription.
publishAllHub.publishAll values hubPublishes values in order, with no other publish interleaved, and sums their results.
tryPublishHub.tryPublish value hubPublishes a value only if that needs no waiting; returns None without publishing otherwise.
tryPublishNowHub.tryPublishNow value hubTries to publish immediately from a synchronous callback.
subscribeHub.subscribe strategy hubSubscribes to the hub with a buffering strategy.
subscriberCountHub.subscriberCount hubReturns the number of current subscriptions.
shutdownHub.shutdown hubShuts the hub and all of its subscriptions down.
isShutdownHub.isShutdown hubReturns whether the hub has been shut down.
awaitShutdownHub.awaitShutdown hubSuspends until the hub is shut down.

Operations

kind:member

make

Hub.make ()
Member
Creates a hub with no subscriptions.

Parameters

NameTypeDescription
unit

Returns

Flow<'env, 'error, Hub<'a>>
kind:member

makeScoped

Hub.makeScoped ()
Member
Creates a hub that is shut down when the current scope closes.

Parameters

NameTypeDescription
unit

Returns

Flow<'env, 'error, Hub<'a>>
kind:member

publish

Hub.publish value hub
Member
Delivers a value to every current subscription.

Parameters

NameTypeDescription
value'a
hubHub<'a>

Returns

Flow<'env, 'error, PublishResult>

Verification Examples

flow {
    let! (readings: Hub<float>) = Hub.make ()
    return! readings |> Hub.publish 21.5
}
kind:member

publishAll

Hub.publishAll values hub
Member
Publishes values in order, with no other publish interleaved, and sums their results.

Parameters

NameTypeDescription
values'a seq
hubHub<'a>

Returns

Flow<'env, 'error, PublishResult>
kind:member

tryPublish

Hub.tryPublish value hub
Member
Publishes a value only if that needs no waiting; returns None without publishing otherwise.

Parameters

NameTypeDescription
value'a
hubHub<'a>

Returns

Flow<'env, 'error, PublishResult option>
kind:member

tryPublishNow

Hub.tryPublishNow value hub
Member
Tries to publish immediately from a synchronous callback.

Parameters

NameTypeDescription
value'a
hubHub<'a>

Returns

HubTryPublishResult
kind:member

subscribe

Hub.subscribe strategy hub
Member
Subscribes to the hub with a buffering strategy. The subscription ends when the current scope closes.

Parameters

NameTypeDescription
strategyQueueStrategy
hubHub<'a>

Returns

Flow<'env, 'error, Dequeue<'a>>

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
kind:member

subscriberCount

Hub.subscriberCount hub
Member
Returns the number of current subscriptions.

Parameters

NameTypeDescription
hubHub<'a>

Returns

Flow<'env, 'error, int>
kind:member

shutdown

Hub.shutdown hub
Member
Shuts the hub and all of its subscriptions down. Calling it again has no effect.

Parameters

NameTypeDescription
hubHub<'a>

Returns

Flow<'env, 'error, unit>
kind:member

isShutdown

Hub.isShutdown hub
Member
Returns whether the hub has been shut down.

Parameters

NameTypeDescription
hubHub<'a>

Returns

Flow<'env, 'error, bool>
kind:member

awaitShutdown

Hub.awaitShutdown hub
Member
Suspends until the hub is shut down.

Parameters

NameTypeDescription
hubHub<'a>

Returns

Flow<'env, 'error, unit>