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

Semaphores

Sixty workers compete for the three permits of a semaphore. A random third are interrupted, some while they wait for a permit and some while they hold one.

Run it with dotnet run --project examples/Axial.TortureTest -- semaphores 100.

What it does

  • Each worker takes a permit with Semaphore.withPermit, counts itself in, sleeps briefly, and counts itself out with Flow.ensuring, so it is counted out however it ends.
  • After the workers have settled, three flows each take a permit and wait on a Deferred until all three hold one at the same time. This only completes if every permit came back.
let semaphores (round: Round) : Flow<Axial.ClockEnvironment, Never, Check list> =
    let permits = 3
    let workers = round.Size 60

    flow {
        let! semaphore = Semaphore.make permits
        let! active = Ref.make 0
        let! busiest = Ref.make 0
        let entered = ResizeArray<int>()

        // Each worker holds a permit while it counts itself in; ensuring counts it out however the worker ends.
        let worker id =
            flow {
                let! now = active |> Ref.updateAndGet ((+) 1)
                do! busiest |> Ref.update (max now)
                lock entered (fun () -> entered.Add id)
                do! Flow.sleep (TimeSpan.FromMilliseconds(float (round.Next 2)))
            }
            |> Flow.ensuring (active |> Ref.update (fun count -> count - 1))
            |> Semaphore.withPermit semaphore

        let! fibers = [ 1..workers ] |> Flow.traverse (worker >> Flow.fork)

        // Interrupt a random third of the workers, whether they are waiting for a permit or holding one.
        let! exits =
            fibers
            |> Flow.traverse (fun fiber -> if round.Chance 33 then Fiber.interrupt fiber else Fiber.await fiber)

        let! finalActive = Ref.get active
        let! mostAtOnce = Ref.get busiest

        // If an interrupted waiter had lost a permit, these three could not all hold one at the same time.
        let! barrier = Deferred.make<Axial.ClockEnvironment, Never, unit> ()
        let! arrived = Ref.make 0

        let holdTogether =
            flow {
                let! count = arrived |> Ref.updateAndGet ((+) 1)
                if count = permits then do! Deferred.succeed () barrier |> Flow.ignore
                do! Deferred.await barrier
            }
            |> Semaphore.withPermit semaphore

        let! allPermitsFree =
            List.replicate permits holdTogether
            |> Flow.sequencePar
            |> Flow.map (fun _ -> true)
            |> Flow.timeoutToOk (TimeSpan.FromSeconds 10.0) false

        let succeeded = exits |> List.filter (function Exit.Success _ -> true | _ -> false) |> List.length

        return
            [ check "no more workers than permits ever held one at once" (mostAtOnce <= permits && mostAtOnce >= 1)
              check "every worker that finished had entered" (succeeded <= entered.Count)
              check "every worker that entered was counted out, even if interrupted inside" (finalActive = 0)
              check "interrupted waiters returned their permits: all permits can be held together" allPermitsFree ]
    }

What each check proves

Check Guarantee
No more workers than permits ever held one at once The semaphore bounds concurrency under contention.
Every worker that finished had entered A worker completes only after taking a permit.
Every worker that entered was counted out, even if interrupted inside Flow.ensuring runs its cleanup on interruption.
Interrupted waiters returned their permits A waiter interrupted as its permit is handed over passes the permit on instead of losing it.