import { readFile } from 'node:fs/promises'; import { existsSync } from 'node:fs'; import { join } from 'node:path'; import { artifactPathFor } from './artifact-store.js'; import { VclawError } from './errors.js'; import { appendProjectEvent } from './events.js'; import { resolveActiveTransport } from './execution-adapter.js'; import { describeExecutionRuns } from './execution-run-marker.js'; import { refreshExecutionStatus } from './execution-status.js'; import { createLaneCoordinator } from './lane-coordinator.js'; import { abandonDreaminaUseApiNativeJob } from './native-dreamina.js'; import { abandonSeedanceModelArkNativeJob } from './native-modelark.js'; import { abandonReapiSeedanceNativeJob } from './native-reapi.js'; import { abandonRunwayUseApiNativeJob } from './native-runway.js'; import { ROUTE_CAPABILITIES } from './provider-platform/route-capabilities.js'; import type { ProviderRouteId } from './provider-platform/types.js'; import { readSceneCandidatesArtifact, sceneCandidatesPathFor, updateSceneCandidatesArtifact } from './scene-candidate-store.js'; import { resolveWorkspaceRootFromEnv } from './workspace-root.js'; import { ensureProjectWorkspace, readProjectManifest, resolveProjectWorkspace } from './workspace.js'; import type { VideoExecutionPollResult, VideoExecutionReport, VideoProductionMode } from './types.js'; /** The transports whose poll ends only on a provider-reported state, and which have no provider-side cancel. */ const ABANDON = { 'native-runway': abandonRunwayUseApiNativeJob, 'native-dreamina': abandonDreaminaUseApiNativeJob, // ModelArk can cancel a QUEUED task but not a running one, so a wedged // running task still needs a way to stop being waited for. 'native-modelark': abandonSeedanceModelArkNativeJob, 'native-reapi': abandonReapiSeedanceNativeJob, } as const; /** Why every other transport is refused, in the words the refusal prints. */ const NOT_NEEDED: Record = { 'native-veo': 'the Flow poll completes from local job state and never waits on a wedged provider job', 'native-seedance': 'this transport has a real provider-side cancel: use `vclaw video execute-cancel`', 'in-tree-engine': 'the free Seedance engine owns its own job state', 'custom-adapter': 'a custom adapter owns its own job state', 'command-shim': 'a command shim owns its own job state', 'native-magnific': 'a Magnific job is short and completes on its own', blocked: 'no transport is active for this route', }; /** * Every provider job this project is still waiting for, job id -> its route: * the last report's job AND every per-scene job whose candidate is still * `pending`. `execute-abandon` and `execute-bind` act on the same set, so they * read it from one place — two copies had already begun to drift. */ export function jobsWaitedFor( report: VideoExecutionReport | null, candidates: { scenes: Array<{ candidates: Array<{ status: string; route: string; source: { externalJobId?: string } }> }> } | null, ): Map { const waitedFor = new Map(); if (report?.routeId && report.submission?.externalJobId) waitedFor.set(report.submission.externalJobId, report.routeId); for (const scene of candidates?.scenes ?? []) { for (const candidate of scene.candidates) { if (candidate.status === 'pending' && candidate.source.externalJobId) waitedFor.set(candidate.source.externalJobId, candidate.route); } } return waitedFor; } export interface AbandonedJob { externalJobId: string; routeId: string; abandonedScenes: number[]; untouchedScenes: Array<{ sceneIndex: number; status: string }>; /** Candidates that were `pending` on this job and are now `failed`. */ candidateIds: string[]; } /** * Stop WAITING on jobs the provider has left in flight. `execute-cancel` is * honest about routes with no provider-side cancel (it answers `unsupported` and * changes nothing), which left no way to end a poll that only ends when the * provider says so. This is that lever under its real name: abandoning makes no * provider call, the job keeps running and keeps billing, its clip is never * collected, and it cannot be undone. Without `confirm` it only reports the plan. * * Which jobs: `jobIds` when given; otherwise the last report's job AND every * candidate still `pending` on one of these routes. A per-scene `produce --scene * N` is its own job, and a later submit overwrites the single report, so the * wedged job is usually one only the candidate list still remembers. * * A whole job is abandoned, never some of its scenes: the poll reports `failed` * as soon as one scene has, so a sibling scene that finished later would be * downloaded and never ingested. */ export async function abandonExecution( projectSlug: string, options: { root?: string; productionMode?: VideoProductionMode; env?: NodeJS.ProcessEnv; /** Abandon exactly these provider jobs; default is every job still being waited for. */ jobIds?: string[]; /** false/omitted = report the plan and change nothing. */ confirm?: boolean; /** Injectable status poll (tests). Default: the real `refreshExecutionStatus`. */ refreshStatus?: typeof refreshExecutionStatus; } = {}, ): Promise<{ reportPath: string; report: VideoExecutionReport | null; /** Present once the abandon happened: the status poll that followed it. */ poll?: VideoExecutionPollResult; /** Set when the abandon happened but the status poll after it could not run (for instance a stale director review). The abandon stands. */ statusPollError?: string; abandonment: { /** false = nothing was changed; pass --confirm-abandon to do it. */ applied: boolean; jobs: AbandonedJob[]; /** Jobs a sweep left alone because their saved state could not be read, each with why; the others were still abandoned. */ notAbandoned?: string[]; providerJobsStillRunning: true; canBeUndone: false; }; }> { const root = options.root ?? resolveWorkspaceRootFromEnv(); const env = options.env ?? process.env; if (!(await readProjectManifest(resolveProjectWorkspace(projectSlug, root)))) { throw new VclawError('asset_not_found', `Execution abandon unavailable for ${projectSlug}: project manifest is missing.`, { projectSlug }); } const workspace = await ensureProjectWorkspace(projectSlug, root); const outputDir = join(workspace.projectDir, 'outputs'); // A run still inside its submit window has no job id yet; abandoning now would // act on an EARLIER job while this one goes on to submit. const liveRuns = (await describeExecutionRuns(workspace).catch(() => [])).filter((run) => run.description.state === 'live'); if (liveRuns.length > 0) { throw new VclawError( 'execution_blocked_by_readiness', `A \`vclaw video produce\` run (pid ${liveRuns[0]!.marker.pid}, started ${liveRuns[0]!.marker.startedAt}) is still submitting and has no provider job id yet. Wait for it to finish, then abandon.`, { pid: liveRuns[0]!.marker.pid }, ); } const reportPath = artifactPathFor(workspace, 'execution-report'); const report = existsSync(reportPath) ? JSON.parse(await readFile(reportPath, 'utf-8')) as VideoExecutionReport : null; const hasCandidates = existsSync(sceneCandidatesPathFor(root, projectSlug)); const candidates = hasCandidates ? await readSceneCandidatesArtifact(root, projectSlug) : null; const waitedFor = jobsWaitedFor(report, candidates); if (waitedFor.size === 0) { throw new VclawError('execution_blocked_by_readiness', `Nothing to abandon for ${projectSlug}: no provider job is being waited for (no live report job and no pending candidate).`, { projectSlug }); } const requested = options.jobIds && options.jobIds.length > 0 ? options.jobIds : [...waitedFor.keys()]; const explicit = Boolean(options.jobIds?.length); const jobs: AbandonedJob[] = []; const refusals: string[] = []; const unreadable: string[] = []; for (const externalJobId of requested) { const routeId = waitedFor.get(externalJobId); if (!routeId) { refusals.push(`${externalJobId}: not a job this project is waiting for (known: ${[...waitedFor.keys()].join(', ')})`); continue; } const resolvedTransport = resolveActiveTransport(routeId as ProviderRouteId, env); // reapi-seedance reads `blocked` whenever VCLAW_REAPI_SEEDANCE_VIA is unset; // that gate exists for a SUBMIT. An abandon touches only the saved job // state, so the selector has no say in it. const transport = routeId === 'reapi-seedance' && resolvedTransport === 'blocked' ? 'native-reapi' : resolvedTransport; const abandon = ABANDON[transport as keyof typeof ABANDON]; if (!abandon) { refusals.push(`${externalJobId} on ${routeId} (${transport}): ${NOT_NEEDED[transport] ?? 'this transport is not one that waits on a provider-reported state'}`); continue; } if (!existsSync(join(outputDir, '.vclaw-jobs', `${externalJobId}.json`))) { refusals.push(`${externalJobId} on ${routeId}: no local job state, so nothing here is waiting on it`); continue; } // One job whose saved state cannot be read must not stop a sweep from // abandoning the others; it is refused and named, and its file untouched. let result: Awaited>; try { result = await abandon({ outputDir, externalJobId, dryRun: !options.confirm }); } catch (error) { const refusal = `${externalJobId} on ${routeId}: ${error instanceof Error ? error.message : String(error)}`; refusals.push(refusal); unreadable.push(refusal); continue; } if (result.abandonedScenes.length === 0) { refusals.push(`${externalJobId} on ${routeId}: no scene is still in flight`); continue; } jobs.push({ externalJobId, routeId, ...result, candidateIds: [] }); } // An explicit request is all-or-nothing about being understood; a sweep just // needs something to act on. if (jobs.length === 0 || (explicit && refusals.length > 0)) { throw new VclawError( 'execution_blocked_by_readiness', `video execute-abandon did nothing: ${refusals.join('; ')}.`, { refusals, appliesTo: Object.keys(ABANDON) }, ); } const pendingOn = (externalJobId: string) => (candidates?.scenes ?? []).flatMap((scene) => scene.candidates) .filter((candidate) => candidate.status === 'pending' && candidate.source.externalJobId === externalJobId); for (const job of jobs) job.candidateIds = pendingOn(job.externalJobId).map((candidate) => candidate.id); const abandonment = { applied: Boolean(options.confirm), jobs, ...(unreadable.length > 0 ? { notAbandoned: unreadable } : {}), providerJobsStillRunning: true as const, canBeUndone: false as const, }; if (!options.confirm) return { reportPath, report, abandonment }; // The review surface reads the CANDIDATE, not the job state: left `pending` it // would go on saying "rendering" about a job nobody is waiting for. const abandonedJobIds = new Set(jobs.map((job) => job.externalJobId)); const laneTickets: Array<{ ticketId: string; laneId: string; route: string }> = []; if (hasCandidates) { const completedAt = new Date().toISOString(); await updateSceneCandidatesArtifact(root, projectSlug, (artifact) => { for (const scene of artifact.scenes) { scene.candidates = scene.candidates.map((candidate) => { if (candidate.status !== 'pending' || !candidate.source.externalJobId || !abandonedJobIds.has(candidate.source.externalJobId)) return candidate; if (candidate.source.laneTicketId && candidate.source.laneId) { laneTickets.push({ ticketId: candidate.source.laneTicketId, laneId: candidate.source.laneId, route: candidate.route }); } return { ...candidate, status: 'failed' as const, completedAt }; }); } return { artifact, result: undefined }; }); } // Our own queue slot, given back the way a terminal poll gives it back. The // PROVIDER's slot is still taken by the running job; nothing here can free that. for (const ticket of laneTickets) { try { const laneCap = ROUTE_CAPABILITIES[ticket.route as ProviderRouteId]?.maxConcurrentJobs ?? null; if (laneCap !== null && env.VCLAW_LANE_DISABLE !== '1') createLaneCoordinator({ env }).release(ticket.ticketId, ticket.laneId, laneCap); } catch { /* queue bookkeeping must not undo an abandon that already happened */ } } await appendProjectEvent(workspace, { type: 'execution.abandoned', payload: { jobs: jobs.map((job) => ({ externalJobId: job.externalJobId, routeId: job.routeId, abandonedScenes: job.abandonedScenes, candidateIds: job.candidateIds })), providerJobsStillRunning: true, }, }); // The report's own job, if it was one of them, is settled by the ordinary // status poll: the report's poll block, the checkpoint and its queue slot. if (report?.submission?.externalJobId && abandonedJobIds.has(report.submission.externalJobId)) { // The abandon has already happened and cannot be undone, so a status poll // that refuses to run (a stale director review, an unreadable report) must // not turn it into an error whose retry says "nothing to abandon". try { const refreshed = await (options.refreshStatus ?? refreshExecutionStatus)(projectSlug, { root, ...(options.productionMode ? { productionMode: options.productionMode } : {}), env, }); return { reportPath: refreshed.reportPath, report: refreshed.report, poll: refreshed.poll, abandonment }; } catch (error) { return { reportPath, report, abandonment, statusPollError: `The jobs were abandoned, but the status poll that records it in the report did not run: ${error instanceof Error ? error.message : String(error)} Run \`vclaw video execute-status\` once that is resolved.`, }; } } return { reportPath, report, abandonment }; }