import { createLaneCoordinator } from './lane-coordinator.js'; import { ROUTE_CAPABILITIES } from './provider-platform/route-capabilities.js'; import { existsSync } from 'node:fs'; import { resolveWorkspaceRootFromEnv } from './workspace-root.js'; import { readFile, readdir } from 'node:fs/promises'; import { join } from 'node:path'; import { artifactPathFor, writeArtifact } from './artifact-store.js'; import { writeStageCheckpoint } from './checkpoints.js'; import { appendProjectEvent } from './events.js'; import { appendGenerationTelemetry, buildGenerationTelemetryFromPoll, } from './generation-telemetry.js'; import { pollExecutionPayload } from './execution-runtime.js'; import { deriveAssetManifestFromSelection } from './scene-candidates.js'; import { readSceneCandidatesArtifact, sceneCandidatesPathFor, withSceneArtifactsLock, writeSceneCandidatesArtifact, } from './scene-candidate-store.js'; import { markPending } from './scene-selection.js'; import { readSceneSelectionArtifact, writeSceneSelectionArtifact, } from './scene-selection-store.js'; import { regenerateRunSurface } from './preview-portal/index.js'; import { buildProjectStatusReport } from './status.js'; import { ensureProjectWorkspace, readProjectManifest, resolveProjectWorkspace, updateProjectManifestState } from './workspace.js'; import type { AssetManifestArtifact } from './artifacts.js'; import type { ProviderRouteId } from './provider-platform/types.js'; import type { SceneCandidate, SceneCandidateOutput, SceneCandidatesArtifact, VideoExecutionPollResult, VideoExecutionReport, VideoProductionMode, } from './types.js'; function mergeAssets( existing: AssetManifestArtifact['assets'], incoming: AssetManifestArtifact['assets'], ): AssetManifestArtifact['assets'] { const byId = new Map(existing.map((asset) => [asset.id, asset])); for (const asset of incoming) { byId.set(asset.id, asset); } return [...byId.values()].sort((left, right) => left.id.localeCompare(right.id)); } /** * Ids of candidate-derived OUTPUT assets are `${candidate.id}-output-${n}` * (see `deriveAssetManifestFromSelection`). Operator/keyframe INPUT assets use * `${kind}-${path}` ids (see `parseAssetSpec` in `asset-spec.ts`), so this * pattern never matches an input keyframe. */ const DERIVED_OUTPUT_ID_RE = /-output-\d+$/; /** The canonical on-disk per-scene keyframe filename: `scene-keyframe.png`. */ const SCENE_KEYFRAME_RE = /^scene(\d+)-keyframe\.png$/; /** * The per-scene image keyframe SOURCE OF TRUTH on disk: * `projects//references/scene-keyframe.png` (the SAME location the * `keyframe-qc` gate reads). Returns one `image` asset per keyframe file found, * with the SAME `image-` id `parseAssetSpec` mints so it merges cleanly * with a manifest-registered keyframe. Used to SELF-HEAL a manifest that a prior * poll ratcheted to video-only: because these are re-read from disk every poll, * a lost keyframe is restored even when the manifest no longer carries it. * Absent `references/` dir → `[]` (byte-identical no-op). */ async function collectDiskKeyframeAssets( workspace: Awaited>, ): Promise { const referencesDir = join(workspace.projectDir, 'references'); let names: string[]; try { names = await readdir(referencesDir); } catch { return []; } const assets: AssetManifestArtifact['assets'] = []; for (const name of names) { const match = SCENE_KEYFRAME_RE.exec(name); if (!match) continue; const sceneIndex = Number(match[1]); if (!Number.isInteger(sceneIndex)) continue; const path = join(referencesDir, name); assets.push({ id: `image-${path}`, kind: 'image', path, sceneIndex }); } return assets; } /** * Candidate-mode `execute-status` re-derives the asset-manifest from the * SELECTED candidate outputs. Before any winner is selected the derived * manifest is empty, so a naive overwrite DELETES the per-scene INPUT image * keyframes that `buildExecutionPayload` needs to keep pending scenes on the * image-to-video path (no keyframe → text-to-video → fresh character every * scene). This preserves the existing manifest's INPUT assets (operator * image keyframe / audio / video reference assets whose ids are NOT candidate * outputs) and merges them under the freshly-derived outputs — the same * merge-preserve shape the non-candidate branch already uses via `mergeAssets`. * * THE RATCHET IT CLOSES: reconstructing inputs ONLY from the current manifest * meant that once the manifest was video-only even ONCE * (`inputs.length === 0` → the old code `return derived`), the per-scene image * keyframes could NEVER be recovered — every subsequent poll re-derived a * video-only manifest, so later `produce` runs submitted `referencePaths: []` * → text-to-video → a fresh generic character each scene (the #1 identity-loss * defect). The fix additionally re-reads the per-scene image keyframe from its * on-disk SOURCE OF TRUTH — `references/scene-keyframe.png` — so a manifest * that was ratcheted to video-only SELF-HEALS on the next poll, without the * operator manually rebuilding an all-image manifest. * * Merge order (last wins on id collision): * 1. disk keyframes (base, so genuinely-missing scenes are restored) * 2. manifest INPUT assets (win over disk keyframes → operator-registered * metadata such as `backend`/`role` is preserved verbatim) * 3. freshly-derived candidate OUTPUTS (the video/image the poll produced) * The three id namespaces never collide by scene, so a scene keeps BOTH its * image keyframe AND any derived video output. When no keyframe exists on disk * and the manifest already carries the inputs, the result is byte-identical to * the prior behavior. */ export async function preserveInputKeyframes( workspace: Awaited>, derived: AssetManifestArtifact, ): Promise { const manifestPath = artifactPathFor(workspace, 'asset-manifest'); let existing: AssetManifestArtifact['assets'] = []; if (existsSync(manifestPath)) { try { const parsed = JSON.parse(await readFile(manifestPath, 'utf-8')) as AssetManifestArtifact; existing = Array.isArray(parsed.assets) ? parsed.assets : []; } catch { // Unparseable manifest → fall through with no surviving inputs; the disk // keyframes below still self-heal the per-scene image bindings. existing = []; } } // Preserve EVERY operator INPUT asset (image keyframe, audio, or video/V2V // reference) — i.e. any asset whose id is NOT a candidate-derived output. // Narrowing to image/audio would re-introduce the same wipe bug for input // `video:` references on the first pre-selection poll; the `-output-` id // guard alone cleanly separates inputs from derived outputs. const inputs = existing.filter((asset) => !DERIVED_OUTPUT_ID_RE.test(asset.id)); // Disk source-of-truth self-heal (closes the video-only ratchet). const diskKeyframes = await collectDiskKeyframeAssets(workspace); // Base = disk keyframes; manifest inputs win over them (preserve metadata); // derived outputs applied last. Empty disk + already-present inputs → identical // to `mergeAssets(inputs, derived.assets)` (byte-identical legacy output). const withKeyframes = mergeAssets(diskKeyframes, inputs); return { projectSlug: derived.projectSlug, assets: mergeAssets(withKeyframes, derived.assets) }; } /** Immutably replace one candidate within the candidates artifact. */ function updateCandidate( artifact: SceneCandidatesArtifact, sceneIndex: number, candidateId: string, fn: (prev: SceneCandidate) => SceneCandidate, ): SceneCandidatesArtifact { return { ...artifact, scenes: artifact.scenes.map((scene) => scene.sceneIndex !== sceneIndex ? scene : { ...scene, candidates: scene.candidates.map((c) => (c.id === candidateId ? fn(c) : c)) }, ), }; } /** * Per-scene candidate poll — the auto-chain / per-scene-submit safety net. * * The single `execution-report` records only the LAST submission's job, but in * candidate mode each scene carries its OWN adapter job id on * `candidate.source.externalJobId`. When the latest report is blocked / job-less * (e.g. a later scene failed to submit and overwrote the report), the legacy * single-job poll bails with `execution-already-blocked` and strands earlier * scenes' in-flight jobs. This polls every still-`pending` candidate that has its * own job id, promotes each independently, and re-derives the asset-manifest — so * one blocked scene never hides the rest of the chain's progress. * * Returns the full refresh result, or `null` when there is no candidate artifact * or no pending candidate carries a job id (the caller then keeps the legacy * blocked-report behavior — preserving today's output for non-candidate runs). */ /** * Pure: the outputs of one poll that belong to `sceneIndex` — those tagged for * it, plus untagged ones (adapters that do not tag attach to every scene of the * run). A completed poll that carries outputs for OTHER scenes only says * nothing about this one; the caller must not mark it completed (#361: two * pool scenes shared one execution report, the poll carried scene 0's clip, * and scene 1 was written `completed` with zero outputs — a state * `waitForSceneVideo` can only time out against). */ export function pollOutputsForScene( outputs: VideoExecutionPollResult['outputs'], sceneIndex: number, ): SceneCandidateOutput[] { const own: SceneCandidateOutput[] = []; for (const out of outputs) { if (out.kind !== 'image' && out.kind !== 'video' && out.kind !== 'audio') continue; if (typeof out.sceneIndex === 'number' && out.sceneIndex !== sceneIndex) continue; own.push({ kind: out.kind, path: out.path }); } return own; } /** The issue text for a completion that carried nothing for a scene (#361). */ export function completedWithoutSceneOutputIssue(sceneIndex: number, outputs: VideoExecutionPollResult['outputs']): string { const tagged = [...new Set(outputs.map((out) => out.sceneIndex).filter((i): i is number => typeof i === 'number'))].sort((a, b) => a - b); return tagged.length > 0 ? `scene ${sceneIndex}: completed but the provider's outputs were tagged for scene${tagged.length === 1 ? '' : 's'} ${tagged.join(', ')} only — nothing for this scene (#361)` : `scene ${sceneIndex}: completed but provider returned no outputs to ingest.`; } async function pollPendingSceneCandidates(args: { workspace: Awaited>; root: string; projectSlug: string; baseReport: VideoExecutionReport; env?: NodeJS.ProcessEnv; }): Promise<{ reportPath: string; report: VideoExecutionReport; poll: VideoExecutionPollResult; assetManifestPath?: string; } | null> { const { workspace, root, projectSlug, baseReport } = args; if (!existsSync(sceneCandidatesPathFor(root, projectSlug))) return null; const candidates = await readSceneCandidatesArtifact(root, projectSlug); const pending: Array<{ sceneIndex: number; candidate: SceneCandidate }> = []; for (const scene of candidates.scenes) { for (const candidate of scene.candidates) { if (candidate.status === 'pending' && candidate.source.externalJobId) { pending.push({ sceneIndex: scene.sceneIndex, candidate }); } } } if (pending.length === 0) return null; const lastCheckedAt = new Date().toISOString(); const outputDir = `${workspace.projectDir}/outputs`; const allOutputs: VideoExecutionPollResult['outputs'] = []; const issues: string[] = []; let completed = 0; let failed = 0; let stillPending = 0; let lastJobId: string | null = null; // Collect per-scene deltas during the SLOW poll loop (kept OUTSIDE the lock). // Each delta is applied to a FRESH re-read of the on-disk artifacts under the // lock below, so a concurrent writer's update is never clobbered. const candidateDeltas: Array<{ sceneIndex: number; candidateId: string; apply: (prev: SceneCandidate) => SceneCandidate; }> = []; const selectionMarks: Array<{ markPendingScene: number; candidateId: string }> = []; const terminalLaneTickets = new Map(); const pendingLaneTickets = new Set(); for (const { sceneIndex, candidate } of pending) { const jobId = candidate.source.externalJobId as string; lastJobId = jobId; let poll: VideoExecutionPollResult; try { poll = await pollExecutionPayload({ projectSlug, routeId: candidate.route as ProviderRouteId, externalJobId: jobId, outputDir, workspaceRoot: workspace.root, }, { env: args.env }); } catch (error) { issues.push(`scene ${sceneIndex} (${candidate.id}): poll failed — ${(error as Error).message}`); if (candidate.source.laneTicketId) pendingLaneTickets.add(candidate.source.laneTicketId); stillPending += 1; continue; } allOutputs.push(...poll.outputs); if (poll.issues.length) issues.push(...poll.issues.map((issue) => `scene ${sceneIndex}: ${issue}`)); // A candidate is completed only by an output that is ITS OWN (tagged for // this scene, or untagged); a completion carrying other scenes' clips only // is a failure for this scene, never a completed-with-nothing (#361). const sceneOutputs = pollOutputsForScene(poll.outputs, sceneIndex); const completedWithoutOutputs = poll.status === 'completed' && sceneOutputs.length === 0; if (poll.status === 'completed' && !completedWithoutOutputs) { candidateDeltas.push({ sceneIndex, candidateId: candidate.id, apply: (prev) => ({ ...prev, status: 'completed', completedAt: lastCheckedAt, outputs: sceneOutputs, }), }); selectionMarks.push({ markPendingScene: sceneIndex, candidateId: candidate.id }); if (candidate.source.laneTicketId && candidate.source.laneId) { terminalLaneTickets.set(candidate.source.laneTicketId, { laneId: candidate.source.laneId, route: candidate.route, }); } completed += 1; } else if (poll.status === 'failed' || completedWithoutOutputs) { if (completedWithoutOutputs) { issues.push(completedWithoutSceneOutputIssue(sceneIndex, poll.outputs)); } candidateDeltas.push({ sceneIndex, candidateId: candidate.id, apply: (prev) => ({ ...prev, status: 'failed', completedAt: lastCheckedAt, }), }); if (candidate.source.laneTicketId && candidate.source.laneId) { terminalLaneTickets.set(candidate.source.laneTicketId, { laneId: candidate.source.laneId, route: candidate.route, }); } failed += 1; } else { if (candidate.source.laneTicketId) pendingLaneTickets.add(candidate.source.laneTicketId); stillPending += 1; } } // Apply the collected deltas to a FRESH re-read under the single artifact lock, // then derive the asset-manifest from the LOCKED result. const locked = await withSceneArtifactsLock(root, projectSlug, async () => { let cur = await readSceneCandidatesArtifact(root, projectSlug); let sel = await readSceneSelectionArtifact(root, projectSlug); for (const delta of candidateDeltas) { cur = updateCandidate(cur, delta.sceneIndex, delta.candidateId, delta.apply); } for (const mark of selectionMarks) { sel = markPending(sel, mark.markPendingScene, [mark.candidateId]); } await writeSceneCandidatesArtifact(root, projectSlug, cur); await writeSceneSelectionArtifact(root, projectSlug, sel); return { cur, sel }; }); const derived = deriveAssetManifestFromSelection(projectSlug, locked.cur, locked.sel); const merged = await preserveInputKeyframes(workspace, derived); const assetManifestPath = await writeArtifact(workspace, 'asset-manifest', merged); // Aggregate: still-pending dominates (keep polling); else completed if any // scene finished; else everything we polled failed. const status: VideoExecutionPollResult['status'] = stillPending > 0 ? 'pending' : completed > 0 ? 'completed' : 'failed'; // Release only terminal attempts' exact tickets. Never release by project: // parallel submissions for the same film may legitimately own other slots. for (const [ticketId, lane] of terminalLaneTickets) { if (pendingLaneTickets.has(ticketId)) continue; try { const laneCap = ROUTE_CAPABILITIES[lane.route as ProviderRouteId]?.maxConcurrentJobs ?? null; if (laneCap !== null && (args.env ?? process.env).VCLAW_LANE_DISABLE !== '1') { createLaneCoordinator({ env: args.env ?? process.env }) .release(ticketId, lane.laneId, laneCap); } } catch { /* lane bookkeeping must not break status */ } } const rawResult = { reason: 'per-scene-candidate-poll', polled: pending.length, completed, failed, stillPending, }; const poll: VideoExecutionPollResult = { status, externalJobId: lastJobId, outputs: allOutputs, issues, rawResult, }; const reportPath = artifactPathFor(workspace, 'execution-report'); const updatedReport: VideoExecutionReport = { ...baseReport, poll: { lastCheckedAt, status, issues, ...(status === 'completed' ? { outputsIngested: allOutputs.length } : {}), rawResult, }, }; if (status === 'completed') { await writeStageCheckpoint(workspace, { stage: 'assets', status: 'completed', generatedAt: lastCheckedAt, artifacts: { 'asset-manifest': assetManifestPath, 'execution-report': reportPath }, summary: 'Pending per-scene candidate jobs completed and outputs were ingested.', issues: [], nextAction: 'Run review on the generated outputs.', }); await updateProjectManifestState(workspace, { updatedAt: lastCheckedAt, currentStage: 'review', lastCompletedStage: 'assets', lastCheckpointStatus: 'completed', }); } else if (status === 'failed') { await writeStageCheckpoint(workspace, { stage: 'assets', status: 'failed', generatedAt: lastCheckedAt, artifacts: { 'execution-report': reportPath }, summary: 'Pending per-scene candidate jobs failed.', issues, nextAction: 'Resolve provider issues and resubmit the affected scenes.', }); await updateProjectManifestState(workspace, { updatedAt: lastCheckedAt, currentStage: 'assets', lastCompletedStage: 'storyboard', lastCheckpointStatus: 'failed', }); } else { await writeStageCheckpoint(workspace, { stage: 'assets', status: 'pending', generatedAt: lastCheckedAt, artifacts: { 'execution-report': reportPath }, summary: 'Per-scene candidate jobs are still pending.', issues, nextAction: 'Poll execution status again later.', }); await updateProjectManifestState(workspace, { updatedAt: lastCheckedAt, currentStage: 'assets', lastCompletedStage: 'storyboard', lastCheckpointStatus: 'pending', }); } await writeArtifact(workspace, 'execution-report', updatedReport); await appendProjectEvent(workspace, { type: 'execution.status.refreshed', recordedAt: lastCheckedAt, payload: { reportPath, status, externalJobId: lastJobId, outputsIngested: allOutputs.length, }, }); await appendGenerationTelemetry(workspace, buildGenerationTelemetryFromPoll({ report: updatedReport, poll, recordedAt: lastCheckedAt, })); // Repaint the live run dashboard (status badges + newly-ingested clips) so it // stays current each poll without a manual portal command. Best-effort. await regenerateRunSurface(root, projectSlug, args.env); return { reportPath, report: updatedReport, poll, assetManifestPath }; } export async function refreshExecutionStatus( projectSlug: string, options: { root?: string; productionMode?: VideoProductionMode; env?: NodeJS.ProcessEnv; } = {}, ): Promise<{ reportPath: string; report: VideoExecutionReport; poll: VideoExecutionPollResult; assetManifestPath?: string; }> { const root = options.root ?? resolveWorkspaceRootFromEnv(); const resolvedWorkspace = resolveProjectWorkspace(projectSlug, root); const projectManifest = await readProjectManifest(resolvedWorkspace); if (!projectManifest) { const now = new Date().toISOString(); const reportPath = artifactPathFor(resolvedWorkspace, 'execution-report'); const issues = [ `Execution status unavailable for ${projectSlug}: project manifest is missing. Run \`vclaw video init ${projectSlug}\` first.`, ]; const poll: VideoExecutionPollResult = { status: 'failed', externalJobId: null, outputs: [], issues, rawResult: { reason: 'missing-project-manifest', }, }; const report: VideoExecutionReport = { projectSlug, productionMode: options.productionMode ?? 'storyboard', operationKind: 'text-to-video', routeId: null, status: 'blocked', dryRun: false, generatedAt: now, blockers: issues, executedSteps: ['execution-status-requested'], taskCount: 0, poll: { lastCheckedAt: now, status: 'failed', issues, rawResult: poll.rawResult, }, }; return { reportPath, report, poll, }; } const workspace = await ensureProjectWorkspace(projectSlug, root); const status = await buildProjectStatusReport(projectSlug, root, options.productionMode ?? 'storyboard'); if (status.productionMode === 'director' && status.storyboardReviewStale) { throw new Error( `Execution status unavailable for ${projectSlug}: storyboard review is stale. Refresh ${status.storyboardReviewPath ?? 'storyboard.md'} before continuing.`, ); } const reportPath = artifactPathFor(workspace, 'execution-report'); if (!existsSync(reportPath)) { const lastCheckedAt = new Date().toISOString(); const issues = ['Execution status unavailable: execution-report artifact is missing. Run execution first.']; const poll: VideoExecutionPollResult = { status: 'failed', externalJobId: null, outputs: [], issues, rawResult: { reason: 'missing-execution-report', }, }; const updatedReport: VideoExecutionReport = { projectSlug, productionMode: status.productionMode, operationKind: 'text-to-video', routeId: null, status: 'blocked', dryRun: false, generatedAt: lastCheckedAt, blockers: issues, executedSteps: ['execution-status-requested'], taskCount: 0, poll: { lastCheckedAt, status: 'failed', issues, rawResult: poll.rawResult, }, }; await writeStageCheckpoint(workspace, { stage: 'assets', status: 'failed', generatedAt: lastCheckedAt, artifacts: { 'execution-report': reportPath, }, summary: 'Execution status refresh failed.', issues, nextAction: 'Run execution before polling execution status.', }); await updateProjectManifestState(workspace, { updatedAt: lastCheckedAt, currentStage: 'assets', lastCompletedStage: 'storyboard', lastCheckpointStatus: 'failed', }); await writeArtifact(workspace, 'execution-report', updatedReport); await appendProjectEvent(workspace, { type: 'execution.status.refreshed', recordedAt: lastCheckedAt, payload: { reportPath, status: 'failed', externalJobId: null, outputsIngested: 0, }, }); return { reportPath, report: updatedReport, poll, }; } const report = JSON.parse(await readFile(reportPath, 'utf-8')) as VideoExecutionReport; if (!report.routeId || !report.submission?.externalJobId) { // The latest report has no pollable job of its own. Before bailing, poll any // still-pending per-scene candidate jobs (auto-chain / per-scene submits each // carry their own job id) so one blocked scene doesn't strand the rest of the // chain. No candidate artifact / no pending job → null → legacy bail below. const candidateResult = await pollPendingSceneCandidates({ workspace, root, projectSlug, baseReport: report, env: options.env, }); if (candidateResult) return candidateResult; const reportHasBlockers = Array.isArray(report.blockers) && report.blockers.length > 0; // A --dry-run report is not a failed live job. Saying "failed" about one is // wrong on its own, and it is actively misleading while a real `produce` is // in flight: `produce` writes nothing until the provider answers, so for the // whole render the disk still holds the PREVIOUS dry-run report, and a status // check from a second shell read it as a dead job. Report it as a distinct // not-submitted reason and — critically — write NOTHING, so a status check // can never mutate a run's own source of truth underneath it. (This does not // yet DETECT an in-flight run; that needs a run marker `produce` writes // before submitting. Tracked separately.) if (report.dryRun === true && !reportHasBlockers) { const issues = [ 'Execution status unavailable: the last execution report is a --dry-run, not a live submission. ' + 'A `vclaw video produce` run that is still in flight has not written its report yet — re-check once it ' + 'finishes, or run `vclaw video produce` without --dry-run to create a live job.', ]; return { reportPath, report, poll: { status: 'failed', externalJobId: null, outputs: [], issues, rawResult: { reason: 'last-run-was-dry-run' }, }, }; } const lastCheckedAt = new Date().toISOString(); const issues = reportHasBlockers ? [...report.blockers] : !report.routeId ? ['Execution status unavailable: last execution report has no provider route id.'] : ['Execution status unavailable: last execution report has no live adapter job id.']; const reason = reportHasBlockers ? 'execution-already-blocked' : !report.routeId ? 'missing-provider-route-id' : 'missing-live-adapter-job-id'; const nextAction = reportHasBlockers ? 'Resolve execution blockers and rerun execution.' : 'Rerun execution to create a live adapter job id before polling status.'; const poll: VideoExecutionPollResult = { status: 'failed', externalJobId: null, outputs: [], issues, rawResult: { reason, }, }; const updatedReport: VideoExecutionReport = { ...report, poll: { lastCheckedAt, status: 'failed', issues, rawResult: poll.rawResult, }, }; await writeStageCheckpoint(workspace, { stage: 'assets', status: 'failed', generatedAt: lastCheckedAt, artifacts: { 'execution-report': reportPath, }, summary: 'Execution status refresh failed.', issues, nextAction, }); await updateProjectManifestState(workspace, { updatedAt: lastCheckedAt, currentStage: 'assets', lastCompletedStage: 'storyboard', lastCheckpointStatus: 'failed', }); await writeArtifact(workspace, 'execution-report', updatedReport); await appendProjectEvent(workspace, { type: 'execution.status.refreshed', recordedAt: lastCheckedAt, payload: { reportPath, status: 'failed', externalJobId: null, outputsIngested: 0, }, }); return { reportPath, report: updatedReport, poll, }; } // Concurrent submits (the render pool, or any parallel `execute --scene`) put // MORE THAN ONE job in flight, but the single execution-report tracks only the // LATEST job (`report.submission`). Polling only that job orphans earlier // still-pending candidates — their finished scenes never download. When a // pending candidate carries a job id OTHER than the latest, poll EVERY pending // candidate's own job instead (this still covers the latest job too). The // single-job fast path below is unchanged when the only pending candidates // belong to the latest job — so the common one-job case is byte-identical. if (existsSync(sceneCandidatesPathFor(root, projectSlug))) { const pendingJobs = new Set(); for (const scene of (await readSceneCandidatesArtifact(root, projectSlug)).scenes) { for (const candidate of scene.candidates) { if (candidate.status === 'pending' && candidate.source.externalJobId) { pendingJobs.add(candidate.source.externalJobId); } } } const latestJobId = report.submission.externalJobId; const hasOrphanedJob = [...pendingJobs].some((jobId) => jobId !== latestJobId); if (hasOrphanedJob) { const multiJobResult = await pollPendingSceneCandidates({ workspace, root, projectSlug, baseReport: report, env: options.env, }); if (multiJobResult) return multiJobResult; } } const poll = await pollExecutionPayload({ projectSlug, routeId: report.routeId, externalJobId: report.submission.externalJobId, outputDir: `${workspace.projectDir}/outputs`, workspaceRoot: workspace.root, }, { env: options.env, }); const completedWithoutOutputs = poll.status === 'completed' && poll.outputs.length === 0; const normalizedPollStatus: VideoExecutionPollResult['status'] = completedWithoutOutputs ? 'failed' : poll.status; const normalizedIssues = completedWithoutOutputs ? [...poll.issues, 'Execution completed but provider returned no outputs to ingest.'] : poll.issues; const pollMetadata = { lastCheckedAt: new Date().toISOString(), status: normalizedPollStatus, issues: normalizedIssues, ...(normalizedPollStatus === 'completed' ? { outputsIngested: poll.outputs.length } : {}), rawResult: poll.rawResult, }; const updatedReport: VideoExecutionReport = { ...report, poll: pollMetadata, }; const lastCheckedAt = pollMetadata.lastCheckedAt; // Candidate-mode detection mirrors executeProject: if a candidate artifact // exists, or the last report carries `candidatesByScene`, treat this poll as // a candidate-mode poll and route output ingestion through the candidate // store. Otherwise we keep the legacy direct-asset-manifest behavior. const candidateMode = existsSync(sceneCandidatesPathFor(root, projectSlug)) || Array.isArray(report.candidatesByScene); let assetManifestPath: string | undefined; if (normalizedPollStatus === 'completed') { if (candidateMode) { // Candidate path — update per-scene candidates, append to // pendingCandidateIds, then re-derive asset-manifest from selection. // Group poll outputs by sceneIndex (pure, on already-fetched poll data — // safe OUTSIDE the lock). Outputs without a sceneIndex are attached to // every candidate we created for this run (preserves the legacy fallback // when adapters don't tag outputs). const candidateIdsByScene = new Map(); for (const entry of report.candidatesByScene ?? []) { candidateIdsByScene.set(entry.sceneIndex, entry.candidateId); } // Apply the completion updates to a FRESH re-read under the single artifact // lock, so a concurrent pool scene's candidate/selection write is not lost. const locked = await withSceneArtifactsLock(root, projectSlug, async () => { let updatedCandidates = await readSceneCandidatesArtifact(root, projectSlug); let updatedSelection = await readSceneSelectionArtifact(root, projectSlug); // This poll is about ONE job. Every candidate still pending under that // job is judged by it — not only the ones the latest report lists, // because a later per-scene submit rewrites the single shared report // and drops its predecessors' entries (#361: scene 0 stayed pending // forever while scene 1 was marked completed with nothing). for (const scene of updatedCandidates.scenes) { for (const candidate of scene.candidates) { if (candidate.status === 'pending' && candidate.source.externalJobId === report.submission?.externalJobId && !candidateIdsByScene.has(scene.sceneIndex)) { candidateIdsByScene.set(scene.sceneIndex, candidate.id); } } } for (const [sceneIndex, candidateId] of candidateIdsByScene) { // Same rule as the per-candidate poll: only an output of this scene's // own completes it. The job finished, and it carried nothing for this // scene — that is a failed scene, not a completed one (#361). const sceneOutputs = pollOutputsForScene(poll.outputs, sceneIndex); if (sceneOutputs.length === 0) { normalizedIssues.push(completedWithoutSceneOutputIssue(sceneIndex, poll.outputs)); updatedCandidates = updateCandidate(updatedCandidates, sceneIndex, candidateId, (prev) => ({ ...prev, status: 'failed', completedAt: lastCheckedAt, })); continue; } updatedCandidates = updateCandidate(updatedCandidates, sceneIndex, candidateId, (prev) => ({ ...prev, status: 'completed', completedAt: lastCheckedAt, outputs: sceneOutputs, })); updatedSelection = markPending(updatedSelection, sceneIndex, [candidateId]); } await writeSceneCandidatesArtifact(root, projectSlug, updatedCandidates); await writeSceneSelectionArtifact(root, projectSlug, updatedSelection); return { updatedCandidates, updatedSelection }; }); // Derive asset-manifest from selection so the legacy review/publish // readers keep seeing a coherent manifest. Before any operator has // selected a winner the derived `assets` is empty, so preserve the // existing manifest's INPUT image/audio keyframes (merge-preserve) — // otherwise the first post-`produce` poll would wipe the per-scene // keyframes that keep pending scenes on the image-to-video path. const derived = deriveAssetManifestFromSelection(projectSlug, locked.updatedCandidates, locked.updatedSelection); const merged = await preserveInputKeyframes(workspace, derived); assetManifestPath = await writeArtifact(workspace, 'asset-manifest', merged); } else { const existingAssetManifest = existsSync(artifactPathFor(workspace, 'asset-manifest')) ? JSON.parse(await readFile(artifactPathFor(workspace, 'asset-manifest'), 'utf-8')) as AssetManifestArtifact : { projectSlug, assets: [] }; const nextAssetManifest: AssetManifestArtifact = { projectSlug, assets: mergeAssets(existingAssetManifest.assets ?? [], poll.outputs), }; assetManifestPath = await writeArtifact(workspace, 'asset-manifest', nextAssetManifest); } await writeStageCheckpoint(workspace, { stage: 'assets', status: 'completed', generatedAt: lastCheckedAt, artifacts: { 'asset-manifest': assetManifestPath, 'execution-report': reportPath, }, summary: 'Live execution completed and outputs were ingested.', issues: [], nextAction: 'Run review on the generated outputs.', }); await updateProjectManifestState(workspace, { updatedAt: lastCheckedAt, currentStage: 'review', lastCompletedStage: 'assets', lastCheckpointStatus: 'completed', }); } else if (normalizedPollStatus === 'failed') { await writeStageCheckpoint(workspace, { stage: 'assets', status: 'failed', generatedAt: lastCheckedAt, artifacts: { 'execution-report': reportPath, }, summary: 'Live execution failed.', issues: normalizedIssues, nextAction: 'Resolve provider issues and resubmit execution.', }); await updateProjectManifestState(workspace, { updatedAt: lastCheckedAt, currentStage: 'assets', lastCompletedStage: 'storyboard', lastCheckpointStatus: 'failed', }); } else { await writeStageCheckpoint(workspace, { stage: 'assets', status: 'pending', generatedAt: lastCheckedAt, artifacts: { 'execution-report': reportPath, }, summary: 'Live execution is still pending.', issues: normalizedIssues, nextAction: 'Poll execution status again later.', }); await updateProjectManifestState(workspace, { updatedAt: lastCheckedAt, currentStage: 'assets', lastCompletedStage: 'storyboard', lastCheckpointStatus: 'pending', }); } await writeArtifact(workspace, 'execution-report', updatedReport); await appendProjectEvent(workspace, { type: 'execution.status.refreshed', recordedAt: lastCheckedAt, payload: { reportPath, status: normalizedPollStatus, externalJobId: poll.externalJobId, outputsIngested: poll.outputs.length, }, }); await appendGenerationTelemetry(workspace, buildGenerationTelemetryFromPoll({ report: updatedReport, poll: { ...poll, status: normalizedPollStatus, issues: normalizedIssues, }, recordedAt: lastCheckedAt, })); // A terminal provider job owns one exact lane ticket captured at submission. // Release that ticket only; older reports without the field safely expire by // TTL instead of risking a broad project-level release. if (normalizedPollStatus !== 'pending' && report.laneLease && (options.env ?? process.env).VCLAW_LANE_DISABLE !== '1') { try { createLaneCoordinator({ env: options.env ?? process.env }).release( report.laneLease.ticketId, report.laneLease.laneId, report.laneLease.limit, ); } catch { /* lane bookkeeping must not break status */ } } // Repaint the live run dashboard (status badges + newly-ingested clips) so it // stays current each poll without a manual portal command. Best-effort. await regenerateRunSurface(root, projectSlug, options.env); return { reportPath, report: updatedReport, poll, ...(assetManifestPath ? { assetManifestPath } : {}), }; }