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.JavaScriptStreams
Each round feeds random input through every stream operator. Deterministic operators must agree
with the List functions they mirror. Concurrent and timed operators must keep the properties they promise, whatever
the timing. Streams stopped early or failing must release what they own.
Run it with dotnet run --project examples/Axial.TortureTest -- streams 100.
What it does
map,filter,choose,scan,indexed,distinctUntilChangedBy,skip,take,chunkBySize,collect,append,zip,mapFlow,unfoldFlow, and therun*consumers run on the same input as theirListequivalents.mapFlowParandmapFlowParUsingrun with a random bound.bufferruns with a lossless, a sliding, and a dropping strategy against a source that pauses at random.mergePar,groupedWithin,throttle,debounce, andswitchMapFlowrun on the same pausing source.- A stream that owns a resource through
FlowStream.usingis stopped early bytake, and fails part way.
let run (round: Round) : Flow<Axial.ClockEnvironment, Never, Check list> =
let input = [ for _ in 1 .. round.Size 300 -> round.Next 20 ]
let source () = FlowStream.fromSeq input
flow {
// Deterministic operators must agree with the List functions they mirror.
let! mapped = source () |> FlowStream.map ((*) 3) |> FlowStream.filter (fun value -> value % 2 = 0) |> FlowStream.runCollect
let! chosen = source () |> FlowStream.choose (fun value -> if value > 10 then Some(value - 10) else None) |> FlowStream.runCollect
let! scanned = source () |> FlowStream.scan (+) 0 |> FlowStream.runCollect
let! indexed = source () |> FlowStream.indexed |> FlowStream.runCollect
let! distinct = source () |> FlowStream.distinctUntilChangedBy (fun value -> value / 5) |> FlowStream.runCollect
let! skipped = source () |> FlowStream.skip 7 |> FlowStream.skipWhile (fun value -> value < 15) |> FlowStream.runCollect
let! taken = source () |> FlowStream.takeWhile (fun value -> value <> 19) |> FlowStream.take 50 |> FlowStream.runCollect
let! chunks = source () |> FlowStream.chunkBySize 7 |> FlowStream.runCollect
let! expanded = source () |> FlowStream.collect (fun value -> FlowStream.fromSeq [ value; -value ]) |> FlowStream.runCollect
let! appended = FlowStream.singleton -1 |> FlowStream.append (source ()) |> FlowStream.append FlowStream.empty |> FlowStream.runCollect
let! zipped = source () |> FlowStream.zip (FlowStream.fromSeq [ "a"; "b"; "c" ]) |> FlowStream.runCollect
let! flowMapped = source () |> FlowStream.mapFlow (fun value -> Flow.ok (value + 1)) |> FlowStream.runCollect
let! unfolded = FlowStream.unfoldFlow (fun state -> Flow.ok (if state < 10 then Some(state, state + 1) else None)) 0 |> FlowStream.runCollect
let! single = FlowStream.fromFlow (Flow.ok 42) |> FlowStream.runCollect
let! localized =
(FlowStream.fromSeq input : FlowStream<unit, Never, int>)
|> FlowStream.localEnv (fun (_: ClockEnvironment) -> ())
|> FlowStream.runCollect
let! folded = source () |> FlowStream.runFold (+) 0
let! counted = source () |> FlowStream.runCount
let! head = source () |> FlowStream.runTryHead
let! last = source () |> FlowStream.runTryLast
let! nothing = FlowStream.empty<Axial.ClockEnvironment, Never, int> |> FlowStream.runTryHead
let seen = ResizeArray<int>()
do! source () |> FlowStream.tapFlow (fun value -> Flow.delay (fun () -> seen.Add value; Flow.ok ())) |> FlowStream.runDrain
do! source () |> FlowStream.runForEach seen.Add
do! source () |> FlowStream.runForEachFlow (fun value -> Flow.delay (fun () -> seen.Add value; Flow.ok ()))
let modelHolds =
mapped = (input |> List.map ((*) 3) |> List.filter (fun value -> value % 2 = 0))
&& chosen = (input |> List.choose (fun value -> if value > 10 then Some(value - 10) else None))
&& scanned = (List.scan (+) 0 input |> List.tail)
&& indexed = List.indexed input
&& distinct = (input |> List.fold (fun kept value -> match kept with previous :: _ when previous / 5 = value / 5 -> kept | _ -> value :: kept) [] |> List.rev)
&& skipped = (input |> List.skip (min 7 input.Length) |> List.skipWhile (fun value -> value < 15))
&& taken = (input |> List.takeWhile (fun value -> value <> 19) |> List.truncate 50)
&& chunks = List.chunkBySize 7 input
&& expanded = (input |> List.collect (fun value -> [ value; -value ]))
&& appended = -1 :: input
&& zipped = List.zip (List.truncate 3 input) (List.truncate input.Length [ "a"; "b"; "c" ])
&& flowMapped = List.map ((+) 1) input
&& unfolded = [ 0..9 ]
&& single = [ 42 ]
&& localized = input
&& folded = List.sum input
&& counted = input.Length
&& head = List.tryHead input
&& last = List.tryLast input
&& nothing = None
&& List.ofSeq seen = input @ input @ input
// Concurrent and timed operators keep the properties they promise, whatever the timing.
let bound = round.Next 5 + 1
let active = ref 0
let peak = ref 0
let slowly value =
flow {
let now = increment active
lock peak (fun () -> peak.Value <- max peak.Value now)
do! Flow.sleep (ms (float (round.Next 2)))
return value
}
|> Flow.ensuring (Flow.delay (fun () -> decrement active |> ignore; Flow.ok ()))
let! parallelResults = source () |> FlowStream.mapFlowPar (Parallelism.bounded bound) slowly |> FlowStream.runCollect
let busy = Collections.Generic.HashSet<int>()
let shared = ref false
let acquired = ref 0
let released = ref 0
let pool =
Resource.ofAsync
(Flow.delay (fun () -> Flow.ok (increment acquired)))
(fun _ _ -> async { increment released |> ignore })
let! pooledResults =
source ()
|> FlowStream.mapFlowParUsing (Parallelism.bounded bound) pool (fun resource value ->
flow {
lock busy (fun () -> if not (busy.Add resource) then shared.Value <- true)
do! Flow.sleep (ms (float (round.Next 2)))
lock busy (fun () -> busy.Remove resource |> ignore)
return value
})
|> FlowStream.runCollect
let timed () = source () |> FlowStream.tapFlow (fun _ -> if round.Chance 5 then Flow.sleep (ms 3.0) else Flow.ok ())
let! lossless = timed () |> FlowStream.buffer (QueueStrategy.BackPressure 4) |> FlowStream.runCollect
let! sliding = timed () |> FlowStream.buffer (QueueStrategy.Sliding 2) |> FlowStream.runCollect
let! dropping = timed () |> FlowStream.buffer (QueueStrategy.Dropping 2) |> FlowStream.runCollect
let! merged = [ source () |> FlowStream.map (fun value -> 1, value); timed () |> FlowStream.map (fun value -> 2, value) ] |> FlowStream.mergePar |> FlowStream.runCollect
let! groups = timed () |> FlowStream.groupedWithin 5 (ms 2.0) |> FlowStream.runCollect
let! throttled = timed () |> FlowStream.indexed |> FlowStream.throttle (ms 1.0) |> FlowStream.runCollect
let! debounced = timed () |> FlowStream.indexed |> FlowStream.debounce (ms 1.0) |> FlowStream.runCollect
let! switched = timed () |> FlowStream.indexed |> FlowStream.switchMapFlow (fun (index, value) -> Flow.sleep (ms (float (round.Next 2))) |> Flow.map (fun () -> index, value)) |> FlowStream.runCollect
let indexedInput = List.indexed input
let byIndex (pairs: (int * int) list) = pairs |> List.map fst |> List.pairwise |> List.forall (fun (a, b) -> a < b)
let concurrent =
[ check "mapFlowPar kept every value and its bound" (List.sort parallelResults = List.sort input && peak.Value <= bound)
check "mapFlowParUsing never shared a resource and released every one" (List.sort pooledResults = List.sort input && not shared.Value && acquired.Value = released.Value)
check "a lossless buffer kept every value in order" (lossless = input)
check "a sliding buffer kept values in order, ending with the last" (isSubsequence sliding input && List.tryLast sliding = List.tryLast input)
check "a dropping buffer kept values in order" (isSubsequence dropping input)
check "mergePar kept every value and each source's order" (List.sort merged = List.sort ((input |> List.map (fun value -> 1, value)) @ (input |> List.map (fun value -> 2, value))) && inOrderPerProducer (merged |> List.indexed |> List.map (fun (position, (source, _)) -> source, position)))
check "groupedWithin grouped every value in order, within its size" (List.concat groups = input && groups |> List.forall (fun group -> not group.IsEmpty && group.Length <= 5))
check "throttle kept values in order, including the first and the last" (byIndex throttled && isSubsequence throttled indexedInput && List.tryHead throttled = List.tryHead indexedInput && List.tryLast throttled = List.tryLast indexedInput)
check "debounce kept values in order, including the last" (byIndex debounced && isSubsequence debounced indexedInput && List.tryLast debounced = List.tryLast indexedInput)
check "switchMapFlow kept results in order, including the last" (byIndex switched && isSubsequence switched indexedInput && List.tryLast switched = List.tryLast indexedInput) ]
// Stopping early, or failing, releases what the stream owns and stops the fibers it started.
let opened = ref 0
let closed = ref 0
let owned () =
FlowStream.using
(Resource.ofAsync (Flow.delay (fun () -> increment opened |> ignore; Flow.ok ())) (fun _ _ -> async { increment closed |> ignore }))
(fun () -> FlowStream.repeatFlow (Flow.ok 1))
let! early = owned () |> FlowStream.mapFlowPar (Parallelism.bounded 3) slowly |> FlowStream.take 5 |> FlowStream.runCollect
let! failed =
owned ()
|> FlowStream.mapFlow (fun _ -> if round.Chance 20 then Flow.fail "broken" else Flow.ok 1)
|> FlowStream.mapError (fun error -> error.Length)
|> FlowStream.runDrain
|> exitOf
let! runningAfter = Flow.sleep (ms 5.0) |> Flow.map (fun () -> active.Value)
return
[ check "deterministic operators agree with the List model" modelHolds ]
@ concurrent
@ [ check "stopping early or failing released every owned resource" (early = [ 1; 1; 1; 1; 1 ] && opened.Value = 2 && closed.Value = 2 && failed = Exit.Failure(Cause.Fail 6))
check "no stream fiber was left running" (runningAfter = 0) ]
}What each check proves
| Check | Guarantee |
|---|---|
| Deterministic operators agree with the List model | Each operator computes what its List counterpart does. |
| mapFlowPar kept every value and its bound | Parallel mapping never exceeds its bound and loses nothing. |
| mapFlowParUsing never shared a resource and released every one | Each mapping holds its own resource. |
| A lossless buffer kept every value in order | BackPressure slows the producer instead of dropping. |
| A sliding buffer kept values in order, ending with the last | Sliding loses older values but always delivers the newest. |
| A dropping buffer kept values in order | Dropping loses newer values but never reorders. |
| mergePar kept every value and each source's order | Merging interleaves sources without losing or reordering any one of them. |
| groupedWithin grouped every value in order, within its size | Groups close on size or time, and no value is lost at a boundary. |
| throttle kept values in order, including the first and the last | Throttling drops values in between, never the first or the last. |
| debounce kept values in order, including the last | The final value always gets through. |
| switchMapFlow kept results in order, including the last | The newest value's flow always finishes, even when it arrives as a timer fires. |
| Stopping early or failing released every owned resource | take and a failure both close the stream's scope. |
| No stream fiber was left running | Every fiber a stream starts is interrupted and awaited before the consumer returns. |

