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

Queues

Four producers and three consumers share one queue with a back-pressure buffer of 8. Every take and offer can be interrupted at the moment a value changes hands, and the queue is shut down while consumers are still waiting.

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

What it does

  • Producers alternate Queue.offer and Queue.offerAll batches. The buffer is full most of the time, so offers wait.
  • Consumers pick a taking style at random for each take: Dequeue.take, takeBetween, poll, or takeUpTo. They stop when the queue is shut down and drained.
  • A saboteur forks takes and interrupts each one at once. When a value is handed over at the same instant as the interruption, the take either keeps it and reports success, or gives it back to the front of the queue.
  • A withdrawer forks offers and interrupts each one at once. An offer interrupted before it was accepted must not enqueue its value.
  • A second part fills a dropping queue and a sliding queue from three producers while one slow reader takes, and checks that the counters account for every value.
let run (round: Round) : Flow<Axial.ClockEnvironment, Never, Check list> =
    let producers = 4
    let perProducer = round.Size 300
    let consumers = 3

    flow {
        let consumed = ResizeArray<int * int>()
        let seenBy = Array.init consumers (fun _ -> ResizeArray<int * int>())
        let stolen = ResizeArray<int * int>()
        let offered = ResizeArray<int * int>()
        let withdrawn = ResizeArray<int * int>()
        let! (queue: Queue<int * int>) = Queue.make (QueueStrategy.BackPressure 8)

        // Producers alternate single offers and batches. Every value they offer is accepted: the queue is lossless.
        let producer id =
            flow {
                let mutable next = 1

                while next <= perProducer do
                    if round.Chance 30 then
                        let batch = [ for sample in next .. min perProducer (next + 4) -> id, sample ]
                        let! _ = queue |> Queue.offerAll batch
                        lock offered (fun () -> offered.AddRange batch)
                        next <- next + batch.Length
                    else
                        let! _ = queue |> Queue.offer (id, next)
                        lock offered (fun () -> offered.Add(id, next))
                        next <- next + 1
            }

        // Consumers mix every way of taking, until the queue is shut down and drained.
        let consumer index =
            flow {
                let mutable running = true

                while running do
                    let! batch =
                        match round.Next 4 with
                        | 0 -> queue |> Dequeue.take |> Flow.map List.singleton |> exitOf
                        | 1 -> queue |> Dequeue.takeBetween 1 6 |> exitOf
                        | 2 -> queue |> Dequeue.poll |> Flow.map Option.toList |> exitOf
                        | _ -> queue |> Dequeue.takeUpTo 3 |> exitOf

                    match batch with
                    | Exit.Success values ->
                        lock consumed (fun () -> consumed.AddRange values)
                        seenBy[index].AddRange values
                    | Exit.Failure _ -> running <- false
            }

        // Takes interrupted as they start: a value handed over in the same instant is either kept by the take
        // (the interruption reports success) or given back to the front of the queue.
        let saboteur =
            flow {
                for _ in 1 .. round.Size 200 do
                    let! taker = queue |> Dequeue.take |> Flow.fork
                    let! exit = Fiber.interrupt taker

                    match exit with
                    | Exit.Success value -> lock stolen (fun () -> stolen.Add value)
                    | Exit.Failure _ -> ()
            }

        // Offers interrupted as they start: an offer interrupted before it was accepted must never enqueue.
        let withdrawer =
            flow {
                for sample in 1 .. round.Size 100 do
                    let! offer = queue |> Queue.offer (0, sample) |> Flow.fork
                    let! exit = Fiber.interrupt offer

                    match exit with
                    | Exit.Success _ -> lock offered (fun () -> offered.Add(0, sample))
                    | Exit.Failure _ -> lock withdrawn (fun () -> withdrawn.Add(0, sample))
            }

        let! consumerFibers = [ 0 .. consumers - 1 ] |> Flow.traverse (consumer >> Flow.fork)
        do! [ for id in 1..producers -> producer id ] @ [ saboteur; withdrawer ] |> Flow.sequencePar |> Flow.ignore
        let! capacity = Flow.ok (Dequeue.capacity queue)
        do! Dequeue.shutdown queue
        do! Dequeue.awaitShutdown queue
        do! consumerFibers |> Flow.traverse Fiber.join |> Flow.ignore
        let! leftover = Dequeue.takeAll queue
        let! stats = Dequeue.stats queue
        let! isShut = Dequeue.isShutdown queue
        let! size = Dequeue.size queue

        let everything = List.ofSeq consumed @ List.ofSeq stolen @ leftover

        return
            [ check "every accepted value is taken exactly once" (List.sort everything = List.sort (List.ofSeq offered))
              check "no withdrawn offer was ever enqueued" (withdrawn |> Seq.forall (fun value -> not (List.contains value everything)))
              check
                  "each consumer saw every producer's values in order"
                  (seenBy |> Array.forall (List.ofSeq >> List.filter (fst >> (<>) 0) >> inOrderPerProducer))
              check "the queue counted every accepted value" (stats.Accepted = int64 offered.Count)
              check "no taker or offerer is left waiting" (stats.WaitingTakers = 0 && stats.WaitingOfferers = 0)
              check "the queue is shut down and empty" (isShut && size = 0 && capacity = Some 8) ]
    }

