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.JavaScriptTime and Latest Value
Most stream operators react only to values. The operators on this page also react to time, or to a newer value arriving, which is what progress reporting, batching, and search-as-you-type need.
Each of them moves upstream values through a one-slot queue from a producer fiber, so it can wait for "the next value
or a timer" at once. The producer runs in the consuming Flow's scope: it keeps upstream back-pressure, and it stops
when the consumer finishes, fails, or stops early with take.
Batch by size or time
FlowStream.groupedWithin size window emits a list when it holds size values or when window has passed since the
list's first value, whichever comes first. A slow trickle is not held back waiting for a full batch:
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 batchSizes (events: FlowStream<ClockEnvironment, Never, int>) : Flow<ClockEnvironment, Never, int list> =
events
|> FlowStream.groupedWithin 100 (TimeSpan.FromSeconds 1.0)
|> FlowStream.map List.length
|> FlowStream.runCollect
batchSizes: FlowStream<ClockEnvironment,Never,int> -> Flow<ClockEnvironment,Never,int list>events: FlowStream<ClockEnvironment,Never,int>Axial.FlowStream`3Represents a cold stream of values that requires an environment, can fail with a typed error, and supports backpressure. The type of the environment dependency. The type of the failure value. The type of the success values in the stream.
Axial.ClockEnvironmentAn environment containing only a clock, for timed flows with no other services.
Axial.NeverRepresents an error channel that cannot occur.
intAn 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.
listThe 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.
(|>): '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.FlowStreamModulegroupedWithin: int -> TimeSpan -> FlowStream<'env,'error,'value> -> FlowStream<'env,'error,'value list>Groups values into lists of at most , emitting early when passes. A group starts with the next value and is emitted when it holds values or when has passed since its first value, whichever comes first. No empty group is emitted. When upstream ends or fails, the partial group is emitted first. Use it to batch writes or progress updates without waiting indefinitely for a full batch. events |> FlowStream.groupedWithin 100 (TimeSpan.FromSeconds 1.0)
System.TimeSpanRepresents a time interval.
FromSeconds: float -> TimeSpanReturns a that represents a specified number of seconds, where the specification is accurate to the nearest millisecond. A number of seconds, accurate to the nearest millisecond. An object that represents . is less than or greater than . -or- is . -or- is . is equal to .
map: ('v -> 'w) -> FlowStream<'env,'error,'v> -> FlowStream<'env,'error,'w>Transforms the successful values of a stream using the provided function. The function to transform each value. The stream whose values should be transformed. A new stream that yields transformed values. let stream = FlowStream.fromSeq [1; 2; 3] |> FlowStream.map (fun n -> n * 2)
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.
length: 'T list -> intReturns the length of the list. The input list. The length of the list. The notation list.Length is preferred. This is an O(n) operation, where n is the length of the list. let inputs = [ "a"; "b"; "c" ] inputs |> List.length Evaluates to 3 This is an O(n) operation, where n is the length of the list.
runCollect: FlowStream<'a,'b,'c> -> Flow<'a,'b,'c list>Collects all emitted values into a list. stream |> FlowStream.runCollect
A burst of 250 events arrives well within a second, so the batches fill by size and the remainder is emitted when the stream ends:
batchSizes (FlowStream.fromSeq [ 1..250 ]) |> Flow.run (ClockEnvironment Clock.live) |> shouldEqual (Exit.Success [ 100; 100; 50 ])
batchSizes: FlowStream<ClockEnvironment,Never,int> -> Flow<ClockEnvironment,Never,int list>Axial.FlowStreamModulefromSeq: '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 ()
(..): ^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.
No empty list is emitted, and the partial list is emitted when upstream ends or before a failure is propagated.
Report progress at a steady rate
FlowStream.throttle interval emits at most one value per interval. The first value is emitted immediately; values
arriving faster replace each other, and the latest is emitted when the interval ends:
let rendered (progress: FlowStream<ClockEnvironment, Never, int>) : Flow<ClockEnvironment, Never, int list> =
progress
|> FlowStream.throttle (TimeSpan.FromMilliseconds 100.0)
|> FlowStream.runCollect
rendered: FlowStream<ClockEnvironment,Never,int> -> Flow<ClockEnvironment,Never,int list>progress: FlowStream<ClockEnvironment,Never,int>Axial.FlowStream`3Represents a cold stream of values that requires an environment, can fail with a typed error, and supports backpressure. The type of the environment dependency. The type of the failure value. The type of the success values in the stream.
Axial.ClockEnvironmentAn environment containing only a clock, for timed flows with no other services.
Axial.NeverRepresents an error channel that cannot occur.
intAn 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.
listThe 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.
(|>): '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.FlowStreamModulethrottle: TimeSpan -> FlowStream<'env,'error,'value> -> FlowStream<'env,'error,'value>Emits at most one value per , keeping the latest value when faster. The first value is emitted immediately. Values that arrive within of the last emission replace each other, and the latest is emitted when the interval ends. The pending value is emitted when upstream ends, and before a failure is propagated. Use it for progress reporting, where only the current state matters. progress |> FlowStream.throttle (TimeSpan.FromMilliseconds 100.0)
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 .
runCollect: FlowStream<'a,'b,'c> -> Flow<'a,'b,'c list>Collects all emitted values into a list. stream |> FlowStream.runCollect
A hundred progress updates in a burst render as the first and the last:
rendered (FlowStream.fromSeq [ 1..100 ]) |> Flow.run (ClockEnvironment Clock.live) |> shouldEqual (Exit.Success [ 1; 100 ])
rendered: FlowStream<ClockEnvironment,Never,int> -> Flow<ClockEnvironment,Never,int list>Axial.FlowStreamModulefromSeq: '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 ()
(..): ^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.
Wait for input to settle
FlowStream.debounce quiet emits a value only once quiet passes without a newer one, so a burst produces only its
last value. The pending value is emitted when upstream ends.
Keep only the latest request
FlowStream.switchMapFlow mapper runs mapper for each value, and interrupts the running flow when a newer value
arrives. Results for stale input are never emitted, and the superseded request's cleanup finishes before the next
one starts:
let search (query: string) : Flow<ClockEnvironment, Never, string> =
Flow.sleep (TimeSpan.FromMilliseconds 5.0) |> Flow.map (fun () -> $"results for {query}")
let shownResults (keystrokes: FlowStream<ClockEnvironment, Never, string>) : Flow<ClockEnvironment, Never, string list> =
keystrokes
|> FlowStream.debounce (TimeSpan.FromMilliseconds 200.0)
|> FlowStream.switchMapFlow search
|> FlowStream.runCollect
search: string -> Flow<ClockEnvironment,Never,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.NeverRepresents an error channel that cannot occur.
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)
shownResults: FlowStream<ClockEnvironment,Never,string> -> Flow<ClockEnvironment,Never,string list>keystrokes: FlowStream<ClockEnvironment,Never,string>Axial.FlowStream`3Represents a cold stream of values that requires an environment, can fail with a typed error, and supports backpressure. The type of the environment dependency. The type of the failure value. The type of the success values in the stream.
listThe 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.
Axial.FlowStreamModuledebounce: TimeSpan -> FlowStream<'env,'error,'value> -> FlowStream<'env,'error,'value>Emits a value only once passes without a newer one. Each value replaces the pending one and restarts the wait, so a burst produces only its last value. The pending value is emitted when upstream ends, and before a failure is propagated. Use it for search-as-you-type input. keystrokes |> FlowStream.debounce (TimeSpan.FromMilliseconds 300.0)
switchMapFlow: ('value -> Flow<'env,'error,'next>) -> FlowStream<'env,'error,'value> -> FlowStream<'env,'error,'next>Maps each value to a flow, interrupting the previous flow when a newer value arrives. Only the latest value's flow is kept: when upstream produces a new value while a flow is running, that flow is interrupted (and its cleanup awaited) and the new value's flow starts. Results are emitted as flows complete. When upstream ends, the running flow is allowed to finish. The first failure stops the stream. Use it for search-as-you-type and autocomplete, where a result for stale input is worthless. queries |> FlowStream.debounce (TimeSpan.FromMilliseconds 200.0) |> FlowStream.switchMapFlow search
runCollect: FlowStream<'a,'b,'c> -> Flow<'a,'b,'c list>Collects all emitted values into a list. stream |> FlowStream.runCollect
Typing "a", "ax", "axi" quickly searches only once, for the settled input:
shownResults (FlowStream.fromSeq [ "a"; "ax"; "axi" ]) |> Flow.run (ClockEnvironment Clock.live) |> shouldEqual (Exit.Success [ "results for axi" ])
shownResults: FlowStream<ClockEnvironment,Never,string> -> Flow<ClockEnvironment,Never,string list>Axial.FlowStreamModulefromSeq: '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 ()
(|>): '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.
When upstream ends, the running flow is allowed to finish. The first failure stops the stream.
These operators wait on real timers, so tests that use them depend on timing; give assertions generous margins.

