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.JavaScriptParallel
This scenario checks the parallel combinators with a random bound and random work durations.
Run it with dotnet run --project examples/Axial.TortureTest -- parallel 100.
What it does
Flow.traverseParandFlow.forEachParrun hundreds of items with a randomParallelismbound, counting how many run at once.Flow.traverseParUsingandFlow.forEachParUsinggive each running item a resource from a pool and record whether two items ever held the same one.- One random item of a traversal fails.
Flow.zipParruns a failing side against one that never ends, and two failing sides.Flow.sequenceParruns items of random length.- Fifty
Flow.races run two sides of random length, each signalling aDeferredwhen it ends.
let run (round: Round) : Flow<Axial.ClockEnvironment, Never, Check list> =
let items = [ 1 .. round.Size 200 ]
let bound = round.Next 6 + 1
flow {
// Bounded traversal: never more than the bound at once, results in input order.
let active = ref 0
let peak = ref 0
let tracked 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 * 2
}
|> Flow.ensuring (Flow.delay (fun () -> decrement active |> ignore; Flow.ok ()))
let! doubled = items |> Flow.traversePar (Parallelism.bounded bound) tracked
let processed = ResizeArray<int>()
do! items |> Flow.forEachPar (Parallelism.ofProcessors id) (fun value -> Flow.delay (fun () -> lock processed (fun () -> processed.Add value); Flow.ok ()))
// Pooled resources: each is used by one item at a time, and every one acquired is released.
let acquired = ref 0
let released = ref 0
let busy = Collections.Generic.HashSet<int>()
let shared = ref false
let pool =
Resource.ofAsync
(Flow.delay (fun () -> Flow.ok (increment acquired)))
(fun _ _ -> async { increment released |> ignore })
let useResource 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
}
let! pooled = items |> Flow.traverseParUsing (Parallelism.bounded bound) pool useResource
do! items |> Flow.forEachParUsing (Parallelism.bounded bound) pool (fun resource value -> useResource resource value |> Flow.ignore)
// A failing item stops the traversal; every item that started still cleaned up.
let started = ref 0
let cleaned = ref 0
let failAt = items.[round.Next items.Length]
let! failed =
items
|> Flow.traversePar (Parallelism.bounded bound) (fun value ->
flow {
increment started |> ignore
do! Flow.sleep (ms (float (round.Next 2)))
if value = failAt then return! Flow.fail value
return value
}
|> Flow.ensuring (Flow.delay (fun () -> increment cleaned |> ignore; Flow.ok ())))
|> exitOf
// zipPar fails fast and interrupts the other side; two failures are both reported.
let otherInterrupted = ref false
let! failFast = Flow.zipPar (Flow.fail "left") (Flow.never |> Flow.onInterrupt (Flow.delay (fun () -> otherInterrupted.Value <- true; Flow.ok ()))) |> exitOf
let! bothFail = Flow.zipPar (Flow.sleep (ms 1.0) |> Flow.bind (fun () -> Flow.fail "a")) (Flow.fail "b") |> exitOf
let! inOrder = [ for value in 1..20 -> Flow.sleep (ms (float (round.Next 3))) |> Flow.map (fun () -> value) ] |> Flow.sequencePar
// Races: the winner's value comes back, and both sides always finish: the winner completes, and the loser
// completes or is interrupted, running its exit handler either way.
let races = round.Size 50
let! finishes = [ for _ in 1 .. 2 * races -> Deferred.make<Axial.ClockEnvironment, Never, unit> () ] |> Flow.sequence
let! winners =
finishes
|> List.chunkBySize 2
|> Flow.traverse (fun pair ->
let side (name: string) (finished: Deferred<Never, unit>) =
Flow.sleep (ms (float (round.Next 3)))
|> Flow.map (fun () -> name)
|> Flow.onExit (fun _ -> Deferred.succeed () finished |> Flow.ignore)
Flow.race (side "left" pair[0]) (side "right" pair[1]))
let! allSidesFinished =
finishes
|> Flow.traverse Deferred.await
|> Flow.map (fun _ -> true)
|> Flow.timeoutToOk (TimeSpan.FromSeconds 10.0) false
let failures =
match bothFail with
| Exit.Failure cause -> Cause.failures cause |> List.sort
| Exit.Success _ -> []
return
[ check "a bounded traversal never ran more than its bound, and kept input order" (peak.Value <= bound && doubled = List.map ((*) 2) items)
check "forEachPar processed every item exactly once" (List.sort (List.ofSeq processed) = items)
check "no pooled resource was used by two items at once, and every one was released" (not shared.Value && pooled = items && acquired.Value = released.Value && acquired.Value <= 2 * bound)
check "a failing item failed the traversal, and every started item cleaned up" (failed = Exit.Failure(Cause.Fail failAt) && started.Value = cleaned.Value)
check "zipPar failed fast and interrupted the other side" (failFast = Exit.Failure(Cause.Fail "left") && otherInterrupted.Value)
check "zipPar reported the failures that happened" (not failures.IsEmpty && failures |> List.forall (fun failure -> failure = "a" || failure = "b"))
check "sequencePar kept input order" (inOrder = [ 1..20 ])
check "every race returned a side's value" (winners |> List.forall (fun winner -> winner = "left" || winner = "right"))
check "both sides of every race finished and ran their exit handler" (allSidesFinished && winners.Length = races) ]
}What each check proves
| Check | Guarantee |
|---|---|
| A bounded traversal never ran more than its bound, and kept input order | The bound is respected, and results come back in input order whatever order they finish in. |
| forEachPar processed every item exactly once | No item is skipped or repeated. |
| No pooled resource was used by two items at once, and every one was released | A resource is held by one item at a time, and the pool releases everything it acquired. |
| A failing item failed the traversal, and every started item cleaned up | Fail-fast interrupts the other items, and their cleanup still runs. |
| zipPar failed fast and interrupted the other side | A failure does not wait for the other side to finish on its own. |
| zipPar reported the failures that happened | Failures from both sides are kept in the cause. |
| sequencePar kept input order | Results are in input order. |
| Every race returned a side's value | A race returns its winner's value. |
| Both sides of every race finished and ran their exit handler | The loser is interrupted and settles; a race never leaves it running. |

