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.JavaScript

Fibers

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

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.