Repository F# setup
open System
open System.IO
open System.Threading
open System.Threading.Tasks
open Axial
open Axial.Layers
open Axial.Console
open Axial.FileSystem
open Axial.Hosting
open Axial.Hosting.Browser
open Axial.Hosting.Node
open Axial.PlatformService
open Axial.State
open Axial.Telemetry
open Axial.Telemetry.JavaScriptHub
Hub<'a> delivers every published value to every current subscription. Each subscription chooses how it buffers, so
one publisher can feed a consumer that must see every value alongside consumers that only want the latest one.
Here one sensor feed goes to three subscribers: a historian that must record every reading, a display that only needs the newest reading, and an alarm panel that keeps the first reading it has not yet handled.
Shared setup
// Setup for the checked examples on this page.
open System
open System.IO
open System.Threading
open System.Threading.Tasks
open Axial
open Axial.Layers
open Axial.Console
open Axial.FileSystem
open Axial.Hosting
open Axial.Hosting.Browser
open Axial.Hosting.Node
open Axial.PlatformService
open Axial.State
open Axial.Telemetry
open Axial.Telemetry.JavaScript
/// Fails the docs test when an example's result differs from the value shown.
let shouldEqual expected actual =
if actual <> expected then failwithf "Expected %A but got %A" expected actual
SystemIOThreadingTasksAxialLayersConsoleFileSystemHostingBrowserNodePlatformServiceStateTelemetryJavaScriptshouldEqual: 'a -> 'a -> unitexpected: 'aactual: 'a(<>): 'T -> 'T -> boolStructural inequality The first parameter. The second parameter. The result of the comparison. 5 <> 5 // Evaluates to false 5 <> 6 // Evaluates to true [1; 2] <> [1; 2] // Evaluates to false
failwithf: Printf.StringFormat<'T,'Result> -> 'TPrint to a string buffer and raise an exception with the given result. Helper printers must return strings. The formatter. The formatted result. See Printf.failwithf (link: ) for examples.
> open Axial.PlatformService;;
> (flow {
- let! (hub: Hub<int>) = Hub.make ()
- let! historian = hub |> Hub.subscribe (QueueStrategy.BackPressure 2)
- let! display = hub |> Hub.subscribe (QueueStrategy.Sliding 1)
- let! alarms = hub |> Hub.subscribe (QueueStrategy.Dropping 1)
- let! recording = historian |> FlowStream.fromDequeue |> FlowStream.runCollect |> Flow.fork
- do! hub |> Hub.publishAll [ 1..5 ] |> Flow.ignore
- do! Hub.shutdown hub
- let! history = Fiber.join recording
- let! latest = Dequeue.takeAll display
- let! pending = Dequeue.takeAll alarms
- return [ history; latest; pending ]
- } : Flow<ClockEnvironment, Never, int list list>)
- |> Flow.run (ClockEnvironment Clock.live);;val it: Exit<int list list,Never> = Success [[1; 2; 3; 4; 5]; [5]; [1]]The three lists are what the historian, the display, and the alarm panel received. The display and alarm subscribers never took anything, yet the publisher was never held up by them. The historian's buffer holds two readings, so the publisher waited whenever the historian fell two readings behind.
Lossless delivery needs either back-pressure or unbounded memory. A BackPressure subscriber that stops taking values
eventually stops the publisher, and an Unbounded one grows without limit instead. While the publisher waits for a full
BackPressure subscriber, no later value reaches any subscriber, lossy ones included: every subscriber sees the same
order, so one stalled lossless subscriber stalls the whole feed. Size a lossless subscriber's buffer for the bursts you
expect, and use Dequeue.size on the subscription to raise an alarm before it fills.
Subscription strategies
A subscription takes the same QueueStrategy as a queue, and Hub.subscribe returns the subscription as
a Dequeue, so a subscriber consumes it with the Dequeue functions or FlowStream.fromDequeue.
| Strategy | Full buffer | Effect on the publisher |
|---|---|---|
QueueStrategy.BackPressure n |
the publisher waits for room | can suspend Hub.publish |
QueueStrategy.Sliding n |
the oldest value is evicted | never |
QueueStrategy.Dropping n |
the new value is discarded | never |
QueueStrategy.Unbounded |
never full | never |
Hub.publish returns a PublishResult that counts the subscriptions that accepted the value, the Dropping
subscriptions that discarded it, and the Sliding subscriptions that evicted an older value to take it, so the
publisher can report losses without inspecting its subscribers. PublishResult.add sums two results and
PublishResult.empty is the starting point, for totalling a run of tryPublish calls.
Use a hub when there can be zero or many consumers. Use a queue when there is exactly one logical consumer.
Publishing without waiting
A control loop that must keep its period cannot wait for a slow historian. Hub.tryPublish publishes only if it can
finish without waiting, and returns None without delivering anything otherwise, so no subscriber sees a value that
the others missed.
> open Axial.PlatformService;;
> (flow {
- let! (hub: Hub<int>) = Hub.make ()
- let! historian = hub |> Hub.subscribe (QueueStrategy.BackPressure 1)
- let! first = hub |> Hub.tryPublish 1
- let! second = hub |> Hub.tryPublish 2
- return [ first.IsSome; second.IsSome ]
- } : Flow<unit, Never, bool list>)
- |> Flow.run ();;val it: Exit<bool list,Never> = Success [true; false]The first value reached the historian. The second found its buffer full, so tryPublish delivered it nowhere.
For a synchronous host callback, Hub.tryPublishNow makes the same immediate attempt without starting a Flow. Its
result distinguishes Published result, Busy (another publisher has the turn), Full (a back-pressure subscriber
has no room), and Shutdown. Published result still reports any drops or evictions in lossy subscriptions. The
callback decides what to do with a value that was not published.
Do not bound Hub.publish with a timeout instead. An interrupted publish leaves the value with the subscribers it has
already reached, so publishing it again delivers it to them twice.
Subscriptions belong to a scope
Hub.subscribe registers the subscription with the current scope. Closing that scope removes the subscription and
releases a publisher waiting on it, so run subscribe inside Flow.scoped, a forked fiber, or an application root. A
forked fiber has a scope of its own, so a consumer fiber's subscription ends when the consumer does. A subscriber can
also leave early with Dequeue.shutdown, which has the same effect.
> open Axial.PlatformService;;
> (flow {
- let! (hub: Hub<string>) = Hub.make ()
- let! received =
- flow {
- let! events = hub |> Hub.subscribe QueueStrategy.Unbounded
- do! hub |> Hub.publish "while subscribed" |> Flow.ignore
- return! Dequeue.takeAll events
- }
- |> Flow.scoped
- do! hub |> Hub.publish "after the scope closed" |> Flow.ignore
- let! count = Hub.subscriberCount hub
- return received @ [ $"subscribers left: {count}" ]
- } : Flow<unit, Never, string list>)
- |> Flow.run ();;val it: Exit<string list,Never> = Success ["while subscribed"; "subscribers left: 0"]Streams
FlowStream.fromHub strategy hub subscribes when the stream starts and unsubscribes when it ends, so a stream consumer
never leaves a subscription behind. FlowStream.runIntoHub hub publishes every value of a stream.
let shown = ResizeArray<float>()
let render (reading: float) : Flow<ClockEnvironment, Never, unit> = Flow.delay (fun () -> shown.Add reading; Flow.ok ())
// A display that follows the latest reading for as long as it runs.
let display (readings: Hub<float>) : Flow<ClockEnvironment, Never, unit> =
readings
|> FlowStream.fromHub (QueueStrategy.Sliding 1)
|> FlowStream.runForEachFlow render
// A sensor stream feeding the hub.
let feed (readings: Hub<float>) : Flow<ClockEnvironment, Never, unit> =
FlowStream.fromSeq [ 20.5; 20.7; 21.0 ] |> FlowStream.runIntoHub readings
shown: ResizeArray<float>``.ctor``: unit -> unitInitializes a new instance of the class that is empty and has the default initial capacity.
floatAn abbreviation for the CLI type . Basic Types
render: float -> Flow<ClockEnvironment,Never,unit>reading: floatAxial.Flow`3Represents a cold workflow that reads an environment, returns a typed result, and is executed explicitly through one of its execution members such as ToTask, ToAsync, or RunSynchronously. The type of the environment dependency. The type of the failure value. The type of the success value.
Axial.ClockEnvironmentAn environment containing only a clock, for timed flows with no other services.
Axial.NeverRepresents an error channel that cannot occur.
unitThe type 'unit', which has only one value "()". This value is special and always uses the representation 'null'. Basic Types
Axial.Flowdelay: (unit -> Flow<'env,'error,'value>) -> Flow<'env,'error,'value>Defers flow construction until execution time. A function that returns the flow to execute. A flow that lazily evaluates the factory when executed. let flow = Flow.delay (fun () -> Flow.succeed 42)
Add: float -> unitAdds an object to the end of the . The object to be added to the end of the . The value can be for reference types.
ok: 'value -> Flow<'env,'error,'value>Creates a successful synchronous flow. The value to wrap in a successful flow. A flow that always succeeds with the provided value.
display: Hub<float> -> Flow<ClockEnvironment,Never,unit>readings: Hub<float>Axial.Hub`1Broadcasts every published value to every current subscription. Create one with Hub.make. Use a hub when there can be zero or many consumers. The type of the published values.
(|>): 'T1 -> ('T1 -> 'U) -> 'UApply a function to a value, the value being on the left, the function on the right The argument. The function. The function result. let doubleIt x = x * 2 3 |> doubleIt // Evaluates to 6
Axial.FlowStreamModulefromHub: QueueStrategy -> Hub<'value> -> FlowStream<'env,'error,'value>Creates a stream of the values published to a hub, subscribing for the life of the stream. The stream subscribes with when it starts and unsubscribes when it ends, so the subscription can never outlive its consumer. It ends normally once the hub is shut down and the backlog is drained. Values published before the stream starts are not delivered. When a consumer is forked and must see values published straight afterwards, subscribe first with Hub.subscribe and consume the subscription with FlowStream.fromDequeue. readings |> FlowStream.fromHub (QueueStrategy.Sliding 1) |> FlowStream.runForEach display
Axial.QueueStrategyWhat a queue, or a hub subscription, does with a value offered while it is full. Lossless delivery needs either back-pressure or unbounded memory. BackPressure keeps every value by making the producer wait; Unbounded keeps every value by growing without limit; Dropping and Sliding never make the producer wait and lose values instead. A non-positive capacity fails the flow that uses the strategy with a defect.
SlidingA full buffer evicts its oldest value to make room for the new one.
runForEachFlow: ('value -> Flow<'env,'error,unit>) -> FlowStream<'env,'error,'value> -> Flow<'env,'error,unit>Runs an effectful action for every stream value. stream |> FlowStream.runForEachFlow save
feed: Hub<float> -> Flow<ClockEnvironment,Never,unit>fromSeq: 'value seq -> FlowStream<'env,'error,'value>Creates a stream from a synchronous sequence of values. The sequence of values to be emitted by the stream. A that yields each value from the sequence. FlowStream.fromSeq [1..10] |> FlowStream.runCollect |> Flow.run ()
runIntoHub: Hub<'value> -> FlowStream<'env,'error,'value> -> Flow<'env,'error,unit>Runs a stream and publishes every value to a hub. Each value is published with Hub.publish, so a full BackPressure subscription slows the stream down. If the hub is shut down first, the flow is interrupted. sensor |> FlowStream.runIntoHub readings
flow {
let! (readings: Hub<float>) = Hub.make ()
let! displaying = display readings |> Flow.fork
// Wait until the display has subscribed, then publish and shut down.
let mutable subscribers = 0
while subscribers = 0 do
let! count = Hub.subscriberCount readings
subscribers <- count
if subscribers = 0 then
do! Flow.sleep (TimeSpan.FromMilliseconds 1.0)
do! feed readings
do! Hub.shutdown readings
do! Fiber.join displaying
}
|> Flow.run (ClockEnvironment Clock.live)
|> shouldEqual (Exit.Success())
shown |> Seq.last |> shouldEqual 21.0
flow: FlowBuilderThe universal flow { } computation expression.
readings: Hub<float>Axial.Hub`1Broadcasts every published value to every current subscription. Create one with Hub.make. Use a hub when there can be zero or many consumers. The type of the published values.
floatAn abbreviation for the CLI type . Basic Types
Axial.HubModuleCreates hubs, publishes to them, and subscribes to them. A publish reaches every subscription that exists when it starts; a subscriber sees only values published after its subscribe completes, so a late joiner that needs current state must get it elsewhere, for example from a SubscriptionRef. Publishes are serialized, so every subscriber observes values in the same order even with concurrent publishers. Dropping, Sliding, and Unbounded subscriptions never delay the publisher. A BackPressure subscription with a full buffer suspends publish until it has room or its subscription ends. While the publisher waits, no later value reaches any subscriber, lossy ones included: keeping one order for every subscriber means a stalled lossless subscriber stalls the feed. Size a lossless subscriber's buffer for the bursts you expect, watch Dequeue.size to raise an alarm before it fills, and use Hub.tryPublish where the publisher must never wait. Shutting a hub down shuts every subscription down with Dequeue.shutdown semantics: subscribers drain their backlog, streams over subscriptions end normally, and later publishes are interrupted.
make: unit -> Flow<'env,'error,Hub<'a>>Creates a hub with no subscriptions.
displaying: Fiber<Never,unit>display: Hub<float> -> Flow<ClockEnvironment,Never,unit>(|>): 'T1 -> ('T1 -> 'U) -> 'UApply a function to a value, the value being on the left, the function on the right The argument. The function. The function result. let doubleIt x = x * 2 3 |> doubleIt // Evaluates to 6
Axial.Flowfork: Flow<'env,'error,'value> -> Flow<'env,'none,Fiber<'error,'value>>Starts a flow in a new fiber without waiting for it to complete. Forking turns a cold flow description into hot child work and returns a handle that can later be joined or interrupted. Prefer zipPar or race when the caller only needs a simple parallel composition. Wait for the handle with Fiber.join or Fiber.await, and stop it with Fiber.interrupt. The flow to fork. A flow that produces a handle.
subscribers: int(=): 'T -> 'T -> boolStructural equality The first parameter. The second parameter. The result of the comparison. 5 = 5 // Evaluates to true 5 = 6 // Evaluates to false [1; 2] = [1; 2] // Evaluates to true (1, 5) = (1, 6) // Evaluates to false
count: intsubscriberCount: Hub<'a> -> Flow<'env,'error,int>Returns the number of current subscriptions.
sleep: TimeSpan -> Flow<'env,'error,unit>Suspends the flow for the specified duration, observing cancellation. The duration to sleep. A flow that completes after the specified delay, or is interrupted if cancelled first.
System.TimeSpanRepresents a time interval.
FromMilliseconds: float -> TimeSpanReturns a that represents a specified number of milliseconds. A number of milliseconds. An object that represents . is less than or greater than . -or- is . -or- is . is equal to .
feed: Hub<float> -> Flow<ClockEnvironment,Never,unit>shutdown: Hub<'a> -> Flow<'env,'error,unit>Shuts the hub and all of its subscriptions down. Calling it again has no effect.
Axial.FiberModuleOperations on a running , the handle returned by Flow.fork. Every operation returns a flow; nothing waits or interrupts until that flow runs. Reading a fiber's outcome through join, await, or interrupt marks it observed, so a defect it died with is not also reported as unobserved.
join: Fiber<'error,'value> -> Flow<'env,'error,'value>Waits for a fiber and returns its value, failing the same way the fiber failed. Joining preserves the child's error channel: a Cause.Fail becomes the same typed error, and interruption and defects remain interruption and defects. Use await to inspect the outcome instead. The fiber to join. A flow that completes with the fiber's value. let loadProfile : Flow<unit, string, string> = Flow.ok "profile" let loadOrders : Flow<unit, string, int list> = Flow.ok [ 1; 2 ] let page = flow { let! fiber = Flow.fork loadProfile let! orders = loadOrders let! profile = Fiber.join fiber return profile, orders }
run: 'env -> Flow<'env,'error,'value> -> Exit<'value,'error>Runs the workflow and blocks until the final exit is available. The environment used by the workflow. The workflow to run. The final workflow exit. let exit = workflow |> Flow.run environment
``.ctor``: IClock -> ClockEnvironmentAxial.PlatformService.ClockHelpers for the clock service.
live: IClockCreates a live clock backed by and a monotonic timer.
shouldEqual: 'a -> 'a -> unitAxial.Exit`2Represents the final outcome of a workflow execution. The type of the success value. The type of the domain-specific failure value.
SuccessThe workflow completed successfully.
shown: ResizeArray<float>Microsoft.FSharp.Collections.SeqModuleContains operations for working with values of type .
last: 'T seq -> 'TReturns the last element of the sequence. The input sequence. The last element of the sequence. Thrown when the input sequence is null. Thrown when the input does not have any elements. ["pear"; "banana"] |> Seq.last Evaluates to banana [] |> Seq.last Throws ArgumentException This is an O(n) operation, where n is the length of the sequence.
A Sliding 1 display may skip readings, but it always ends on the newest one.
fromHub does not see values published before the stream starts. When you fork a consumer and publish straight
afterwards, subscribe first with Hub.subscribe and consume the subscription with FlowStream.fromDequeue, as the
first example on this page does.
Ordering and late subscribers
Publishes are serialized: every subscriber observes values in the same order, even when several fibers publish
concurrently. A subscriber sees only values published after its subscribe completes. A consumer that joins late and
needs the current state first should follow a SubscriptionRef instead, which starts each
stream with the current value.
Shutdown
Hub.shutdown shuts every subscription down with the same semantics as Dequeue.shutdown. Subscribers can still take
their backlog, FlowStream.fromDequeue streams end normally once drained, and later publishes are interrupted.
Hub.makeScoped creates a hub that is shut down when the current scope closes. Hub.isShutdown reports the state, and
Hub.awaitShutdown suspends until the hub is shut down.
The hub torture test runs a hub with lossless, sliding, and dropping subscribers joining and leaving throughout, and checks every guarantee on this page.

