/** * #468 — the in-flight short-circuit for `execute-status`. * * `produce` writes its execution report only after the provider answers, so * during the submit window the report on disk belongs to the PREVIOUS run. The * status ladder in execution-status.ts reads that report and, in two branches * (missing report; blocked / job-less report), rewrites it. This module runs * first: if a run marker with a fresh heartbeat exists AND its scene scope * overlaps what the report on disk covers, it returns a pending, non-mutating * answer that names the run. A live run scoped to scenes disjoint from the * report's (a parallel `produce --scene` from another shell) falls through, so * the already-submitted scenes keep getting polled. Stale markers (a killed * `produce`) are reaped and recorded as `execution.run.abandoned` on the way. */ import { existsSync } from 'node:fs'; import { readFile } from 'node:fs/promises'; import { appendProjectEvent } from './events.js'; import { clearExecutionRunMarker, describeExecutionRuns } from './execution-run-marker.js'; import { isProviderRouteId } from './provider-platform/types.js'; import type { VideoProjectWorkspace } from './workspace.js'; import type { VideoExecutionPollResult, VideoExecutionReport, VideoProductionMode } from './types.js'; export interface InFlightExecutionResult { reportPath: string; report: VideoExecutionReport; poll: VideoExecutionPollResult; } /** * Does a live run's scene scope overlap the scenes the on-disk report covers? * `null` on either side means "the whole storyboard"; a report with no * `candidatesByScene` is a legacy single-job run, which is whole-storyboard too. * Exported for tests. */ export function runScopeOverlaps( markerScenes: number[] | null, reportScenes: number[] | null, ): boolean { if (markerScenes === null || reportScenes === null) return true; const reportSet = new Set(reportScenes); return markerScenes.some((index) => reportSet.has(index)); } export async function resolveInFlightExecution( workspace: VideoProjectWorkspace, input: { projectSlug: string; productionMode: VideoProductionMode; reportPath: string }, ): Promise { // The marker READ stays loud: an unreadable state dir is a real problem. const runs = await describeExecutionRuns(workspace); for (const run of runs.filter((r) => r.description.state === 'stale')) { // Reaping is housekeeping — it must never take down a status check that // worked before markers existed (read-only state dir, EROFS, …). try { await clearExecutionRunMarker(run); await appendProjectEvent(workspace, { type: 'execution.run.abandoned', payload: { pid: run.marker.pid, startedAt: run.marker.startedAt, heartbeatAt: run.marker.heartbeatAt, routeId: run.marker.routeId, sceneIndexes: run.marker.sceneIndexes, reason: run.description.reason, }, }); } catch { // best-effort — the next status check will try again } } const liveRuns = runs.filter((r) => r.description.state === 'live'); if (liveRuns.length === 0) return null; const existingReport = existsSync(input.reportPath) ? JSON.parse(await readFile(input.reportPath, 'utf-8')) as VideoExecutionReport : null; const reportScenes = existingReport?.candidatesByScene ? existingReport.candidatesByScene.map((entry) => entry.sceneIndex) : null; const overlapping = liveRuns.filter((run) => runScopeOverlaps(run.marker.sceneIndexes, reportScenes)); if (overlapping.length === 0) return null; const issues = overlapping.map((run) => { const scope = run.marker.sceneIndexes ? `scenes ${run.marker.sceneIndexes.join(',')}` : 'all scenes'; return `A \`vclaw video produce\` run started at ${run.marker.startedAt} (pid ${run.marker.pid}, ` + `${run.marker.routeId ?? 'route unresolved'}, ${scope}) is still submitting ` + `(${run.description.elapsedSeconds}s elapsed). The execution report on disk is the PREVIOUS run's; ` + 'nothing was polled and nothing was written. Re-check once it finishes. If it is wedged, confirm the ' + `process first (\`ps -p ${run.marker.pid} -o command\`) and only then kill it — pids are reused.`; }); const firstRoute = overlapping[0].marker.routeId; const report: VideoExecutionReport = existingReport ?? { projectSlug: input.projectSlug, productionMode: input.productionMode, operationKind: 'text-to-video', routeId: isProviderRouteId(firstRoute) ? firstRoute : null, status: 'blocked', dryRun: false, generatedAt: new Date().toISOString(), blockers: [], executedSteps: ['execution-status-requested'], taskCount: 0, }; return { reportPath: input.reportPath, report, poll: { status: 'pending', externalJobId: null, outputs: [], issues, rawResult: { reason: 'execution-in-flight', runs: overlapping.map((run) => ({ pid: run.marker.pid, startedAt: run.marker.startedAt, heartbeatAt: run.marker.heartbeatAt, routeId: run.marker.routeId, sceneIndexes: run.marker.sceneIndexes, elapsedSeconds: run.description.elapsedSeconds, })), }, }, }; }