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
This scenario covers the fiber tools that decide which work keeps running, what forked work inherits, and how a running program reports its fibers.
Run it with dotnet run --project examples/Axial.TortureTest -- fibers 100.
What it does
- Sixty searches are forked into one
FiberSlotwithFlow.forkReplacing, so each replaces the one before it, as when a user keeps typing. The same runs per key withFiberSlot.makeKeyedandFlow.forkReplacingKey. - Ten workers are forked under
Flow.annotate,Flow.withTraceId, annotation sinks, andFlow.localEnv, and each reports what it inherited. - Forty named fibers that succeed, fail, or sleep run under a
FiberRegistryand aFiberObserver. The sleepers are interrupted through the registry by id and by name. A discarded fiber and a detached fiber both die.
let run (round: Round) : Flow<Axial.ClockEnvironment, Never, Check list> =
let requests = round.Size 60
flow {
// Latest wins: each request replaces the one before it. Only the last is sure to finish; every other one
// either finished first or was interrupted, and every one of them cleaned up.
let! cleanedUp = Ref.make 0
let! slot = FiberSlot.make ()
let search query =
flow {
do! Flow.sleep (TimeSpan.FromMilliseconds(float (round.Next 3)))
return query
}
|> Flow.ensuring (cleanedUp |> Ref.update ((+) 1))
let! searches = [ 1..requests ] |> Flow.traverse (fun query -> slot |> Flow.forkReplacing <| search query)
let! lastResult = searches |> List.last |> Fiber.join
do! FiberSlot.interrupt slot
let! searchExits = searches |> Flow.traverse Fiber.await
let! polled = searches |> Flow.traverse Fiber.poll
let! cleaned = Ref.get cleanedUp
// Keyed slots: the same rule per key, and interruptAll leaves nothing running.
let! keyed = FiberSlot.makeKeyed ()
let! previews =
[ for index in 1..requests -> index % 5, index ]
|> Flow.traverse (fun (key, index) -> keyed |> Flow.forkReplacingKey key <| search index)
let lastPerKey = [ for key in 0..4 -> previews |> List.indexed |> List.filter (fun (index, _) -> (index + 1) % 5 = key) |> List.last |> snd ]
let! lastPreviews = lastPerKey |> Flow.traverse Fiber.await
do! FiberSlot.interruptAll keyed
let! remaining = FiberSlot.count keyed
let! previewPolls = previews |> Flow.traverse Fiber.poll
// Annotations, the trace id, and the environment are inherited by forked fibers, and siblings do not see
// each other's annotations.
// Observer and sink callbacks are synchronous, so they count with plain atomic counters.
let sinkCount = ref 0
let! inherited =
[ for worker in 1..10 ->
flow {
let! annotations = Flow.annotations
let! trace = Flow.traceId
let! id = Flow.fiberId
let! tenant = Flow.envWith _.Tenant
return worker, annotations, trace, id, tenant
}
|> Flow.annotate "worker" (string worker)
|> Flow.fork ]
|> Flow.sequence
|> Flow.bind (Flow.traverse Fiber.join)
|> Flow.annotate "request" "r1"
|> Flow.withTraceId "trace-1"
// addAnnotationSink composes with the sink installed outside it; withAnnotationSink replaces it.
|> Flow.addAnnotationSink (fun _ _ -> increment sinkCount |> ignore)
|> Flow.withAnnotationSink (fun _ _ -> ())
|> Flow.localEnv (fun (env: ClockEnvironment) -> { Tenant = "acme"; Clock = env.Clock })
let annotationsHeld =
inherited
|> List.forall (fun (worker, annotations, trace, _, tenant) ->
annotations.TryFind "request" = Some "r1"
&& annotations.TryFind "worker" = Some(string worker)
&& trace = Some "trace-1"
&& tenant = "acme")
let distinctIds = inherited |> List.map (fun (_, _, _, id, _) -> id) |> List.distinct |> List.length
let annotated = sinkCount.Value
return
[ check "the last request's result is delivered" (lastResult = requests)
check "every replaced request finished or was interrupted, and none is left running" (polled |> List.forall settledOk)
check "every request cleaned up" (cleaned = requests && searchExits.Length = requests)
check "the last preview per key finished, and interruptAll left none running" (lastPreviews |> List.forall (function Exit.Success _ -> true | _ -> false) && remaining = 0 && previewPolls |> List.forall settledOk)
check "forked fibers inherit annotations, the trace id, and the environment, but not their siblings' annotations" annotationsHeld
check "every fiber has its own id" (distinctIds = 10)
check "annotation sinks saw the annotations" (annotated >= 10) ]
}
/// The fiber registry and observers under a burst of named fibers that succeed, fail, and are interrupted.
let diagnostics (round: Round) : Flow<Axial.ClockEnvironment, Never, Check list> =
let workers = round.Size 40
flow {
let! clock = Flow.envWith (fun (env: ClockEnvironment) -> env.Clock)
let registry = FiberRegistry(10)
let starts = ref 0
let ends = ref 0
let startDumps = ResizeArray<FiberDump>()
let counting =
{ FiberObserver.none with
OnStart =
fun metadata ->
increment starts |> ignore
lock startDumps (fun () -> startDumps.Add(FiberDump.ofMetadata metadata))
OnEnd = fun _ _ -> increment ends |> ignore }
let body =
flow {
let! fibers =
[ for index in 1..workers ->
match index % 3 with
| 0 -> Flow.sleep (TimeSpan.FromMinutes 5.0) |> Flow.forkNamed "sleeper" |> Flow.map Choice1Of2
| 1 -> (Flow.fail "expected" : Flow<Axial.ClockEnvironment, string, unit>) |> Flow.forkNamed "failer" |> Flow.map Choice2Of2
| _ -> Flow.ok () |> Flow.forkNamed "worker" |> Flow.map Choice1Of2 ]
|> Flow.sequence
let sleepers = fibers |> List.choose (function Choice1Of2 fiber when (Fiber.dump fiber).Name = Some "sleeper" -> Some fiber | _ -> None)
let firstSleeper = List.head sleepers
let interruptedOne = registry.Interrupt (Fiber.dump firstSleeper).Id
let interruptedByName = registry.InterruptByName "sleeper"
let! exits = fibers |> Flow.traverse (function Choice1Of2 fiber -> Fiber.await fiber |> Flow.map ignore | Choice2Of2 fiber -> Fiber.await fiber |> Flow.map ignore)
// A discarded fiber that dies is reported when its scope closes; a detached one is not.
do!
flow {
let! _ = Flow.die (InvalidOperationException "lost") |> Flow.fork
let! _ = Flow.die (InvalidOperationException "deliberate") |> Flow.forkDetached
do! Flow.sleep (TimeSpan.FromMilliseconds 5.0)
}
|> Flow.scoped
return interruptedOne, interruptedByName, exits.Length
}
let! interruptedOne, interruptedByName, settled =
body
|> Flow.withFiberRegistry registry
|> Flow.withFiberObserver counting
|> Flow.addFiberObserver (FiberObserver.compose FiberObserver.none FiberObserver.none)
let stats = registry.Stats() |> List.map (fun stats -> stats.Name, stats) |> Map.ofList
let sleepers = workers / 3
let started, ended = starts.Value, ends.Value
let dump = registry.DumpAt clock
let observedDumps = lock startDumps (fun () -> List.ofSeq startDumps)
let tree = FiberDump.renderTreeAt clock observedDumps
let lines = observedDumps |> List.map (FiberDump.renderAt clock)
return
// The fiber interrupted by id may not have settled yet when InterruptByName runs, so it can be counted twice.
[ check "interrupting by id and by name reached every sleeper" (interruptedOne && interruptedByName >= sleepers - 1)
check "the registry counted every named fiber and how each ended" (stats["sleeper"].Interrupted = sleepers && stats["failer"].Failed = stats["failer"].Count && settled = workers)
check "nothing is live afterwards, and history stays within its capacity" (registry.LiveFiberCount = 0 && registry.Snapshot().IsEmpty && registry.Settled().Length <= registry.HistoryCapacity)
check "the registry saw every fork, and dumps render" (registry.StartedCount >= int64 workers && not (isNull dump))
check "the observer saw a dump of every start, and each renders" (observedDumps.Length = started && lines |> List.forall (fun line -> tree.Contains line) && lines |> List.exists (fun line -> line.Contains "\"sleeper\""))
check "a discarded fiber's defect was reported, a detached one's was not" (registry.UnobservedDefects() |> List.map _.Defect |> List.exists (fun text -> text.Contains "lost") && not (registry.UnobservedDefects() |> List.exists (fun defect -> defect.Defect.Contains "deliberate")))
check "the observer saw every fiber start and end" (started = ended && started >= workers) ]
}What each check proves
| Check | Guarantee |
|---|---|
| The last request's result is delivered | Replacing never interrupts the newest fiber. |
| Every replaced request finished or was interrupted, and none is left running | A slot holds one fiber, and interrupting the slot stops it. |
| Every request cleaned up | A replaced fiber runs its Flow.ensuring cleanup. |
| The last preview per key finished, and interruptAll left none running | Keyed slots apply the same rule per key. |
| Forked fibers inherit annotations, the trace id, and the environment, but not their siblings' annotations | A fork copies its parent's context; a sibling's annotation stays with that sibling. |
| Every fiber has its own id | Flow.fiberId identifies each fork. |
| Annotation sinks saw the annotations | Flow.addAnnotationSink composes with the sink outside it. |
| Interrupting by id and by name reached every sleeper | FiberRegistry.Interrupt and InterruptByName find live fibers. |
| The registry counted every named fiber and how each ended | FiberRegistry.Stats groups outcomes by fiber name. |
| Nothing is live afterwards, and history stays within its capacity | Settled fibers leave the live set, and the settled history is bounded. |
| The observer saw a dump of every start, and each renders | FiberDump renders what the observer saw. |
| A discarded fiber's defect was reported, a detached one's was not | A defect nobody observed is reported when its scope closes, unless the fiber was forked with Flow.forkDetached. |
| The observer saw every fiber start and end | Every start has a matching end. |