/// Lossy strategies under concurrent producers: a dropping queue and a sliding queue each read by a slow consumer.
let runLossy (round: Round) : Flow<Axial.ClockEnvironment, Never, Check list> =
    let perProducer = round.Size 500

    let exercise (make: Flow<Axial.ClockEnvironment, Never, Queue<int * int>>) =
        flow {
            let! queue = make
            let seen = ResizeArray<int * int>()
            let accepted = ref 0

            let! reader =
                flow {
                    let mutable running = true

                    while running do
                        let! next = queue |> Dequeue.take |> exitOf

                        match next with
                        | Exit.Success value -> seen.Add value
                        | Exit.Failure _ -> running <- false
                }
                |> Flow.fork

            do!
                [ for id in 1..3 ->
                      flow {
                          for sample in 1..perProducer do
                              let! wasAccepted = queue |> Queue.offer (id, sample)
                              if wasAccepted then lock accepted (fun () -> accepted.Value <- accepted.Value + 1)
                      } ]
                |> Flow.sequencePar
                |> Flow.ignore

            do! Dequeue.shutdown queue
            do! Fiber.join reader
            let! stats = Dequeue.stats queue
            return List.ofSeq seen, accepted.Value, stats
        }

    flow {
        let! dropSeen, dropAccepted, dropStats = exercise (Queue.dropping 4)
        let! slideSeen, slideAccepted, slideStats = exercise (Queue.sliding 4)
        let offered = 3 * perProducer

        // A queue created under a scope is shut down when the scope closes.
        let! scoped = Queue.makeScoped (QueueStrategy.Unbounded) |> Flow.scoped
        let! scopedShut = Dequeue.isShutdown scoped
        let! (unbounded: Queue<int>) = Queue.unbounded ()
        do! unbounded |> Queue.offerAll [ 1..1000 ] |> Flow.ignore
        let! unboundedAll = Dequeue.takeAll unbounded
        let! (callbackInbox: Queue<int>) = Queue.bounded 1
        let first = callbackInbox |> Queue.tryOffer 7
        let full = callbackInbox |> Queue.tryOffer 8
        let! accepted = Dequeue.take callbackInbox
        do! Dequeue.shutdown callbackInbox
        let shut = callbackInbox |> Queue.tryOffer 9

        return
            [ check "a dropping queue accounts for every offer as accepted or dropped" (int64 dropAccepted + dropStats.Dropped = int64 offered)
              check "a dropping queue delivers what it accepted, in each producer's order" (dropSeen.Length = dropAccepted && inOrderPerProducer dropSeen)
              check "a sliding queue accepts every offer" (slideAccepted = offered)
              check "a sliding queue delivers or evicts every value, in each producer's order" (int64 slideSeen.Length + slideStats.Evicted = int64 offered && inOrderPerProducer slideSeen)
              check "a scoped queue is shut down with its scope" scopedShut
              check "a synchronous callback sees acceptance, full, and shutdown without losing the accepted value"
                  (first = QueueTryOfferResult.Accepted && full = QueueTryOfferResult.Full && accepted = 7
                   && shut = QueueTryOfferResult.Shutdown)
              check "an unbounded queue never refuses" (unboundedAll = [ 1..1000 ]) ]
    }

What each check proves

Check Guarantee
Every accepted value is taken exactly once An interrupted take never loses a value and never duplicates one: a value handed over in the same instant as the interruption is kept or given back.
No withdrawn offer was ever enqueued An offer interrupted while it waits for space is withdrawn completely.
Each consumer saw every producer's values in order Waiting takers are served one batch at a time, so a value given back by an interrupted take returns ahead of later values.
The queue counted every accepted value Dequeue.stats agrees with what the producers were told.
No taker or offerer is left waiting Shutdown releases every suspended take and offer.
The queue is shut down and empty Dequeue.shutdown and Dequeue.awaitShutdown end the queue, and draining with takeAll leaves nothing behind.
A dropping queue accounts for every offer Every offer is either accepted or counted as dropped, and a dropped offer reports false.
A sliding queue accepts every offer Evictions make room, and every value is delivered or counted as evicted.
A scoped queue is shut down with its scope Queue.makeScoped ties the queue to the enclosing scope.