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 { jobsWaitedFor } from './execution-abandon.js'; import { resolveActiveTransport } from './execution-adapter.js'; import { describeExecutionRuns } from './execution-run-marker.js'; import { refreshExecutionStatus } from './execution-status.js'; import { bindSeedanceModelArkLostCreate, type ModelArkBindResult } from './native-modelark.js'; import type { ProviderRouteId } from './provider-platform/types.js'; import { readSceneCandidatesArtifact, sceneCandidatesPathFor } 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 that can bind a lost create to a named task. Only ModelArk * can: binding is only honest when the task can be READ back before anything * is written (same model, not already some job's), and ModelArk is the one * route whose vendor exposes a task by id AND records a create intent. */ const BIND = { 'native-modelark': bindSeedanceModelArkLostCreate, } as const; /** Why every other transport is refused, in the words the refusal prints. */ const NOT_BINDABLE: Record = { 'native-reapi': 'reAPI publishes no endpoint that reads a task by id, so a named task cannot be checked before it is bound; use `vclaw video execute-abandon`', 'native-runway': 'this transport records no create intent, so it has no lost create to bind', 'native-dreamina': 'this transport records no create intent, so it has no lost create to bind', 'native-veo': 'the Flow poll completes from local job state and never loses a create', 'native-seedance': 'this transport records no create intent, so it has no lost create to bind', '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', }; /** * Bind a scene whose create answer was lost to the ModelArk task the operator * names. The poll binds a lost create by itself when exactly one unbound task * sits in its window; this is the exit for the cases it will not decide — * several tasks in the window, a window it could not search, or no intent time * to search by — where today's only answer is `execute-abandon`. * * Nothing is submitted and nothing is re-submitted: the task is read with one * GET, and only a task that exists, is on this job's model and is not already * owned by a job in this output directory is bound. The scene then polls as an * ordinary submitted one, so its clip and its bill are collected as usual. * Without `confirm` it reports the plan — including what ModelArk says about * the task — and changes nothing. */ export async function bindExecutionTask( projectSlug: string, options: { root?: string; productionMode?: VideoProductionMode; env?: NodeJS.ProcessEnv; /** The ModelArk task id to bind. */ taskId: string; /** The job that holds the lost create; default is the one job being waited for on a bindable route. */ jobId?: string; /** Which scene, when the job holds more than one lost create. */ sceneIndex?: number; /** false/omitted = report the plan and change nothing. */ confirm?: boolean; /** Injectable status poll (tests). Default: the real `refreshExecutionStatus`. */ refreshStatus?: typeof refreshExecutionStatus; /** Injectable transport options (tests): a scripted fetch, a retry policy. */ transportOptions?: Parameters[1]; }, ): Promise<{ reportPath: string; report: VideoExecutionReport | null; /** Present once the bind happened: the status poll that followed it. */ poll?: VideoExecutionPollResult; /** Set when the bind happened but the status poll after it could not run. The bind stands. */ statusPollError?: string; binding: ModelArkBindResult & { routeId: string }; }> { const root = options.root ?? resolveWorkspaceRootFromEnv(); const env = options.env ?? process.env; const taskId = options.taskId.trim(); if (!taskId) { throw new VclawError('missing_required_flag', 'video execute-bind requires --task : the id of the provider task to bind.', { missing: ['--task'] }); } if (!(await readProjectManifest(resolveProjectWorkspace(projectSlug, root)))) { throw new VclawError('asset_not_found', `Execution bind 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 may be about to record this very // scene's task itself; binding now would race its own job state. 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. Wait for it to finish, then bind.`, { 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 candidates = existsSync(sceneCandidatesPathFor(root, projectSlug)) ? await readSceneCandidatesArtifact(root, projectSlug) : null; const waitedFor = jobsWaitedFor(report, candidates); if (waitedFor.size === 0) { throw new VclawError('execution_blocked_by_readiness', `Nothing to bind for ${projectSlug}: no provider job is being waited for (no live report job and no pending candidate).`, { projectSlug }); } const transportFor = (routeId: string) => resolveActiveTransport(routeId as ProviderRouteId, env); const hasJobState = (externalJobId: string) => existsSync(join(outputDir, '.vclaw-jobs', `${externalJobId}.json`)); // A candidate left `pending` by a run that never wrote job state is a ghost: // it is nothing to bind to, and counting it would ask for a --job that has // only one real answer. const onBindableRoute = [...waitedFor.entries()].filter(([, routeId]) => BIND[transportFor(routeId) as keyof typeof BIND]); const bindable = onBindableRoute.filter(([externalJobId]) => hasJobState(externalJobId)); const noJobState = onBindableRoute.filter(([externalJobId]) => !hasJobState(externalJobId)); let externalJobId: string; if (options.jobId) { externalJobId = options.jobId; const routeId = waitedFor.get(externalJobId); if (!routeId) { throw new VclawError('execution_blocked_by_readiness', `${externalJobId}: not a job ${projectSlug} is waiting for (known: ${[...waitedFor.keys()].join(', ')}). Nothing was bound.`, { externalJobId }); } const transport = transportFor(routeId); if (!BIND[transport as keyof typeof BIND]) { throw new VclawError( 'execution_blocked_by_readiness', `${externalJobId} on ${routeId} (${transport}): ${NOT_BINDABLE[transport] ?? 'this transport cannot bind a task by id'}. Nothing was bound.`, { externalJobId, routeId }, ); } } else if (bindable.length === 1) { externalJobId = bindable[0]![0]; } else { throw new VclawError( 'execution_blocked_by_readiness', bindable.length > 0 ? `${projectSlug} is waiting for ${bindable.length} jobs that could take a task (${bindable.map(([job, routeId]) => `${job} on ${routeId}`).join('; ')}); name one with --job . Nothing was bound.` : noJobState.length > 0 ? `Nothing to bind for ${projectSlug}: ${noJobState.map(([job, routeId]) => `${job} on ${routeId}`).join('; ')} is on a route that can bind a task by id, but has no local job state, so there is no lost create here to bind.` : `Nothing to bind for ${projectSlug}: no job it is waiting for is on a route that can bind a task by id (${[...waitedFor.entries()].map(([job, routeId]) => `${job} on ${routeId}`).join('; ')}).`, { projectSlug }, ); } const routeId = waitedFor.get(externalJobId)!; if (!hasJobState(externalJobId)) { throw new VclawError('asset_not_found', `${externalJobId} on ${routeId}: no local job state, so there is no lost create here to bind. Nothing was bound.`, { externalJobId }); } const bind = BIND[transportFor(routeId) as keyof typeof BIND]!; let binding: ModelArkBindResult; try { binding = await bind({ outputDir, externalJobId, workspaceRoot: root, taskId, ...(options.sceneIndex !== undefined ? { sceneIndex: options.sceneIndex } : {}), dryRun: !options.confirm, }, { env, ...(options.transportOptions ?? {}) }); } catch (error) { // The transport's refusals are the operator's to act on (a task another // job owns, a task on another model, an occupied output path), so they // exit as a gate, not as an internal error. throw new VclawError('execution_blocked_by_readiness', error instanceof Error ? error.message : String(error), { externalJobId, routeId, taskId }); } const result = { ...binding, routeId }; if (!options.confirm) return { reportPath, report, binding: result }; await appendProjectEvent(workspace, { type: 'execution.task.bound', payload: { externalJobId, routeId, sceneIndex: binding.sceneIndex, taskId: binding.taskId, createdInsideWindow: binding.createdInsideWindow }, }); // The bind has already happened and the scene is now an ordinary submitted // one, so a status poll that refuses to run (a stale director review, an // unreadable report) must not read as "nothing was bound". 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, binding: result }; } catch (error) { return { reportPath, report, binding: result, statusPollError: `Scene ${binding.sceneIndex} is bound to ModelArk task ${binding.taskId}, but the status poll that would read it did not run: ${error instanceof Error ? error.message : String(error)}. Run \`vclaw video execute-status\` once that is resolved.`, }; } }