namespace SomhairlesDream.Server open System open System.Collections.Concurrent open System.Collections.Generic open System.IO open System.Security.Cryptography 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 | ArtifactMismatch of string type ArtifactFile = { Stream: Stream RelativePath: 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 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, String.Join("/", segments)) else Error "artifact path escapes run directory" let checkpointDigests = ConcurrentDictionary>() let recordCheckpointDigest (ids: RunIds) (relativePath: string) = try match safeArtifactPath (runDirectory ids.ProjectId ids.RunId) relativePath with | Ok(fullPath, normalized) when File.Exists(fullPath) -> use stream = File.OpenRead(fullPath) use sha = SHA256.Create() let digest = (Convert.ToHexString(sha.ComputeHash(stream)).ToLowerInvariant(), stream.Length) let digests = checkpointDigests.GetOrAdd( (ids.ProjectId, ids.RunId), fun _ -> ConcurrentDictionary() ) digests[normalized] <- digest | _ -> () with _ -> () let rejectLinks (runDirectory: string) (fullPath: string) = let root = Path.GetFullPath(runDirectory) let relative = fullPath.Substring(root.Length + 1) let mutable current = root let mutable outcome = Ok () for segment in relative.Split( [| Path.DirectorySeparatorChar; Path.AltDirectorySeparatorChar |], StringSplitOptions.RemoveEmptyEntries ) do match outcome with | Error _ -> () | Ok () -> current <- Path.Combine(current, segment) let info: FileSystemInfo = if Directory.Exists(current) && not (File.Exists(current)) then DirectoryInfo(current) :> FileSystemInfo else FileInfo(current) :> FileSystemInfo if not (isNull info.LinkTarget) then outcome <- Error "artifact path crosses a filesystem link" outcome let openValidated (runDirectory: string) (fullPath: string) (normalized: string) (expectedSha256: string) (expectedBytes: int64) : Result = match rejectLinks runDirectory fullPath with | Error message -> Error(InvalidArtifactPath message) | Ok () -> try if not (File.Exists(fullPath)) then Error ArtifactNotFound else let stream = File.Open(fullPath, FileMode.Open, FileAccess.Read, FileShare.Read) try if stream.Length <> expectedBytes then Error(ArtifactMismatch $"artifact byte count mismatch: {normalized}") else use sha = SHA256.Create() let sha256 = Convert.ToHexString(sha.ComputeHash(stream)).ToLowerInvariant() if sha256 <> expectedSha256 then Error(ArtifactMismatch $"artifact hash mismatch: {normalized}") else stream.Seek(0L, SeekOrigin.Begin) |> ignore Ok { Stream = stream :> Stream; RelativePath = normalized } with _ -> stream.Dispose() reraise () with | :? FileNotFoundException | :? DirectoryNotFoundException -> Error ArtifactNotFound | :? IOException as ex -> Error(InvalidArtifactPath $"artifact unavailable: {ex.Message}") let verifiedManifest (projectId: string) (runId: string) = let path = Path.Combine(runDirectory projectId runId, "manifest.json") match ArtifactVerifier.verify path with | Ok manifest -> Ok manifest | Error message -> Error(InvalidManifest message) let memberDigests (manifest: ArtifactManifest) = let map = Dictionary() let add (relativePath: string) (sha256: string) (bytes: int64) = map[relativePath] <- (sha256, bytes) for step in manifest.Steps do add step.ArtifactPath step.Sha256 step.Bytes step.Render |> Option.iter (fun render -> add render.ArtifactPath render.Sha256 render.Bytes) map let completedDigest (projectId: string) (runId: string) (normalized: string) = match verifiedManifest projectId runId with | Error lookupError -> Error lookupError | Ok manifest -> match (memberDigests manifest).TryGetValue(normalized) with | true, digest -> Ok digest | false, _ -> Error ArtifactNotFound let checkpointDigest (projectId: string) (runId: string) (normalized: string) = match checkpointDigests.TryGetValue((projectId, runId)) with | true, digests -> match digests.TryGetValue(normalized) with | true, digest -> Ok digest | false, _ -> Error ArtifactNotFound | false, _ -> Error ArtifactNotFound 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 event with | Checkpoint value -> recordCheckpointDigest value.Ids value.ArtifactPath | _ -> ()) 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 member _.ArtifactRoot = root member _.Start(ids: RunIds, render: bool) : Result = 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 = 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 verifiedManifest projectId runId member _.Artifact(projectId: string, runId: string, relativePath: string) : Result = 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(fullPath, normalized) -> match rejectLinks directory fullPath with | Error message -> Error(InvalidArtifactPath message) | Ok () -> match completedDigest projectId runId normalized with | Error lookupError -> Error lookupError | Ok(sha256, bytes) -> openValidated directory fullPath normalized sha256 bytes member _.StepArtifact(projectId: string, runId: string, relativePath: string) : Result = match registry.TryFind(projectId, runId) with | None -> Error RunNotFound | Some store -> let directory = runDirectory projectId runId match safeArtifactPath directory relativePath with | Error message -> Error(InvalidArtifactPath message) | Ok(_, normalized) when not (normalized.EndsWith(".glb", StringComparison.OrdinalIgnoreCase)) -> Error(InvalidArtifactPath "step artifacts must be .glb files") | Ok(_, normalized) when not (normalized.StartsWith("steps/", StringComparison.Ordinal)) -> Error(InvalidArtifactPath "step artifacts must live under steps/") | Ok(fullPath, normalized) -> let snapshot = store.Observe(clock ()) let digest = match rejectLinks directory fullPath with | Error message -> Error(InvalidArtifactPath message) | Ok () -> if snapshot.Status = "complete" then completedDigest projectId runId normalized elif snapshot.CompletedSteps |> Array.exists (fun recorded -> recorded = normalized) then checkpointDigest projectId runId normalized else Error ArtifactNotFound match digest with | Error lookupError -> Error lookupError | Ok(sha256, bytes) -> openValidated directory fullPath normalized sha256 bytes