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.JavaScriptRefs
Eight fibers update one Ref and one SubscriptionRef concurrently, using every update operation.
Run it with dotnet run --project examples/Axial.TortureTest -- refs 100.
What it does
- Each fiber adds one to a counter hundreds of times, choosing
Ref.update,modify,getAndUpdate, orupdateAndGetat random, and records the value each operation replaced. - The fibers then pass tokens through a second ref with
Ref.getAndSet: each call puts a new token in and takes the previous one out. - A
SubscriptionRefis incremented the same way withupdateandmodify, while achangesstream that subscribed before the first update collects every value.
let refs (round: Round) : Flow<Axial.ClockEnvironment, Never, Check list> =
let fibers = 8
let perFiber = round.Size 250
flow {
let! counter = Ref.make 0
// Every operation adds one. Those that report the value they replaced must each see a different one.
let increments =
flow {
let observed = ResizeArray<int>()
for _ in 1..perFiber do
match round.Next 4 with
| 0 -> do! counter |> Ref.update ((+) 1)
| 1 ->
let! before = counter |> Ref.modify (fun value -> value, value + 1)
observed.Add before
| 2 ->
let! before = counter |> Ref.getAndUpdate ((+) 1)
observed.Add before
| _ ->
let! after = counter |> Ref.updateAndGet ((+) 1)
observed.Add(after - 1)
return List.ofSeq observed
}
let! observed = List.replicate fibers increments |> Flow.sequencePar |> Flow.map List.concat
let! total = Ref.get counter
// Tokens passed around with getAndSet: every token is handed on exactly once.
let! slot = Ref.make 0
let! returned =
[ for fiber in 1..fibers -> [ for index in 1..perFiber -> slot |> Ref.getAndSet (fiber * 100000 + index) ] |> Flow.sequence ]
|> Flow.sequencePar
|> Flow.map List.concat
let! last = Ref.get slot
do! slot |> Ref.set -1
let! reset = Ref.get slot
let placed = [ for fiber in 1..fibers do for index in 1..perFiber -> fiber * 100000 + index ]
// A SubscriptionRef updated concurrently: a stream that started first sees every value in turn.
// The stream subscribes when it starts; waiting for its first value means it is subscribed before any update.
let! level = SubscriptionRef.make 0
let! subscribed = Deferred.make<Axial.ClockEnvironment, Never, unit> ()
let! history =
level
|> SubscriptionRef.changes QueueStrategy.Unbounded
|> FlowStream.tapFlow (fun _ -> Deferred.succeed () subscribed |> Flow.ignore)
|> FlowStream.take (fibers * perFiber + 1)
|> FlowStream.runCollect
|> Flow.fork
do! Deferred.await subscribed
do!
List.replicate fibers (flow {
for index in 1..perFiber do
if index % 2 = 0 then do! level |> SubscriptionRef.update ((+) 1)
else do! level |> SubscriptionRef.modify (fun value -> (), value + 1)
})
|> Flow.sequencePar
|> Flow.ignore
let! finalLevel = SubscriptionRef.get level
do! level |> SubscriptionRef.set 0
let! seen = Fiber.join history
return
[ check "no increment was lost" (total = fibers * perFiber)
check "every increment that reported a value saw a different one" (List.length (List.distinct observed) = observed.Length)
check "getAndSet handed every token on exactly once" (List.sort (last :: returned) = List.sort (0 :: placed) && reset = -1)
check "a SubscriptionRef ends at the total" (finalLevel = fibers * perFiber)
check "its change stream saw every value, in order" (seen = [ 0 .. fibers * perFiber ]) ]
}What each check proves
| Check | Guarantee |
|---|---|
| No increment was lost | Every update is atomic. |
| Every increment that reported a value saw a different one | No two updates observed the same state, so none of them overlapped. |
| getAndSet handed every token on exactly once | A swap never loses or duplicates the value it replaces. |
| A SubscriptionRef ends at the total | Its updates are atomic too. |
| Its change stream saw every value, in order | Each update is delivered to the stream before the next update starts, with no gap and no duplicate. |

