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.JavaScriptFibers
Fibers represent running child workflows.
In Axial, a Fiber is a handle to a running Flow. A flow is a cold description of work. A fiber is the hot execution that exists after that work has been started in the background.
The Mental Model
While a Flow is cold (a description of work that hasn't started yet), a Fiber is hot (the work is currently being executed).
When you fork a flow, you are saying: start this work now, give me a typed handle to it, and let the current workflow continue. That handle is the fiber.
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.
let loadBoth left right =
flow {
let! leftFiber = Flow.fork left
let! rightValue = right
let! leftValue = Fiber.join leftFiber
return leftValue, rightValue
}
loadBoth: Flow<'a,'b,'c> -> Flow<'a,'b,'d> -> Flow<'a,'b,('c * 'd)>left: Flow<'a,'b,'c>right: Flow<'a,'b,'d>flow: FlowBuilderThe universal flow { } computation expression.
leftFiber: Fiber<'b,'c>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.
rightValue: 'dleftValue: 'cAxial.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 }
loadBoth (Flow.ok 1 : Flow<ClockEnvironment, string, int>) (Flow.ok "two") |> Flow.run (ClockEnvironment Clock.live) |> shouldEqual (Exit.Success(1, "two"))
loadBoth: Flow<'a,'b,'c> -> Flow<'a,'b,'d> -> Flow<'a,'b,('c * 'd)>Axial.Flowok: '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.
Axial.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.
stringAn abbreviation for the CLI type . Basic Types
intAn abbreviation for the CLI type . Basic Types
(|>): '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
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.
The example starts left in the background, runs right in the current workflow, then joins the child fiber before returning.
Structured Concurrency
Fibers are the foundation of Structured Concurrency in Axial. Unlike "fire-and-forget" background tasks, Fibers allow you to maintain a parent-child relationship between workflows, ensuring that background work is always accounted for and safely cleaned up.
The primary operations for managing fibers are:
Flow.fork: starts a flow in the background and returns aFiber<'error, 'value>handle.Fiber.join: waits for the fiber and resumes with its successful value or typed failure.Fiber.await: waits for the fiber and returns itsExitwithout failing, for callers that decide what the outcome means.Fiber.poll: returns the fiber'sExitif it has settled, without waiting.Fiber.interrupt: asks the fiber to stop, then waits for the child workflow to report its finalExit.Flow.forkDetached: starts deliberate fire-and-forget work whose defects are never reported as unobserved.Flow.forkNamed: forks with a diagnostic name that carries into dumps and telemetry fiber spans, so long-lived background fibers are recognizable instead of bare ids.Fiber.dump: returns a diagnostic snapshot of one fiber handle.Flow.forkGracefulstop grace: forks a fiber that, when its scope closes, is asked to stop withstopand given up tograceto finish before it is interrupted. See stopping a consumer gracefully.
Joining, awaiting, or interrupting a fiber marks its outcome as observed. A fiber whose handle is simply discarded and that later dies with a defect is reported through the runtime's fiber observer; use Flow.forkDetached when the silence is intentional, and Flow.supervise to restart background work that dies with defects.
Why Fibers?
Fibers provide several advantages over raw Task or Async values:
Interruption
In ordinary .NET code, cancellation often depends on manually threading a CancellationToken through every layer. In Axial, interruption is part of the execution model. Fiber.interrupt signals the child fiber and waits for it to finish, so callers can observe the final Exit<'value, 'error>.
Typed Outcomes
A Fiber<'error, 'value> remembers the error type and success type of the workflow it is running. When you Fiber.join a fiber, the joined flow has the same typed failure channel as the child.
Clear Ownership
Fibers make background work visible in the workflow that started it. If the parent needs the result, it joins. If the parent no longer needs the result, it interrupts. An untracked task, by contrast, fails unnoticed unless some other layer checks it.
Diagnostics
Every forked fiber carries metadata:
FiberId: A unique runtime id for the child fiber.Name: The diagnostic name fromFlow.forkNamed, if one was given.ParentId: The id of the fiber that calledFlow.fork.Annotations: The runtime annotations (Flow.annotate) in scope at the fork site.StartedAt/SettledAt: UTC timestamps for fork and settle.Status:Running,Succeeded,Failed, orInterrupted.
Use Fiber.dump when logging or debugging one fiber. The dump is a snapshot, so a running fiber can report Running before Fiber.join and Succeeded, Failed, or Interrupted afterward. To see every live fiber at once as a parent/child tree, install a FiberRegistry with Flow.withFiberRegistry and call registry.DumpAt(clock); see Observability.
Underlying Implementation
On .NET, a fiber wraps a Task<Exit<'value, 'error>>, a CancellationTokenSource, and diagnostic metadata. On Fable, it wraps an Async<Exit<'value, 'error>> with the same public model. Only Metadata is public; the task and cancellation source are reached through the Fiber functions, so observation and interruption always go through the runtime's bookkeeping.
Concurrency Primitives
Most code should not manage fibers manually. Prefer high-level parallel combinators when they express the whole relationship:
Flow.zipPar: Runs two flows concurrently in separate fibers and waits for both.Flow.race: Runs two flows concurrently and returns the result of the winner, interrupting the loser.Flow.traversePar: Maps many values with bounded concurrency and returns the results in input order.Flow.forEachPar: Runs a flow for each value with bounded concurrency, discarding the results.
let fetchPage (url: string) : Flow<ClockEnvironment, string, int> =
Flow.sleep (TimeSpan.FromMilliseconds 1.0) |> Flow.map (fun () -> url.Length)
let pageSizes (urls: string list) : Flow<ClockEnvironment, string, int list> =
urls |> Flow.traversePar (Parallelism.bounded 8) fetchPage
let indexed = ResizeArray<string>()
let indexFile (file: string) : Flow<ClockEnvironment, string, unit> =
Flow.delay (fun () -> lock indexed (fun () -> indexed.Add file); Flow.ok ())
let indexAll (files: string list) : Flow<ClockEnvironment, string, unit> =
files |> Flow.forEachPar (Parallelism.ofProcessors id) indexFile
fetchPage: string -> Flow<ClockEnvironment,string,int>url: stringstringAn abbreviation for the CLI type . Basic Types
Axial.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.
intAn abbreviation for the CLI type . Basic Types
Axial.Flowsleep: 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 .
(|>): '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
map: ('value -> 'next) -> Flow<'env,'error,'value> -> Flow<'env,'error,'next>Transforms the successful value of a flow. If the source fails, the is not executed. The original failure cause is preserved, including typed failures, interruption, and defects. Use map for pure value transformations after an effect has succeeded. A function of type 'value -> 'next to transform the successful value. The source flow of type to transform. A new with the transformed success value of type 'next. let flow = Flow.succeed 1 |> Flow.map (fun x -> x + 1)
Length: intGets the number of characters in the current object. The number of characters in the current string.
pageSizes: string list -> Flow<ClockEnvironment,string,int list>urls: string listlistThe type of immutable singly-linked lists. See the module for further operations related to lists. Use the constructors [] and :: (infix) to create values of this type, or the notation [1; 2; 3]. Use the values in the List module to manipulate values of this type, or pattern match against the values directly. See also F# Language Guide - Lists.
traversePar: Parallelism -> ('value -> Flow<'env,'error,'next>) -> 'value seq -> Flow<'env,'error,'next list>Maps values to flows and runs them with bounded concurrency, returning results in input order. At most mappings run at once; as each finishes, its worker starts the next value. The first failure interrupts the mappings still running and waits for their cleanup before the flow fails, so no sibling is left running in the background. Size CPU-bound work with Parallelism.ofProcessors. When each worker needs its own connection or handle, use traverseParUsing. The maximum number of mappings running at once. Maps each value to a flow. The values to map. A flow containing the mapped values in the order of . let! pages = urls |> Flow.traversePar (Parallelism.bounded 8) fetchPage
Axial.ParallelismModuleCreates bounds for parallel Flow and stream operators.
bounded: int -> ParallelismCreates a positive concurrency bound. Thrown when is not positive.
indexed: ResizeArray<string>``.ctor``: unit -> unitInitializes a new instance of the class that is empty and has the default initial capacity.
indexFile: string -> Flow<ClockEnvironment,string,unit>file: stringunitThe type 'unit', which has only one value "()". This value is special and always uses the representation 'null'. Basic Types
delay: (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)
lock: 'Lock -> (unit -> 'T) -> 'TExecute the function as a mutual-exclusion region using the input value as a lock. The object to be locked. The action to perform during the lock. The resulting value. open System.Linq /// A counter object, supporting unlocked and locked increment type TestCounter () = let mutable count = 0 /// Increment the counter, unlocked member this.IncrementWithoutLock() = count <- count + 1 /// Increment the counter, locked member this.IncrementWithLock() = lock this (fun () -> count <- count + 1) /// Get the count member this.Count = count let counter = TestCounter() // Create a parallel sequence to that uses all our CPUs (seq {1..100000}).AsParallel() .ForAll(fun _ -> counter.IncrementWithoutLock()) // Evaluates to a number between 1-100000, non-deterministically because there is no locking counter.Count let counter2 = TestCounter() // Create a parallel sequence to that uses all our CPUs (seq {1..100000}).AsParallel() .ForAll(fun _ -> counter2.IncrementWithLock()) // Evaluates to 100000 deterministically because the increment to the counter object is locked counter2.Count
Add: string -> 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.
indexAll: string list -> Flow<ClockEnvironment,string,unit>files: string listforEachPar: Parallelism -> ('value -> Flow<'env,'error,unit>) -> 'value seq -> Flow<'env,'error,unit>Runs a flow for each value with bounded concurrency, discarding the results. Runs like traversePar: at most at once, and the first failure interrupts the rest. The maximum number of flows running at once. The flow to run for each value. The values to process. do! files |> Flow.forEachPar (Parallelism.ofProcessors id) indexFile
ofProcessors: (int -> int) -> ParallelismSizes a concurrency bound from the number of processors available to the process. Use this for CPU-bound work. The result is clamped to at least 1, so fun n -> n / 2 is safe on a single-core machine. The processor count is read once, when this is called. On JavaScript the count is 1. Computes the bound from the processor count. files |> Flow.traversePar (Parallelism.ofProcessors id) hashFile
id: 'T -> 'TThe identity function The input value. The same value. id 12 // Evaluates to 12 id "abc" // Evaluates to "abc"
pageSizes [ "a"; "bb"; "ccc" ] |> Flow.run (ClockEnvironment Clock.live) |> shouldEqual (Exit.Success [ 1; 2; 3 ])
indexAll [ "b.fs"; "a.fs"; "c.fs" ] |> Flow.run (ClockEnvironment Clock.live) |> shouldEqual (Exit.Success())
indexed |> Seq.sort |> List.ofSeq |> shouldEqual [ "a.fs"; "b.fs"; "c.fs" ]
pageSizes: string list -> Flow<ClockEnvironment,string,int list>(|>): '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.Flowrun: '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.
indexAll: string list -> Flow<ClockEnvironment,string,unit>indexed: ResizeArray<string>Microsoft.FSharp.Collections.SeqModuleContains operations for working with values of type .
sort: 'T seq -> 'T seqYields a sequence ordered by keys. This function returns a sequence that digests the whole initial sequence as soon as that sequence is iterated. As a result this function should not be used with large or infinite sequences. The function makes no assumption on the ordering of the original sequence and uses a stable sort, that is the original order of equal elements is preserved. This is an O(n log n) operation, where n is the length of the sequence. The input sequence. The result sequence. Thrown when the input sequence is null. let input = seq { 8; 4; 3; 1; 6; 1 } Seq.sort input Evaluates to a sequence yielding the same results as seq { 1; 1 3; 4; 6; 8 }.
Microsoft.FSharp.Collections.ListModuleContains operations for working with values of type . Operations for collections such as lists, arrays, sets, maps and sequences. See also F# Collection Types in the F# Language Guide.
ofSeq: 'T seq -> 'T listBuilds a new list from the given enumerable object. The input sequence. The list of elements from the sequence. let inputs = seq { 1; 2; 5 } inputs |> List.ofSeq Evaluates to [ 1; 2; 5 ]. This is an O(n) operation, where n is the length of the sequence.
traversePar returns results in input order, whatever order the work finishes in.
At most the given number of flows run at once, and each worker starts the next value as soon as it finishes one. The first failure interrupts the flows still running and waits for their cleanup, so no sibling keeps running after the traversal has failed. Size CPU-bound work with Parallelism.ofProcessors, which clamps to at least 1.
When each worker needs its own connection or handle, use Flow.traverseParUsing or Flow.forEachParUsing with a
Resource. Each worker acquires the resource once when it starts, reuses it for every value it takes, and releases it
when it finishes or the traversal fails, so at most parallelism resources exist at once:
let readersOpened = ref 0
let readersClosed = ref 0
let openReader : Resource<ClockEnvironment, string, int> =
Resource.create
(Flow.delay (fun () -> Flow.ok (Interlocked.Increment &readersOpened.contents)))
(fun _ _ ->
Interlocked.Increment &readersClosed.contents |> ignore
Task.CompletedTask)
let searchCommit (reader: int) (commit: int) : Flow<ClockEnvironment, string, bool> = Flow.ok (commit % 2 = 0)
let evenCommits (commits: int list) : Flow<ClockEnvironment, string, bool list> =
commits |> Flow.traverseParUsing (Parallelism.bounded 2) openReader searchCommit
readersOpened: int refref: 'T -> 'T refCreate a mutable reference cell The value to contain in the cell. The created reference cell. let count = ref 0 // Creates a reference cell object with a mutable Value property count.Value // Evaluates to 0 count.Value <- 1 // Updates the value count.Value // Evaluates to 1
readersClosed: int refopenReader: Resource<ClockEnvironment,string,int>Axial.Resource`3Describes acquisition of a value together with registration of its release in the current Flow scope. The environment required to acquire the value. The typed acquisition failure. The acquired value.
Axial.ClockEnvironmentAn environment containing only a clock, for timed flows with no other services.
stringAn abbreviation for the CLI type . Basic Types
intAn abbreviation for the CLI type . Basic Types
Axial.ResourceModulecreate: Flow<'env,'error,'value> -> ('value -> CancellationToken -> Task) -> Resource<'env,'error,'value>Describes acquisition together with a task-based release registered in the current Flow scope.
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)
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.
System.Threading.InterlockedProvides atomic operations for variables that are shared by multiple threads.
Increment: byref<int> -> intIncrements a specified variable and stores the result, as an atomic operation. The variable whose value is to be incremented. The incremented value. The address of is a null pointer.
(~&): 'T -> byref<'T>Address-of. Uses of this value may result in the generation of unverifiable code. The input object. The managed pointer.
contents: 'TThe current value of the reference cell
(|>): '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
ignore: 'T -> unitIgnore the passed value. This is often used to throw away results of a computation. The value to ignore. ignore 55555 // Evaluates to ()
System.Threading.Tasks.TaskRepresents an asynchronous operation.
CompletedTask: TaskGets a task that has already completed successfully. The successfully completed task.
searchCommit: int -> int -> Flow<ClockEnvironment,string,bool>reader: intcommit: intAxial.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.
boolAn abbreviation for the CLI type . Basic Types
(%): ^T1 -> ^T2 -> ^T3Overloaded modulo operator The first parameter. The second parameter. The result of the operation. 29 % 5 // Evaluates to 4
(=): '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
evenCommits: int list -> Flow<ClockEnvironment,string,bool list>commits: int listlistThe type of immutable singly-linked lists. See the module for further operations related to lists. Use the constructors [] and :: (infix) to create values of this type, or the notation [1; 2; 3]. Use the values in the List module to manipulate values of this type, or pattern match against the values directly. See also F# Language Guide - Lists.
traverseParUsing: Parallelism -> Resource<'env,'error,'resource> -> ('resource -> 'value -> Flow<'env,'error,'next>) -> 'value seq -> Flow<'env,'error,'next list>Like traversePar, giving each worker its own resource, acquired once and reused for every value it maps. Each worker acquires when it starts and releases it when it has no values left, or when the traversal fails or is interrupted. At most resources exist at once. Use it for work that needs an expensive handle per thread of work, such as a database connection or a repository reader that is not safe to share. The maximum number of workers, and so of resources. The resource each worker acquires. Maps a value to a flow, given the worker's resource. The values to map. A flow containing the mapped values in the order of . let! matches = commits |> Flow.traverseParUsing (Parallelism.ofProcessors id) openReader (fun reader commit -> search reader commit)
Axial.ParallelismModuleCreates bounds for parallel Flow and stream operators.
bounded: int -> ParallelismCreates a positive concurrency bound. Thrown when is not positive.
evenCommits [ 1..6 ] |> Flow.run (ClockEnvironment Clock.live) |> shouldEqual (Exit.Success [ false; true; false; true; false; true ])
readersOpened.Value <= 2 |> shouldEqual true
readersClosed.Value |> shouldEqual readersOpened.Value
evenCommits: int list -> Flow<ClockEnvironment,string,bool list>(..): ^T -> ^T -> ^T seqThe standard overloaded range operator, e.g. [n..m] for lists, seq {n..m} for sequences The start value of the range. The end value of the range. The sequence spanning the range. [1..4] // Evaluates to [1; 2; 3; 4] [1.5..4.4] // Evaluates to [1.5; 2.5; 3.5] ['a'..'d'] // Evaluates to ['a'; 'b'; 'c'; 'd'] [|1..4|] // Evaluates to an array [|1; 2; 3; 4|] { 1..4 } // Evaluates to a sequence [1; 2; 3; 4])
(|>): '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.Flowrun: '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.
readersOpened: int refValue: intThe current value of the reference cell
(<=): 'T -> 'T -> boolStructural less-than-or-equal comparison The first parameter. The second parameter. The result of the comparison. 5 <= 1 // Evaluates to false 5 <= 5 // Evaluates to true [1; 5] <= [1; 6] // Evaluates to true
readersClosed: int refFlowStream.mapFlowParUsing does the same for a stream, lending each running mapping a resource from a pool of at
most parallelism.
Use explicit fibers when the parent workflow needs to start child work, do something else, and decide later whether to join or interrupt it.
Latest Wins
When only the newest request matters, such as a search box or autocomplete, hold the running fiber in a FiberSlot and
fork with Flow.forkReplacing. Each fork signals the previous fiber in the slot to stop and does not wait for it, so
the new request starts at once:
let search (query: string) : Flow<ClockEnvironment, string, string> =
Flow.sleep (TimeSpan.FromMilliseconds 50.0) |> Flow.map (fun () -> $"results for {query}")
/// Forks a search per keystroke into one slot, and reports the last result and whether each earlier search stopped.
let typeAhead (queries: string list) : Flow<ClockEnvironment, string, string * bool list> =
flow {
let! slot = FiberSlot.make ()
let! searches = queries |> Flow.traverse (fun query -> search query |> Flow.forkReplacing slot)
let! last = searches |> List.last |> Fiber.join
let! earlier = searches |> List.take (searches.Length - 1) |> Flow.traverse Fiber.await
return last, earlier |> List.map (fun exit -> match exit with Exit.Failure cause -> Cause.isInterrupted cause | _ -> false)
}
search: string -> Flow<ClockEnvironment,string,string>query: stringstringAn abbreviation for the CLI type . Basic Types
Axial.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.Flowsleep: 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 .
(|>): '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
map: ('value -> 'next) -> Flow<'env,'error,'value> -> Flow<'env,'error,'next>Transforms the successful value of a flow. If the source fails, the is not executed. The original failure cause is preserved, including typed failures, interruption, and defects. Use map for pure value transformations after an effect has succeeded. A function of type 'value -> 'next to transform the successful value. The source flow of type to transform. A new with the transformed success value of type 'next. let flow = Flow.succeed 1 |> Flow.map (fun x -> x + 1)
typeAhead: string list -> Flow<ClockEnvironment,string,(string * bool list)>queries: string listlistThe type of immutable singly-linked lists. See the module for further operations related to lists. Use the constructors [] and :: (infix) to create values of this type, or the notation [1; 2; 3]. Use the values in the List module to manipulate values of this type, or pattern match against the values directly. See also F# Language Guide - Lists.
boolAn abbreviation for the CLI type . Basic Types
flow: FlowBuilderThe universal flow { } computation expression.
slot: FiberSlot<string,string>Axial.FiberSlotModuleCreates and clears and values.
make: unit -> Flow<'env,'none,FiberSlot<'error,'value>>Creates an empty slot for Flow.forkReplacing.
searches: Fiber<string,string> listtraverse: ('value -> Flow<'env,'error,'next>) -> 'value seq -> Flow<'env,'error,'next list>Transforms a sequence of values into a flow and stops at the first failure. A function that maps each value to a flow. The sequence of values to transform. A flow containing a list of the successful mapped values. let flows = [1; 2; 3] |> Flow.traverse (fun x -> Flow.succeed (x * 2))
forkReplacing: FiberSlot<'error,'value> -> Flow<'env,'error,'value> -> Flow<'env,'none,Fiber<'error,'value>>Starts a flow in a new fiber held by , interrupting the fiber it replaces. Latest wins: the previous fiber in the slot is signalled to stop and is not waited for, so the new request starts at once. Its outcome is marked observed. Use this for UI requests where a result for stale input is worthless, such as search-as-you-type. The slot that holds the current fiber. The flow to run. A flow that produces the new fiber's handle. let onQueryChanged slot query = search query |> Flow.forkReplacing slot |> Flow.ignore
last: stringMicrosoft.FSharp.Collections.ListModuleContains operations for working with values of type . Operations for collections such as lists, arrays, sets, maps and sequences. See also F# Collection Types in the F# Language Guide.
last: 'T list -> 'TReturns the last element of the list. The input list. The last element of the list. Thrown when the input does not have any elements. [ "pear"; "banana" ] |> List.last Evaluates to banana [ ] |> List.last Throws ArgumentException Lists are represented as linked lists so this is an O(n) operation, where n is the length of the list.
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 }
earlier: Exit<string,string> listtake: int -> 'T list -> 'T listReturns the first N elements of the list. Throws InvalidOperationException if the count exceeds the number of elements in the list. List.truncate returns as many items as the list contains instead of throwing an exception. This is an O(count) operation. The number of items to take. The input list. The result list. Thrown when the input list is empty. Thrown when count exceeds the number of elements in the list. let inputs = ["a"; "b"; "c"; "d"] inputs |> List.take 2 Evaluates to ["a"; "b"] let inputs = ["a"; "b"; "c"; "d"] inputs |> List.take 6 Throws InvalidOperationException. let inputs = ["a"; "b"; "c"; "d"] inputs |> List.take 0 Evaluates to the empty list.
Length: intGets the number of items contained in the list
(-): ^T1 -> ^T2 -> ^T3Overloaded subtraction operator The first parameter. The second parameter. The result of the operation. 10 - 2 // Evaluates to 8
await: Fiber<'error,'value> -> Flow<'env,'none,Exit<'value,'error>>Waits for a fiber and returns its , never failing. Use this when the caller decides what the fiber's outcome means, for example a cache that shares one computation between callers, or a supervisor that restarts children. Unlike Deferred.await, which returns the completed value, this returns the whole exit. Interrupting the awaiting flow stops the wait, not the fiber. The fiber to wait for. A flow that succeeds with the fiber's exit. flow { match! Fiber.await fiber with | Exit.Success value -> return Some value | Exit.Failure _ -> return None }
map: ('T -> 'U) -> 'T list -> 'U listBuilds a new collection whose elements are the results of applying the given function to each of the elements of the collection. The function to transform elements from the input list. The input list. The list of transformed elements. let inputs = [ "a"; "bbb"; "cc" ] inputs |> List.map (fun x -> x.Length) Evaluates to [ 1; 3; 2 ] This is an O(n) operation, where n is the length of the list.
exit: Exit<string,string>Axial.Exit`2Represents the final outcome of a workflow execution. The type of the success value. The type of the domain-specific failure value.
FailureThe workflow failed due to a specific cause.
cause: Cause<string>Axial.CauseisInterrupted: Cause<'error> -> boolReturns whether the cause tree contains an interruption signal.
typeAhead [ "a"; "ax"; "axi" ] |> Flow.run (ClockEnvironment Clock.live) |> shouldEqual (Exit.Success("results for axi", [ true; true ]))
typeAhead: string list -> Flow<ClockEnvironment,string,(string * bool list)>(|>): '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.Flowrun: '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.
Flow.forkReplacingKey key slots does the same per key, with a slot from FiberSlot.makeKeyed, so a new preview for
one document replaces only that document's previous load. A key's entry is removed when its fiber settles.
FiberSlot.interrupt and FiberSlot.interruptAll stop what is running and wait for its cleanup, for example on
shutdown. FiberSlot.count returns how many keyed fibers are running. For a stream of inputs,
FlowStream.switchMapFlow applies the same rule inside the stream.
What a fiber owns
Each forked fiber runs in its own scope, a child of the scope that forked it. Whatever the fiber acquires, such as a
resource registered with Flow.scopeAcquireRelease, a hub subscription, or a
scoped queue, is released when the fiber settles, not when the parent scope
eventually closes. A consumer fiber that ends therefore cannot leave a subscription behind that holds up a publisher.
Closing the parent scope interrupts fibers that are still running, then releases what they acquired.
Stopping a consumer gracefully
Interrupting a consumer the moment its scope closes throws away its backlog. Flow.forkGraceful changes what closing
the scope does: it runs a stop request, waits up to a grace period for the fiber to finish, and interrupts it only if it
is still running after that. For a queue consumer, the stop request is Dequeue.shutdown: the consumer's stream then
ends normally once it has drained the queue.
> open Axial.PlatformService;;
> (flow {
- let written = ResizeArray<int>()
- do!
- flow {
- let! (samples: Queue<int>) = Queue.bounded 100
- let! _ =
- samples
- |> FlowStream.fromDequeue
- |> FlowStream.runForEach written.Add
- |> Flow.forkGraceful (Dequeue.shutdown samples) (System.TimeSpan.FromSeconds 5.0)
- do! samples |> Queue.offerAll [ 1..5 ] |> Flow.ignore
- }
- |> Flow.scoped
- return List.ofSeq written
- } : Flow<ClockEnvironment, Never, int list>)
- |> Flow.run (ClockEnvironment Clock.live);;val it: Exit<int list,Never> = Success [1; 2; 3; 4; 5]The scope closed right after the offers, before the consumer had necessarily taken any of them. Closing it shut the
queue down, and the consumer finished all five values before the scope finished closing. With a plain Flow.fork, the
consumer would have been interrupted with part of the backlog still queued.
The pipeline torture test stops a control loop and a historian this way and checks that no sample is lost.
A graceful fiber is also not interrupted when the flow that forked it is interrupted: its stop request runs when that
flow's scope closes, so a consumer still flushes when the application is cancelled. Fiber.interrupt still interrupts
it immediately.

