/** * Execution run markers — the one thing `produce` writes to disk BEFORE it * submits to a provider. * * WHY THIS EXISTS (#468) * ---------------------- * `executeProject` writes nothing until the provider answers: candidates, the * run contract and `execution-report.json` all land AFTER `submitExecutionPayload` * returns (~95 s on Flow, minutes on Seedance). During that window a status check * from a second shell read the PREVIOUS run's report, called a healthy in-flight * job "failed", and in two branches rewrote the report underneath the running * process. Nothing else on disk could prove a run was live: the lane ticket only * exists for routes with a declared concurrency cap, and `.vclaw-jobs/*.json` is * written after the render on the Flow and Seedance transports. * * SHAPE * ----- * One small JSON file PER RUN under `projects//state/execution-runs/`, * named `--.json` (the nonce keeps two runs started by * one process in the same millisecond — `render-scenes`, parallel `--scene` * calls from a library caller — from sharing a file). Per run, not per project, because parallel * `produce --scene` from separate shells is a supported path — two runs must * coexist and each must clear only its own marker. The writer heartbeats the * file every `EXECUTION_RUN_HEARTBEAT_MS` and removes it in a `finally` once the * report is on disk, so a marker that is still there means a process that was * killed, never a completed run. * * This is runtime state, not a canonical artifact: it is deliberately NOT under * `artifacts/` and has no JSON schema (same class as `batch-queue.json`). The * fields mirror what the Cinema queue's task lease carries (`payloadHash`, * `routeId`, `heartbeatAt`) so the marker can become that lease record when the * direct `produce` path converges on the queue (ADR 0008). * * LIVENESS — one rule * ------------------- * A run is `live` iff its heartbeat is within `EXECUTION_RUN_LIVE_WINDOW_MS`. * A pid probe is only an ACCELERATOR for staleness: `ESRCH` (no such process) * short-circuits to `stale`; `EPERM` or any other error is not evidence either * way and defers to the heartbeat. A live pid never proves a live run — pids are * reused. Pure `describeExecutionRun` takes the probe as a parameter so tests * never need a real process. */ import { mkdir, readdir, readFile, rm } from 'node:fs/promises'; import { join } from 'node:path'; import { writeTextFileAtomic } from './atomic-write.js'; import type { VideoProjectWorkspace } from './workspace.js'; export const EXECUTION_RUNS_DIRNAME = 'execution-runs'; /** A heartbeat older than this means the writer is gone (killed, hung, or the machine slept). */ export const EXECUTION_RUN_LIVE_WINDOW_MS = 10 * 60 * 1000; /** How often a running `produce` touches its marker. Well inside the live window. */ export const EXECUTION_RUN_HEARTBEAT_MS = 60 * 1000; export interface ExecutionRunMarker { schemaVersion: 1; pid: number; /** ISO time the run wrote its marker, i.e. immediately before the provider submit. */ startedAt: string; /** ISO time of the last heartbeat; equals `startedAt` until the first touch. */ heartbeatAt: string; routeId: string | null; productionMode: string; /** Scene subset of a `--scene` run; `null` means the whole storyboard. */ sceneIndexes: number[] | null; /** The lane request hash of the exact payload about to be submitted. */ payloadHash: string; } export interface ExecutionRunMarkerHandle { path: string; marker: ExecutionRunMarker; /** The heartbeat write currently in flight, if any — `clear` awaits it. */ pending?: Promise; /** Set by `clear`; a heartbeat that started before it must not resurrect the file. */ cleared?: boolean; } export type ExecutionRunLiveness = 'live' | 'stale'; export interface ExecutionRunDescription { state: ExecutionRunLiveness; /** Why the verdict landed where it did — surfaced verbatim to the operator. */ reason: 'heartbeat-fresh' | 'heartbeat-expired' | 'process-gone'; elapsedSeconds: number; heartbeatAgeSeconds: number; } export type PidProbe = (pid: number) => void; /** Default probe: `kill(pid, 0)` sends no signal but throws ESRCH when the pid is gone. */ export function defaultPidProbe(pid: number): void { process.kill(pid, 0); } export function executionRunsDir(workspace: Pick): string { return join(workspace.stateDir, EXECUTION_RUNS_DIRNAME); } function isIsoTime(value: unknown): value is string { return typeof value === 'string' && value.length > 0 && !Number.isNaN(new Date(value).getTime()); } export function isExecutionRunMarker(value: unknown): value is ExecutionRunMarker { if (!value || typeof value !== 'object' || Array.isArray(value)) return false; const record = value as Record; if (record.schemaVersion !== 1) return false; if (typeof record.pid !== 'number' || !Number.isInteger(record.pid) || record.pid <= 0) return false; if (!isIsoTime(record.startedAt) || !isIsoTime(record.heartbeatAt)) return false; if (record.routeId !== null && typeof record.routeId !== 'string') return false; if (typeof record.productionMode !== 'string') return false; if (record.sceneIndexes !== null && !(Array.isArray(record.sceneIndexes) && record.sceneIndexes.every((n) => typeof n === 'number'))) { return false; } return typeof record.payloadHash === 'string' && record.payloadHash.length > 0; } export async function writeExecutionRunMarker( workspace: Pick, input: { routeId: string | null; productionMode: string; sceneIndexes: number[] | null; payloadHash: string; pid?: number; now?: Date; }, ): Promise { const now = input.now ?? new Date(); const pid = input.pid ?? process.pid; const startedAt = now.toISOString(); const marker: ExecutionRunMarker = { schemaVersion: 1, pid, startedAt, heartbeatAt: startedAt, routeId: input.routeId, productionMode: input.productionMode, sceneIndexes: input.sceneIndexes ? [...input.sceneIndexes] : null, payloadHash: input.payloadHash, }; const dir = executionRunsDir(workspace); await mkdir(dir, { recursive: true }); const nonce = Math.random().toString(36).slice(2, 8); const path = join(dir, `${pid}-${now.getTime()}-${nonce}.json`); await writeTextFileAtomic(path, `${JSON.stringify(marker, null, 2)}\n`); return { path, marker }; } /** * Refresh `heartbeatAt`. Best-effort: a heartbeat that fails to write must never * fail the render it is describing — the worst case is the marker reads stale * ten minutes early, which the reader handles. (The INITIAL write in * `writeExecutionRunMarker` is deliberately NOT best-effort: without a marker * the in-flight guarantee is void, so a run that cannot write one refuses * before submitting rather than submitting blind.) * * Ordered against `clearExecutionRunMarker`: `clearInterval` cannot cancel a * touch that has already started, and the atomic write is `writeFile(tmp)` then * `rename` — a rename landing after the clear's `rm` would resurrect the marker * with a fresh heartbeat (live for ten minutes in a long-lived process such as * `render-scenes`). So a touch records itself on the handle for `clear` to * await, skips entirely once `cleared` is set, and re-removes the file if the * clear raced past it anyway. */ export function touchExecutionRunMarker( handle: ExecutionRunMarkerHandle, now: Date = new Date(), ): Promise { if (handle.cleared) return Promise.resolve(); handle.marker = { ...handle.marker, heartbeatAt: now.toISOString() }; const write = (async () => { try { await writeTextFileAtomic(handle.path, `${JSON.stringify(handle.marker, null, 2)}\n`); } catch { // advisory — see docblock } if (handle.cleared) await rm(handle.path, { force: true }); })(); handle.pending = write; return write; } export async function clearExecutionRunMarker( handle: Pick & Partial>, ): Promise { handle.cleared = true; if (handle.pending) await handle.pending; await rm(handle.path, { force: true }); } /** * Every marker currently on disk for the project. Tolerates a missing directory * (no run has ever been started) and skips files that are not a valid marker (a * torn write, or something else that landed in the directory) — one bad file * must not hide a real live run. */ export async function readExecutionRunMarkers( workspace: Pick, ): Promise { const dir = executionRunsDir(workspace); let names: string[]; try { names = await readdir(dir); } catch (error) { if ((error as NodeJS.ErrnoException).code === 'ENOENT') return []; throw error; } const handles: ExecutionRunMarkerHandle[] = []; for (const name of names.filter((n) => n.endsWith('.json')).sort()) { const path = join(dir, name); try { const parsed: unknown = JSON.parse(await readFile(path, 'utf-8')); if (isExecutionRunMarker(parsed)) handles.push({ path, marker: parsed }); } catch { continue; } } return handles; } /** Pure liveness verdict — see the module docblock for the one rule. */ export function describeExecutionRun( marker: ExecutionRunMarker, nowMs: number = Date.now(), probePid: PidProbe = defaultPidProbe, ): ExecutionRunDescription { const startedMs = new Date(marker.startedAt).getTime(); const heartbeatMs = new Date(marker.heartbeatAt).getTime(); const elapsedSeconds = Math.max(0, Math.round((nowMs - startedMs) / 1000)); const heartbeatAgeSeconds = Math.max(0, Math.round((nowMs - heartbeatMs) / 1000)); let processGone = false; try { probePid(marker.pid); } catch (error) { processGone = (error as NodeJS.ErrnoException).code === 'ESRCH'; } if (processGone) { return { state: 'stale', reason: 'process-gone', elapsedSeconds, heartbeatAgeSeconds }; } const fresh = nowMs - heartbeatMs <= EXECUTION_RUN_LIVE_WINDOW_MS; return { state: fresh ? 'live' : 'stale', reason: fresh ? 'heartbeat-fresh' : 'heartbeat-expired', elapsedSeconds, heartbeatAgeSeconds, }; } /** Convenience: every marker on disk, each paired with its liveness verdict. */ export async function describeExecutionRuns( workspace: Pick, options: { nowMs?: number; probePid?: PidProbe } = {}, ): Promise> { const handles = await readExecutionRunMarkers(workspace); return handles.map((handle) => ({ ...handle, description: describeExecutionRun(handle.marker, options.nowMs, options.probePid), })); }