diff options
Diffstat (limited to 'src/SomhairlesDream.Server/Store.fs')
| -rw-r--r-- | src/SomhairlesDream.Server/Store.fs | 259 |
1 files changed, 259 insertions, 0 deletions
diff --git a/src/SomhairlesDream.Server/Store.fs b/src/SomhairlesDream.Server/Store.fs new file mode 100644 index 0000000..b4a2293 --- /dev/null +++ b/src/SomhairlesDream.Server/Store.fs @@ -0,0 +1,259 @@ +namespace SomhairlesDream.Server + +open System +open System.Collections.Generic +open System.Threading.Channels +open SomhairlesDream.Shared + +type RunStore(projectId: string, runId: string, staleAfter: TimeSpan) = + let gate = obj () + let mutable state = RunState.initial projectId runId DateTimeOffset.UtcNow + let subscribers = ResizeArray<Channel<LiveSnapshot>>() + + let update transition = + lock gate (fun () -> + state <- transition state + for subscriber in subscribers do + subscriber.Writer.TryWrite(state) |> ignore + state) + + member _.Snapshot() = lock gate (fun () -> state) + + member _.Start(now: DateTimeOffset) = update (RunState.start now) + + member _.Heartbeat(now: DateTimeOffset) = update (RunState.heartbeat now) + + member _.Observe(now: DateTimeOffset) = update (RunState.statusAt staleAfter now) + + member _.Publish(now: DateTimeOffset, mesh: MeshSnapshot) = update (RunState.publish now mesh) + + member _.Subscribe() : ChannelReader<LiveSnapshot> = + let channel = Channel.CreateUnbounded<LiveSnapshot>() + + lock gate (fun () -> + channel.Writer.TryWrite(state) |> ignore + subscribers.Add(channel) + channel.Reader) + +type ArtifactRunSubscription(reader: ChannelReader<RunSnapshot>, dispose: unit -> unit) = + let mutable disposed = false + + member _.Reader = reader + + member _.Dispose() = + if not disposed then + disposed <- true + dispose () + + interface IDisposable with + member this.Dispose() = this.Dispose() + +type ArtifactRunStore(ids: RunIds, staleAfter: TimeSpan, initialAt: DateTimeOffset) = + let gate = obj () + let subscribers = ResizeArray<Channel<RunSnapshot>>() + let mutable state : RunSnapshot = + { ProjectId = ids.ProjectId + RunId = ids.RunId + BaseId = ids.BaseId + TargetId = ids.TargetId + Status = "idle" + CurrentStepId = None + CurrentStepIndex = -1 + TotalSteps = 0 + LastHeartbeatAt = None + UpdatedAt = initialAt + Message = "waiting for design run" + ManifestPath = None + ManifestUrl = None + Error = None + IsStale = false } + + let publish next = + for subscriber in subscribers do + subscriber.Writer.TryWrite(next) |> ignore + + let eventIds = function + | RunEvent.Start value -> value.Ids + | RunEvent.Heartbeat value -> value.Ids + | RunEvent.Checkpoint value -> value.Ids + | RunEvent.Complete value -> value.Ids + | RunEvent.Fail value -> value.Ids + + let eventName = function + | RunEvent.Start _ -> "start" + | RunEvent.Heartbeat _ -> "heartbeat" + | RunEvent.Checkpoint _ -> "checkpoint" + | RunEvent.Complete _ -> "complete" + | RunEvent.Fail _ -> "fail" + + let rejectIfNotRunning event = + if state.Status <> "running" then + Some $"cannot apply {eventName event} to run in status {state.Status}" + else + None + + let checkpointCount () = + if state.CurrentStepIndex < 0 then 0 else state.CurrentStepIndex + 1 + + let commit next = + state <- next + publish next + Ok next + + member _.Snapshot() = lock gate (fun () -> state) + + member _.Apply(event: RunEvent) : Result<RunSnapshot, string> = + lock gate (fun () -> + if eventIds event <> ids then + Error "run event identity does not match the selected run" + else + match event with + | RunEvent.Start value -> + if state.Status <> "idle" then + Error $"cannot apply start to run in status {state.Status}" + elif value.TotalSteps <= 0 then + Error "run must contain at least one step" + else + commit + { state with + Status = "running" + TotalSteps = value.TotalSteps + CurrentStepId = None + CurrentStepIndex = -1 + LastHeartbeatAt = Some value.At + UpdatedAt = value.At + Message = "design run started" + Error = None + IsStale = false } + | RunEvent.Heartbeat value -> + match rejectIfNotRunning event with + | Some message -> Error message + | None -> + commit + { state with + LastHeartbeatAt = Some value.At + UpdatedAt = value.At + Message = value.Message + IsStale = false } + | RunEvent.Checkpoint value -> + match rejectIfNotRunning event with + | Some message -> Error message + | None when value.TotalSteps <> state.TotalSteps -> + Error "checkpoint total steps do not match run" + | None when value.StepIndex < 0 || value.StepIndex >= state.TotalSteps -> + Error "checkpoint step index is outside the run" + | None when value.StepIndex <> checkpointCount () -> + Error "checkpoint step index is out of order" + | None -> + commit + { state with + CurrentStepId = Some value.StepId + CurrentStepIndex = value.StepIndex + LastHeartbeatAt = Some value.At + UpdatedAt = value.At + Message = value.Message + IsStale = false } + | RunEvent.Complete value -> + match rejectIfNotRunning event with + | Some message -> Error message + | None when checkpointCount () <> state.TotalSteps -> + Error "cannot complete run before every step is checkpointed" + | None -> + commit + { state with + Status = "complete" + CurrentStepId = None + LastHeartbeatAt = Some value.At + UpdatedAt = value.At + Message = "design run complete" + ManifestPath = Some value.ManifestPath + ManifestUrl = Some $"/api/artifacts/manifest?projectId={ids.ProjectId}&runId={ids.RunId}" + Error = None + IsStale = false } + | RunEvent.Fail value -> + if state.Status <> "idle" && state.Status <> "running" then + Error $"cannot apply fail to run in status {state.Status}" + else + commit + { state with + Status = "failed" + CurrentStepId = value.StepId + LastHeartbeatAt = Some value.At + UpdatedAt = value.At + Message = value.Message + Error = Some value.Message + IsStale = false }) + + member _.Observe(now: DateTimeOffset) = + lock gate (fun () -> + let next = + match state.Status, state.LastHeartbeatAt with + | "running", Some heartbeat when now - heartbeat > staleAfter && not state.IsStale -> + Some + { state with + UpdatedAt = now + Message = "heartbeat timeout" + IsStale = true } + | "running", Some heartbeat when now - heartbeat <= staleAfter && state.IsStale -> + Some + { state with + UpdatedAt = now + Message = "heartbeat resumed" + IsStale = false } + | _ -> None + + match next with + | Some value -> + state <- value + publish value + value + | None -> state) + + member _.Subscribe() = + let channel = Channel.CreateUnbounded<RunSnapshot>() + + lock gate (fun () -> + channel.Writer.TryWrite(state) |> ignore + subscribers.Add(channel)) + + new ArtifactRunSubscription( + channel.Reader, + fun () -> + lock gate (fun () -> + subscribers.Remove(channel) |> ignore + channel.Writer.TryComplete() |> ignore)) + +type ArtifactRunRegistry(staleAfter: TimeSpan, clock: unit -> DateTimeOffset) = + let gate = obj () + let runs = Dictionary<string * string, ArtifactRunStore>() + + member _.TryCreate(ids: RunIds) = + lock gate (fun () -> + let key = ids.ProjectId, ids.RunId + + match runs.TryGetValue(key) with + | true, _ -> None + | false, _ -> + let store = ArtifactRunStore(ids, staleAfter, clock ()) + runs.Add(key, store) + Some store) + + member _.GetOrCreate(ids: RunIds) = + lock gate (fun () -> + let key = ids.ProjectId, ids.RunId + + match runs.TryGetValue(key) with + | true, store -> store + | false, _ -> + let store = ArtifactRunStore(ids, staleAfter, clock ()) + runs.Add(key, store) + store) + + member _.TryFind(projectId: string, runId: string) = + lock gate (fun () -> + match runs.TryGetValue((projectId, runId)) with + | true, store -> Some store + | false, _ -> None) + + member _.ObserveAll(now: DateTimeOffset) = + lock gate (fun () -> runs.Values |> Seq.map (fun store -> store.Observe(now)) |> Seq.toArray) |
