diff options
| author | Somhairle H. Marisol <[email protected]> | 2026-09-21 00:58:15 +0800 |
|---|---|---|
| committer | Somhairle H. Marisol <[email protected]> | 2026-09-21 00:58:15 +0800 |
| commit | 7b29dbff7e43e87741476f47df17ad1582b16546 (patch) | |
| tree | 4c9d0832b4e4edd77e59e9cb0a88abf5e729a9f9 /src/SomhairlesDream.Server | |
| parent | 47add3088d50f12c5f48591f26c53cdb59d8283b (diff) | |
| download | somhairles-dream-fsharp-7b29dbff7e43e87741476f47df17ad1582b16546.tar.gz | |
feat(pipeline): add runnable artifact pipeline increment
[Change Nature]
- This commit adds the first runnable increment of the F# artifact
pipeline: CLI, server, frontend bundle, stage fixture, and tests.
It is feature work, not a bug fix.
[New Capability]
- CLI `run` executes the three-step Blender fixture through the process
bridge and writes GLBs, step logs, events, and a manifest; `verify`
validates a manifest without Blender.
- Server exposes run start, run-scoped SSE, artifact manifest/file
routes, and serves the Fable frontend from its build output.
- Frontend starts runs, follows run-scoped SSE progress, and links the
artifact manifest on completion.
[Implementation]
- Shared layer holds domain, run state, replay, and artifact contracts;
Modeling holds the pipeline, 120s bridge timeout, and result-file
protocol; stage/blender/fixture.py owns geometry in Python (bpy owns
geometry, F# owns orchestration).
- Fable bundle is regenerated into public/ and copied by the server
project; solution includes all projects with warnings-as-errors.
- README documents build, Fable, server, interpreter/env, fixture
limitations, and current verification status.
[Impact]
- Verified: 40/40 solution tests pass; Release build 0 warnings/errors;
real CLI run with bpy 5.0.1 produced 3 GLBs, 3 logs, and a
verify-clean manifest; isolated-port server smoke flow succeeded.
- Not yet verified: browser end-to-end behavior and production
fidelity. Artifacts and caches are gitignored.
Diffstat (limited to 'src/SomhairlesDream.Server')
| -rw-r--r-- | src/SomhairlesDream.Server/ArtifactRunApi.fs | 161 | ||||
| -rw-r--r-- | src/SomhairlesDream.Server/ArtifactRunCoordinator.fs | 141 | ||||
| -rw-r--r-- | src/SomhairlesDream.Server/Program.fs | 97 | ||||
| -rw-r--r-- | src/SomhairlesDream.Server/SomhairlesDream.Server.fsproj | 21 | ||||
| -rw-r--r-- | src/SomhairlesDream.Server/Store.fs | 259 |
5 files changed, 679 insertions, 0 deletions
diff --git a/src/SomhairlesDream.Server/ArtifactRunApi.fs b/src/SomhairlesDream.Server/ArtifactRunApi.fs new file mode 100644 index 0000000..8d76ffc --- /dev/null +++ b/src/SomhairlesDream.Server/ArtifactRunApi.fs @@ -0,0 +1,161 @@ +namespace SomhairlesDream.Server + +open System +open System.IO +open System.Text.Json +open System.Threading.Tasks +open Microsoft.AspNetCore.Builder +open Microsoft.AspNetCore.Http +open SomhairlesDream.Shared + +[<CLIMutable>] +type RunStartRequest = + { ProjectId: string + RunId: string + BaseId: string + TargetId: string + Render: bool } + +type ArtifactRunApiOptions = + { Coordinator: ArtifactRunCoordinator + Clock: unit -> DateTimeOffset + JsonOptions: JsonSerializerOptions } + +module ArtifactRunApi = + let private json options value = Results.Json(value, options.JsonOptions) + + let private jsonWithStatus options statusCode value = + Results.Json(value, options.JsonOptions, statusCode = Nullable<int>(statusCode)) + + let private error options statusCode message = + jsonWithStatus options statusCode {| error = message |} + + let private readStartRequest options (request: HttpRequest) = + task { + try + let! value = request.ReadFromJsonAsync<RunStartRequest>(options.JsonOptions) + + if isNull (box value) then + return Error "request body must be a JSON object" + else + return Ok value + with + | :? JsonException as ex -> return Error $"invalid JSON: {ex.Message}" + } + + let private start options (context: HttpContext) : Task<IResult> = + task { + let! request = readStartRequest options context.Request + + match request with + | Error message -> return error options 400 message + | Ok request -> + let ids : RunIds = + { ProjectId = request.ProjectId + RunId = request.RunId + BaseId = request.BaseId + TargetId = request.TargetId } + + match options.Coordinator.Start(ids, request.Render) with + | Ok(Accepted snapshot) -> return jsonWithStatus options 202 snapshot + | Ok(Conflict snapshot) -> return jsonWithStatus options 409 snapshot + | Error message when message = "run registry entry disappeared" -> + return error options 500 message + | Error message -> return error options 400 message + } + + let private queryValue (context: HttpContext) name = + let value = context.Request.Query[name].ToString() + + if String.IsNullOrWhiteSpace(value) then + None + else + Some value + + let private selector context = + match queryValue context "projectId", queryValue context "runId" with + | Some projectId, Some runId -> + match Ids.validateResourceId projectId, Ids.validateResourceId runId with + | Error message, _ -> Error $"invalid projectId: {message}" + | _, Error message -> Error $"invalid runId: {message}" + | Ok (), Ok () -> Ok(projectId, runId) + | None, _ -> Error "missing query parameter 'projectId'" + | _, None -> Error "missing query parameter 'runId'" + + let private manifest options (context: HttpContext) = + match selector context with + | Error message -> error options 400 message + | Ok(projectId, runId) -> + match options.Coordinator.Manifest(projectId, runId) with + | Ok value -> json options value + | Error RunNotFound -> error options 404 "run not found" + | Error(RunNotComplete snapshot) -> jsonWithStatus options 409 snapshot + | Error(InvalidManifest message) -> error options 500 message + | Error(InvalidArtifactPath message) -> error options 500 message + | Error ArtifactNotFound -> error options 500 "manifest not found" + + let private artifactContentType (path: string) = + match Path.GetExtension(path).ToLowerInvariant() with + | ".glb" -> "model/gltf-binary" + | ".json" -> "application/json" + | ".png" -> "image/png" + | ".log" -> "text/plain" + | _ -> "application/octet-stream" + + let private artifact options (context: HttpContext) = + match selector context, queryValue context "path" with + | Error message, _ -> error options 400 message + | _, None -> error options 400 "missing query parameter 'path'" + | Ok(projectId, runId), Some relativePath -> + match options.Coordinator.Artifact(projectId, runId, relativePath) with + | Ok path -> Results.File(path, artifactContentType path) + | Error RunNotFound -> error options 404 "run not found" + | Error(RunNotComplete snapshot) -> jsonWithStatus options 409 snapshot + | Error(InvalidArtifactPath message) -> error options 400 message + | Error ArtifactNotFound -> error options 404 "artifact not found" + | Error(InvalidManifest message) -> error options 500 message + + let private writeError options (context: HttpContext) statusCode message = + task { + context.Response.StatusCode <- statusCode + context.Response.ContentType <- "application/json" + let payload = JsonSerializer.Serialize({| error = message |}, options.JsonOptions) + do! context.Response.WriteAsync(payload) + } + + let private events options (context: HttpContext) : Task = + task { + match selector context with + | Error message -> do! writeError options context 400 message + | Ok(projectId, runId) -> + match options.Coordinator.Subscribe(projectId, runId) with + | None -> do! writeError options context 404 "run not found" + | Some subscription -> + use _lifetime = subscription :> IDisposable + context.Response.StatusCode <- 200 + context.Response.ContentType <- "text/event-stream" + context.Response.Headers.CacheControl <- "no-cache" + context.Response.Headers["X-Accel-Buffering"] <- "no" + + try + while not context.RequestAborted.IsCancellationRequested do + let! snapshot = subscription.Reader.ReadAsync(context.RequestAborted).AsTask() + let payload = JsonSerializer.Serialize(snapshot, options.JsonOptions) + do! context.Response.WriteAsync($"data: {payload}\n\n", context.RequestAborted) + do! context.Response.Body.FlushAsync(context.RequestAborted) + with + | :? OperationCanceledException -> () + } + + let register (app: WebApplication) options = + app.MapPost("/api/runs/start", Func<HttpContext, Task<IResult>>(fun context -> start options context)) + |> ignore + + app.MapGet("/api/runs/events", Func<HttpContext, Task>(fun context -> events options context)) + |> ignore + + app.MapGet("/api/artifacts/manifest", Func<HttpContext, IResult>(fun context -> manifest options context)) + |> ignore + + app.MapGet("/api/artifacts/file", Func<HttpContext, IResult>(fun context -> artifact options context)) + |> ignore diff --git a/src/SomhairlesDream.Server/ArtifactRunCoordinator.fs b/src/SomhairlesDream.Server/ArtifactRunCoordinator.fs new file mode 100644 index 0000000..78d5dee --- /dev/null +++ b/src/SomhairlesDream.Server/ArtifactRunCoordinator.fs @@ -0,0 +1,141 @@ +namespace SomhairlesDream.Server + +open System +open System.IO +open System.Threading.Tasks +open SomhairlesDream.Modeling +open SomhairlesDream.Shared + +type RunStartOutcome = + | Accepted of RunSnapshot + | Conflict of RunSnapshot + +type ArtifactLookupError = + | RunNotFound + | RunNotComplete of RunSnapshot + | InvalidArtifactPath of string + | ArtifactNotFound + | InvalidManifest of string + +type ArtifactRunCoordinator( + artifactRoot: string, + bridgeFactory: unit -> IArtifactBridge, + clock: unit -> DateTimeOffset, + staleAfter: TimeSpan +) = + let root = Path.GetFullPath(artifactRoot) + let registry = ArtifactRunRegistry(staleAfter, clock) + + let validateIds (ids: RunIds) = + [| ids.ProjectId; ids.RunId; ids.BaseId; ids.TargetId |] + |> Array.tryFind (Ids.validateResourceId >> Result.isError) + |> function + | Some invalid -> Error $"invalid resource id: {invalid}" + | None -> Ok () + + let runDirectory (projectId: string) (runId: string) = + Path.Combine(root, projectId, runId) + + let failIfNeeded (store: ArtifactRunStore) (ids: RunIds) message = + let snapshot = store.Snapshot() + + if snapshot.Status = "idle" || snapshot.Status = "running" then + store.Apply( + RunEvent.Fail + { Ids = ids + At = clock () + StepId = snapshot.CurrentStepId + Message = message } + ) + |> ignore + + let runPipeline (store: ArtifactRunStore) (ids: RunIds) render = + try + let options : PipelineOptions = + { ArtifactRoot = root + Ids = ids + Render = render + Bridge = bridgeFactory () + Clock = clock + OnEvent = + fun event -> + match store.Apply event with + | Ok _ -> () + | Error message -> invalidOp message } + + match Pipeline.run options with + | Ok _ -> () + | Error message -> failIfNeeded store ids message + with ex -> + failIfNeeded store ids ex.Message + + let safeArtifactPath (runDirectory: string) (relativePath: string) = + if String.IsNullOrWhiteSpace(relativePath) || Path.IsPathRooted(relativePath) then + Error "artifact path must be relative" + else + let normalized = relativePath.Replace('\\', '/') + let segments = normalized.Split('/', StringSplitOptions.RemoveEmptyEntries) + + if segments |> Array.exists (fun segment -> segment = ".." || segment = ".") then + Error "artifact path contains traversal" + else + let fullPath = Path.GetFullPath(Path.Combine(runDirectory, normalized.Replace('/', Path.DirectorySeparatorChar))) + let prefix = Path.GetFullPath(runDirectory).TrimEnd(Path.DirectorySeparatorChar) + string Path.DirectorySeparatorChar + + if fullPath.StartsWith(prefix, StringComparison.Ordinal) then + Ok fullPath + else + Error "artifact path escapes run directory" + + member _.ArtifactRoot = root + + member _.Start(ids: RunIds, render: bool) : Result<RunStartOutcome, string> = + match validateIds ids with + | Error message -> Error message + | Ok () -> + match registry.TryCreate(ids) with + | Some store -> + Task.Run(fun () -> runPipeline store ids render) |> ignore + Ok(Accepted(store.Snapshot())) + | None -> + match registry.TryFind(ids.ProjectId, ids.RunId) with + | Some store -> Ok(Conflict(store.Observe(clock ()))) + | None -> Error "run registry entry disappeared" + + member _.TryFind(projectId: string, runId: string) = registry.TryFind(projectId, runId) + + member _.Subscribe(projectId: string, runId: string) = + registry.TryFind(projectId, runId) |> Option.map (fun store -> store.Subscribe()) + + member _.ObserveAll() = registry.ObserveAll(clock ()) + + member _.Manifest(projectId: string, runId: string) : Result<ArtifactManifest, ArtifactLookupError> = + match registry.TryFind(projectId, runId) with + | None -> Error RunNotFound + | Some store -> + let snapshot = store.Observe(clock ()) + + if snapshot.Status <> "complete" then + Error(RunNotComplete snapshot) + else + let path = Path.Combine(runDirectory projectId runId, "manifest.json") + + match ArtifactVerifier.verify path with + | Ok manifest -> Ok manifest + | Error message -> Error(InvalidManifest message) + + member _.Artifact(projectId: string, runId: string, relativePath: string) : Result<string, ArtifactLookupError> = + match registry.TryFind(projectId, runId) with + | None -> Error RunNotFound + | Some store -> + let snapshot = store.Observe(clock ()) + + if snapshot.Status <> "complete" then + Error(RunNotComplete snapshot) + else + let directory = runDirectory projectId runId + + match safeArtifactPath directory relativePath with + | Error message -> Error(InvalidArtifactPath message) + | Ok path when not (File.Exists(path)) -> Error ArtifactNotFound + | Ok path -> Ok path diff --git a/src/SomhairlesDream.Server/Program.fs b/src/SomhairlesDream.Server/Program.fs new file mode 100644 index 0000000..53cfef6 --- /dev/null +++ b/src/SomhairlesDream.Server/Program.fs @@ -0,0 +1,97 @@ +open System +open System.IO +open System.Text.Json +open System.Threading +open System.Threading.Tasks +open Microsoft.AspNetCore.Builder +open Microsoft.AspNetCore.Hosting +open Microsoft.AspNetCore.Http +open Microsoft.Extensions.FileProviders +open Microsoft.Extensions.Hosting +open SomhairlesDream.Server +open SomhairlesDream.Shared + +let builder = WebApplication.CreateBuilder(Environment.GetCommandLineArgs() |> Array.skip 1) +let publicRoot = Path.Combine(AppContext.BaseDirectory, "public") +builder.Environment.WebRootPath <- publicRoot +builder.Environment.WebRootFileProvider <- new PhysicalFileProvider(publicRoot) + +let app = builder.Build() +let store = RunStore("heritage-001", "run-001", TimeSpan.FromMinutes 2.) +let jsonOptions = JsonSerializerOptions(JsonSerializerDefaults.Web) + +let environmentOrDefault name fallback = + match Environment.GetEnvironmentVariable(name) with + | value when not (String.IsNullOrWhiteSpace(value)) -> value + | _ -> fallback + +let artifactRoot = environmentOrDefault "SOMHAIRLES_ARTIFACT_ROOT" (Path.Combine(Environment.CurrentDirectory, ".artifacts")) +let stageRoot = environmentOrDefault "SOMHAIRLES_STAGE_ROOT" (Path.Combine(Environment.CurrentDirectory, "stage")) +let python = environmentOrDefault "SOMHAIRLES_BPYTHON" "python3" +let artifactClock () = DateTimeOffset.UtcNow +let coordinator = + ArtifactRunCoordinator( + artifactRoot, + (fun () -> SomhairlesDream.Modeling.BlenderProcessBridge(python, stageRoot, TimeSpan.FromSeconds 120.) :> SomhairlesDream.Modeling.IArtifactBridge), + artifactClock, + TimeSpan.FromMinutes 2. + ) + +let json value = Results.Json(value, jsonOptions) + +let events (context: HttpContext) : Task = + (task { + context.Response.ContentType <- "text/event-stream" + context.Response.Headers.CacheControl <- "no-cache" + + let reader = store.Subscribe() + + try + while not context.RequestAborted.IsCancellationRequested do + let! snapshot = reader.ReadAsync(context.RequestAborted).AsTask() + let payload = JsonSerializer.Serialize(snapshot, jsonOptions) + do! context.Response.WriteAsync($"data: {payload}\n\n", context.RequestAborted) + do! context.Response.Body.FlushAsync(context.RequestAborted) + with + | :? OperationCanceledException -> () + } :> Task) + +app.UseDefaultFiles() |> ignore +app.UseStaticFiles() |> ignore + +app.MapGet("/health", Func<IResult>(fun () -> Results.Ok({| status = "ok" |}))) |> ignore + +ArtifactRunApi.register + app + { Coordinator = coordinator + Clock = artifactClock + JsonOptions = jsonOptions } + +let observeRuns = + Task.Run( + Func<Task>(fun () -> + task { + try + while not app.Lifetime.ApplicationStopping.IsCancellationRequested do + coordinator.ObserveAll() |> ignore + do! Task.Delay(TimeSpan.FromSeconds 1., app.Lifetime.ApplicationStopping) + with + | :? OperationCanceledException -> () + } :> Task) + ) + +app.MapGet( + "/api/state", + Func<IResult>(fun () -> store.Observe(DateTimeOffset.UtcNow) |> json) +) +|> ignore + +app.MapPost( + "/api/runs/heartbeat", + Func<IResult>(fun () -> store.Heartbeat(DateTimeOffset.UtcNow) |> json) +) +|> ignore + +app.MapGet("/api/events", Func<HttpContext, Task>(events)) |> ignore + +app.Run() diff --git a/src/SomhairlesDream.Server/SomhairlesDream.Server.fsproj b/src/SomhairlesDream.Server/SomhairlesDream.Server.fsproj new file mode 100644 index 0000000..17c159d --- /dev/null +++ b/src/SomhairlesDream.Server/SomhairlesDream.Server.fsproj @@ -0,0 +1,21 @@ +<Project Sdk="Microsoft.NET.Sdk.Web"> + <PropertyGroup> + <TargetFramework>net8.0</TargetFramework> + <RootNamespace>SomhairlesDream.Server</RootNamespace> + <AssemblyName>SomhairlesDream.Server</AssemblyName> + <OutputType>Exe</OutputType> + </PropertyGroup> + <ItemGroup> + <ProjectReference Include="../SomhairlesDream.Shared/SomhairlesDream.Shared.fsproj" /> + <ProjectReference Include="../SomhairlesDream.Modeling/SomhairlesDream.Modeling.fsproj" /> + </ItemGroup> + <ItemGroup> + <Compile Include="Store.fs" /> + <Compile Include="ArtifactRunCoordinator.fs" /> + <Compile Include="ArtifactRunApi.fs" /> + <Compile Include="Program.fs" /> + </ItemGroup> + <ItemGroup> + <Content Include="../../public/**/*" Link="public/%(RecursiveDir)%(Filename)%(Extension)" CopyToOutputDirectory="PreserveNewest" CopyToPublishDirectory="PreserveNewest" /> + </ItemGroup> +</Project> 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) |
