import { join } from 'node:path'; import { beginCinemaTaskSubmission, claimCinemaTask, completeCinemaTaskDownload, handoffCinemaTaskLease, markCinemaTaskSubmissionUnknown, readCinemaProductionQueue, reconcileCinemaProviderTask, recordCinemaTaskDownloadStarted, recordCinemaTaskSubmitted, type CinemaLaneQueuePort, } from './cinema-production-queue.js'; import type { CinemaQueueMoney, CinemaQueueTask } from './cinema-queue-types.js'; import type { CinemaSpendAuthorizationPort } from './cinema-provider-types.js'; import { appendCinemaAttemptOutcome, appendCinemaGenerationAttempt, } from './cinema-store.js'; import { sha256Text, stableCinemaJson } from './cinema-shot-compiler.js'; import { appendProjectEvent } from './events.js'; import { createLaneCoordinator } from './lane-coordinator.js'; import { readSceneCandidatesArtifact, withSceneArtifactsLock, writeSceneCandidatesArtifact } from './scene-candidate-store.js'; import { appendCandidate, maxRoundForScene, nextCandidateId } from './scene-candidates.js'; import { markPending } from './scene-selection.js'; import { readSceneSelectionArtifact, writeSceneSelectionArtifact } from './scene-selection-store.js'; import { PROVIDER_ROUTE_IDS, type ProviderRouteId } from './provider-platform/types.js'; import type { SceneCandidate, VideoExecutionPayload, VideoExecutionTask } from './types.js'; import { ensureProjectWorkspace, resolveProjectWorkspace } from './workspace.js'; import { lookupExecutionSubmission, pollExecutionPayload, submitExecutionPayload } from './execution-adapter.js'; import { advanceCinemaChainFallback, type CinemaChainFallbackTransition } from './cinema-chain-fallback.js'; // Every live core route — derived, not re-listed, so a new provider route id is // recognised here the moment it lands in PROVIDER_ROUTE_IDS. const PROVIDER_ROUTES = new Set(PROVIDER_ROUTE_IDS); export interface CinemaCompatibilityWorkerObservation { state: 'not-found' | 'accepted' | 'processing' | 'completed' | 'failed' | 'unknown'; authoritative: boolean; providerJobId: string | null; providerStatus: string; envelope: Record; outputs: Array<{ id: string; kind: 'image' | 'video' | 'audio' | 'subtitle' | 'other'; path: string; sceneIndex?: number; backend?: string }>; actualCost: CinemaQueueMoney | null; failure?: { code: string; message: string; retryable: boolean }; } export interface CinemaCompatibilityWorkerTransport { submit(payload: VideoExecutionPayload): Promise<{ providerJobId: string; providerStatus: string; envelope: Record }>; lookup(input: { payload: VideoExecutionPayload; providerJobId: string | null; submissionKey: string }): Promise; } export interface CinemaCompatibilityWorkerResult { disposition: 'submitted' | 'pending' | 'completed' | 'waiting-lane' | 'already-succeeded'; task: CinemaQueueTask; candidate: SceneCandidate | null; providerCalls: number; generationCalls: number; spendAuthorized: boolean; fallback?: CinemaChainFallbackTransition; } async function lookupCompatibilityAdapter( payload: VideoExecutionPayload, providerJobId: string | null, submissionKey: string, env: NodeJS.ProcessEnv, ): Promise> { if (providerJobId) { const result = await pollExecutionPayload({ projectSlug: payload.projectSlug, routeId: payload.routeId, externalJobId: providerJobId, outputDir: payload.outputDir, workspaceRoot: payload.workspaceRoot, }, { env }); return { state: result.status === 'completed' ? 'completed' : result.status === 'failed' ? 'failed' : 'processing', authoritative: true, providerJobId: result.externalJobId ?? providerJobId, providerStatus: result.status, envelope: { status: result.status, externalJobId: result.externalJobId ?? providerJobId, outputCount: result.outputs.length, issues: result.issues, rawResult: result.rawResult }, outputs: result.outputs, ...(result.status === 'failed' ? { failure: { code: 'provider-failed', message: result.issues.join('; ') || 'Provider reported failure', retryable: false } } : {}), }; } const result = await lookupExecutionSubmission({ payload, submissionKey }, { env }); return { state: result.state, authoritative: result.authoritative, providerJobId: result.externalJobId, providerStatus: result.providerStatus, envelope: { state: result.state, authoritative: result.authoritative, externalJobId: result.externalJobId, providerStatus: result.providerStatus, outputCount: result.outputs.length, issues: result.issues, rawResult: result.rawResult, }, outputs: result.outputs, ...(result.state === 'failed' ? { failure: { code: 'provider-failed', message: result.issues.join('; ') || 'Provider reported failure', retryable: false } } : {}), }; } export function createCinemaCompatibilityAdapterTransport( task: CinemaQueueTask, env: NodeJS.ProcessEnv = process.env, ): CinemaCompatibilityWorkerTransport { if (task.authorizationRequirement !== 'none' || task.spendAuthorizationId) { throw new Error('Generic compatibility adapter transport cannot execute paid work; bind and verify an exact provider quote instead'); } if (task.routeId !== 'runway-useapi') { throw new Error(`Generic provider-free compatibility execution is limited to runway-useapi explore mode; ${task.routeId} requires an exact provider quote`); } const runwayMode = (env.VCLAW_RUNWAY_MODE ?? '').trim().toLowerCase(); if (runwayMode === 'credits' || runwayMode === 'credit') { throw new Error('Provider-free compatibility execution refuses VCLAW_RUNWAY_MODE=credits; quote and authorize the paid lane instead'); } if (env.VCLAW_RUNWAY_USEAPI_ADAPTER?.trim() || env.VCLAW_RUNWAY_USEAPI_SUBMIT_CMD?.trim()) { throw new Error('Provider-free compatibility execution refuses custom Runway submit adapters because their billing mode cannot be verified'); } return { async submit(payload) { const result = await submitExecutionPayload(payload, { env }); return { providerJobId: result.externalJobId!, providerStatus: 'submitted', envelope: { status: 'submitted', externalJobId: result.externalJobId }, }; }, async lookup({ payload, providerJobId, submissionKey }) { const result = await lookupCompatibilityAdapter(payload, providerJobId, submissionKey, env); return { ...result, outputs: result.outputs, actualCost: result.state === 'completed' ? { currency: task.estimatedCost.currency, amount: 0 } : null, }; }, }; } function paidActualCost(value: unknown, currency: string): CinemaQueueMoney | null { if (!value || typeof value !== 'object' || Array.isArray(value)) return null; const raw = (value as Record).actualCost; if (!raw || typeof raw !== 'object' || Array.isArray(raw)) return null; const cost = raw as Record; if (cost.currency !== currency || typeof cost.amount !== 'number' || !Number.isFinite(cost.amount) || cost.amount < 0) { throw new Error(`Paid compatibility adapter returned invalid actualCost; expected finite non-negative ${currency}`); } return { currency, amount: cost.amount }; } export function createCinemaCompatibilityPaidAdapterTransport( task: CinemaQueueTask, env: NodeJS.ProcessEnv = process.env, ): CinemaCompatibilityWorkerTransport { if (task.authorizationRequirement !== 'exact-quote' || !task.spendAuthorizationId) { throw new Error('Paid compatibility adapter requires a bound exact spend authorization'); } return { async submit(payload) { const result = await submitExecutionPayload(payload, { env }); return { providerJobId: result.externalJobId!, providerStatus: 'submitted', envelope: { status: 'submitted', externalJobId: result.externalJobId }, }; }, async lookup({ payload, providerJobId, submissionKey }) { const result = await lookupCompatibilityAdapter(payload, providerJobId, submissionKey, env); const rawResult = (result.envelope.rawResult ?? result.envelope) as unknown; const actualCost = result.state === 'completed' ? paidActualCost(rawResult, task.estimatedCost.currency) : null; if (result.state === 'completed' && !actualCost) { throw new Error('Paid compatibility adapter completed without authoritative actualCost'); } return { ...result, envelope: { ...result.envelope, outputCount: result.outputs.length, ...(actualCost ? { actualCost } : {}), }, outputs: result.outputs, actualCost, }; }, }; } function record(value: unknown, label: string): Record { if (!value || typeof value !== 'object' || Array.isArray(value)) throw new Error(`${label} must be an object`); return value as Record; } function chainFallbackEnabled(task: CinemaQueueTask): boolean { const binding = task.payload.parameters.chainBinding; return !!binding && typeof binding === 'object' && !Array.isArray(binding) && (binding as Record).failFast === false; } function strings(value: unknown): string[] { return Array.isArray(value) ? value.filter((entry): entry is string => typeof entry === 'string' && entry.trim().length > 0) : []; } function isStringArray(value: unknown): value is string[] { return Array.isArray(value) && value.every((entry) => typeof entry === 'string'); } function executionTask(value: unknown): VideoExecutionTask | null { if (!value || typeof value !== 'object' || Array.isArray(value)) return null; const task = value as Record; if (!Number.isInteger(task.sceneIndex) || Number(task.sceneIndex) < 0 || typeof task.prompt !== 'string' || !task.prompt.trim() || !['text', 'image', 'video'].includes(String(task.inputKind)) || !isStringArray(task.referencePaths) || !isStringArray(task.sourceAssetIds) || !isStringArray(task.backendHints) || !isStringArray(task.characters) || (task.durationSeconds !== undefined && (typeof task.durationSeconds !== 'number' || !Number.isFinite(task.durationSeconds) || task.durationSeconds <= 0))) return null; return structuredClone(task) as unknown as VideoExecutionTask; } function promptGuidance(value: unknown): VideoExecutionPayload['promptGuidance'] { if (value === undefined) return []; if (!Array.isArray(value)) throw new Error('Compatibility prompt guidance must be an array'); if (!value.every((entry) => { if (!entry || typeof entry !== 'object' || Array.isArray(entry)) return false; const guidance = entry as Record; return typeof guidance.name === 'string' && guidance.name.trim().length > 0 && typeof guidance.reason === 'string' && guidance.reason.trim().length > 0 && ['provider', 'framework'].includes(String(guidance.category)); })) throw new Error('Compatibility prompt guidance contains an invalid entry'); return structuredClone(value) as VideoExecutionPayload['promptGuidance']; } function profile(value: unknown): VideoExecutionPayload['executionProfile'] { const input = record(value, 'Compatibility execution profile'); if (!['16:9', '9:16', '1:1'].includes(String(input.aspectRatio)) || !['fast', 'quality'].includes(String(input.quality)) || !['720p', '1080p'].includes(String(input.resolution)) || typeof input.generateAudio !== 'boolean' || !Number.isInteger(input.outputCount) || Number(input.outputCount) < 1) { throw new Error('Compatibility execution profile is invalid'); } return structuredClone(input) as unknown as VideoExecutionPayload['executionProfile']; } function sceneIndexFor(task: CinemaQueueTask, nativeTask: VideoExecutionTask | null): number | null { if (nativeTask) return nativeTask.sceneIndex; const match = task.shotId?.match(/^scene-(\d+)$/); return match ? Number(match[1]) : null; } function mographTask(task: CinemaQueueTask, parameters: Record): VideoExecutionTask { const refs = strings(parameters.refs); const videoSource = typeof parameters.videoSource === 'string' && parameters.videoSource.trim() ? parameters.videoSource.trim() : null; const sceneIndex = sceneIndexFor(task, null) ?? 0; const prompt = typeof parameters.prompt === 'string' ? parameters.prompt.trim() : ''; if (!prompt) throw new Error('Compatibility mograph prompt must be a non-empty string'); if (parameters.durationSec !== undefined && (typeof parameters.durationSec !== 'number' || !Number.isFinite(parameters.durationSec) || parameters.durationSec <= 0)) { throw new Error('Compatibility mograph duration must be a positive finite number'); } return { sceneIndex, prompt, inputKind: videoSource ? 'video' : refs.length ? 'image' : 'text', referencePaths: videoSource ? [videoSource, ...refs] : refs, sourceAssetIds: [...task.payload.artifactIds], backendHints: ['cinema-compatibility-worker'], characters: [], ...(typeof parameters.durationSec === 'number' ? { durationSeconds: parameters.durationSec } : {}), }; } export function buildCinemaCompatibilityExecutionPayload( root: string, projectSlug: string, task: CinemaQueueTask, ): VideoExecutionPayload { if (task.kind !== 'video-generate') throw new Error(`Compatibility worker cannot execute ${task.kind} tasks`); if (!PROVIDER_ROUTES.has(task.routeId as ProviderRouteId)) throw new Error(`Compatibility worker has no provider transport for ${task.routeId}`); const parameters = record(task.payload.parameters, 'Compatibility queue parameters'); const nativeTask = executionTask(parameters.executionTask); if (parameters.executionTask !== undefined && !nativeTask) { throw new Error('Compatibility native execution task is invalid'); } const resolvedTask = nativeTask ?? mographTask(task, parameters); const operation = String(parameters.operationKind ?? (resolvedTask.inputKind === 'video' ? 'video-to-video' : resolvedTask.inputKind === 'image' ? 'image-to-video' : 'text-to-video')); if (!['text-to-video', 'image-to-video', 'video-to-video'].includes(operation)) throw new Error(`Compatibility operation kind is invalid: ${operation}`); const executionProfile = parameters.executionProfile ? profile(parameters.executionProfile) : { aspectRatio: ['16:9', '9:16', '1:1'].includes(String(parameters.aspect)) ? parameters.aspect as '16:9' | '9:16' | '1:1' : '16:9', quality: 'fast' as const, resolution: '720p' as const, generateAudio: false, outputCount: 1, }; return { workspaceRoot: root, projectSlug, productionMode: 'storyboard', routeId: task.routeId as ProviderRouteId, operationKind: operation as VideoExecutionPayload['operationKind'], executionProfile, generatedAt: task.createdAt, outputDir: join(resolveProjectWorkspace(projectSlug, root).projectDir, 'outputs', 'durable-queue', task.taskId), tasks: [resolvedTask], promptGuidance: promptGuidance(parameters.promptGuidance), }; } export async function resolveCinemaCompatibilityExecutionPayload( root: string, projectSlug: string, task: CinemaQueueTask, ): Promise { const payload = buildCinemaCompatibilityExecutionPayload(root, projectSlug, task); const rawBinding = task.payload.parameters.chainBinding; if (rawBinding === undefined) return payload; const binding = record(rawBinding, 'Compatibility chain binding'); if (binding.resolveSelectedVideoAtQuote !== true) return payload; if (!Array.isArray(binding.sourcePolicy)) throw new Error('Compatibility chain source policy must be an array'); const activeIndex = binding.activeSourcePolicyIndex === undefined ? 0 : Number(binding.activeSourcePolicyIndex); if (!Number.isInteger(activeIndex) || activeIndex < 0 || activeIndex >= binding.sourcePolicy.length) { throw new Error(`Compatibility chain task ${task.taskId} has an invalid active source-policy index`); } const active = record(binding.sourcePolicy[activeIndex], `Compatibility chain task ${task.taskId} active source policy`); if (active.kind === 'image-only') return payload; if (active.kind !== 'chain' || !Number.isInteger(active.sourceSceneIndex) || Number(active.sourceSceneIndex) < 0) { throw new Error(`Compatibility chain task ${task.taskId} active source policy is invalid`); } const sourceSceneIndex = Number(active.sourceSceneIndex); const selection = await readSceneSelectionArtifact(root, projectSlug); const candidates = await readSceneCandidatesArtifact(root, projectSlug); const selectedId = selection.scenes.find((entry) => entry.sceneIndex === sourceSceneIndex)?.selectedCandidateId; if (!selectedId) { throw new Error(`Compatibility chain task ${task.taskId} requires a reviewed selected candidate for active source scene ${sourceSceneIndex}`); } const candidate = candidates.scenes.find((entry) => entry.sceneIndex === sourceSceneIndex) ?.candidates.find((entry) => entry.id === selectedId); if (!candidate || candidate.status !== 'completed') { throw new Error(`Compatibility chain source ${selectedId} is not a completed candidate`); } const video = candidate.outputs.find((output) => output.kind === 'video'); if (!video?.path) throw new Error(`Compatibility chain source ${selectedId} has no video output`); const resolved = structuredClone(payload); const executionTask = resolved.tasks[0]; executionTask.inputKind = 'video'; executionTask.referenceRole = 'keyframe'; executionTask.referencePaths = [video.path, ...executionTask.referencePaths.filter((path) => path !== video.path)]; executionTask.chainedFromCandidateId = candidate.id; resolved.operationKind = 'video-to-video'; return resolved; } async function appendCompatibilitySubmissionAttempt( root: string, projectSlug: string, task: CinemaQueueTask, providerJobId: string, now: string, ): Promise { if (!task.execution) throw new Error(`Compatibility task ${task.taskId} has no provider execution receipt`); await appendCinemaGenerationAttempt(root, projectSlug, { attemptId: task.attemptId ?? `attempt-${task.taskId}`, shotId: task.shotId ?? task.taskId, parentAttemptId: null, payloadHash: task.payloadHash, routeId: task.routeId, status: 'submitted', preparedAt: task.createdAt, providerJobId, changedVariables: [], failureTags: [], evidenceIds: [...new Set([...task.evidenceIds, task.execution.providerStatusHash])], estimatedCost: task.estimatedCost, actualCost: { currency: task.estimatedCost.currency, amount: 0 }, }, now); } async function bridgeCompletedTask( root: string, projectSlug: string, task: CinemaQueueTask, payload: VideoExecutionPayload, ): Promise { const execution = task.execution; if (!execution?.providerJobId || !execution.download.finalPath || !execution.download.contentHash || execution.download.byteLength === null || !execution.download.completedAt) { throw new Error(`Compatibility task ${task.taskId} lacks a terminal download receipt`); } const providerJobId = execution.providerJobId; const finalPath = execution.download.finalPath; const completedAt = execution.download.completedAt; const attemptId = task.attemptId ?? `attempt-${task.taskId}`; const outcomeInput = { outcomeId: `outcome-${task.taskId}`, attemptId, shotId: task.shotId ?? task.taskId, queueTaskId: task.taskId, executionHash: task.executionHash, payloadHash: task.payloadHash, status: 'completed' as const, observedAt: execution.download.completedAt, providerJobId: execution.providerJobId, media: { path: execution.download.finalPath, contentHash: execution.download.contentHash, byteLength: execution.download.byteLength }, actualCost: task.actualCost, evidenceIds: [...task.evidenceIds], }; const outcome = { ...outcomeInput, contentHash: sha256Text(stableCinemaJson(outcomeInput)) }; await appendCinemaAttemptOutcome(root, projectSlug, outcome, execution.download.completedAt); const sceneIndex = sceneIndexFor(task, payload.tasks[0] ?? null); if (sceneIndex === null) return null; let candidate!: SceneCandidate; await withSceneArtifactsLock(root, projectSlug, async () => { const candidates = await readSceneCandidatesArtifact(root, projectSlug); const selection = await readSceneSelectionArtifact(root, projectSlug); const existing = candidates.scenes.flatMap((scene) => scene.candidates) .find((entry) => entry.source.externalJobId === providerJobId); if (existing) { candidate = existing; await writeSceneSelectionArtifact(root, projectSlug, markPending(selection, sceneIndex, [existing.id])); return; } const generationRound = maxRoundForScene(candidates, sceneIndex) + 1; candidate = { id: nextCandidateId(candidates, sceneIndex), generationRound, prompt: payload.tasks[0]?.prompt ?? '', route: task.routeId, submittedAt: execution.submittedAt ?? task.createdAt, completedAt, status: 'completed', outputs: [{ kind: 'video', path: finalPath, ...(payload.tasks[0]?.durationSeconds ? { durationSec: payload.tasks[0].durationSeconds } : {}) }], source: { executionRound: generationRound, adapter: 'native', externalJobId: providerJobId, chainedFromCandidateId: payload.tasks[0]?.chainedFromCandidateId ?? null }, }; await writeSceneCandidatesArtifact(root, projectSlug, appendCandidate(candidates, sceneIndex, candidate)); await writeSceneSelectionArtifact(root, projectSlug, markPending(selection, sceneIndex, [candidate.id])); }); const workspace = await ensureProjectWorkspace(projectSlug, root); await appendProjectEvent(workspace, { type: 'cinema.compatibility.candidate.ingested', recordedAt: execution.download.completedAt, payload: { taskId: task.taskId, attemptId, outcomeId: outcome.outcomeId, sceneIndex, candidateId: candidate.id, contentHash: outcome.contentHash }, }); return candidate; } export async function runCinemaCompatibilityWorkerOnce(options: { root: string; projectSlug: string; workerId: string; taskId: string; transport: CinemaCompatibilityWorkerTransport; confirmProviderCall: boolean; laneQueue?: CinemaLaneQueuePort; spendAuthorization?: CinemaSpendAuthorizationPort; now?: string; }): Promise { const now = options.now ?? new Date().toISOString(); const queue = await readCinemaProductionQueue(options.root, options.projectSlug); const before = queue?.tasks.find((entry) => entry.taskId === options.taskId); if (!before) throw new Error(`Compatibility queue task ${options.taskId} does not exist`); const payload = await resolveCinemaCompatibilityExecutionPayload(options.root, options.projectSlug, before); if (before.status === 'succeeded') { return { disposition: 'already-succeeded', task: before, candidate: await bridgeCompletedTask(options.root, options.projectSlug, before, payload), providerCalls: 0, generationCalls: 0, spendAuthorized: before.spendAuthorizationId !== null }; } if (['ready', 'waiting-lane', 'retry-ready'].includes(before.status) && !options.confirmProviderCall) { throw new Error('Compatibility worker requires explicit provider-call confirmation before claiming submission work'); } const laneQueue = options.laneQueue ?? createLaneCoordinator({ env: process.env, dbPath: join(options.root, 'lanes.db') }); const claimed = await claimCinemaTask(options.root, options.projectSlug, { taskId: options.taskId, workerId: options.workerId, laneQueue, ...(options.spendAuthorization ? { spendAuthorization: options.spendAuthorization } : {}), now, }); if (claimed.disposition === 'waiting-lane' && claimed.task) { return { disposition: 'waiting-lane', task: claimed.task, candidate: null, providerCalls: 0, generationCalls: 0, spendAuthorized: claimed.task.spendAuthorizationId !== null }; } if (!claimed.task?.lease) throw new Error(`Compatibility task ${options.taskId} was not leased`); const leaseId = claimed.task.lease.leaseId; if (claimed.task.status === 'leased') { await beginCinemaTaskSubmission(options.root, options.projectSlug, { taskId: options.taskId, leaseId, now }); let submitted; try { submitted = await options.transport.submit(payload); } catch (error) { await markCinemaTaskSubmissionUnknown(options.root, options.projectSlug, { taskId: options.taskId, leaseId, reason: `Compatibility provider acceptance is unknown: ${error instanceof Error ? error.message : String(error)}`, now }); throw error; } const result = await recordCinemaTaskSubmitted(options.root, options.projectSlug, { taskId: options.taskId, leaseId, providerJobId: submitted.providerJobId, providerStatus: submitted.providerStatus, envelope: submitted.envelope, now, }); await appendCompatibilitySubmissionAttempt(options.root, options.projectSlug, result.task, submitted.providerJobId, now); const handedOff = await handoffCinemaTaskLease(options.root, options.projectSlug, { taskId: options.taskId, leaseId, now }); return { disposition: 'submitted', task: handedOff.task, candidate: null, providerCalls: 1, generationCalls: 1, spendAuthorized: handedOff.task.spendAuthorizationId !== null }; } const providerJobId = claimed.task.execution?.providerJobId ?? null; const submissionKey = claimed.task.execution?.submissionKey; if (!submissionKey) throw new Error(`Compatibility task ${options.taskId} has no submission key to reconcile`); let observation: CinemaCompatibilityWorkerObservation | null = null; const reconciled = await reconcileCinemaProviderTask(options.root, options.projectSlug, { taskId: options.taskId, leaseId, laneQueue, now, port: { async lookup() { observation = await options.transport.lookup({ payload, providerJobId, submissionKey }); return observation; } }, }); if (!providerJobId && reconciled.task.execution?.providerJobId) { await appendCompatibilitySubmissionAttempt( options.root, options.projectSlug, reconciled.task, reconciled.task.execution.providerJobId, now, ); } if (reconciled.task.status !== 'downloading') { const active = ['polling', 'reconciling'].includes(reconciled.task.status) ? await handoffCinemaTaskLease(options.root, options.projectSlug, { taskId: options.taskId, leaseId, now }) : reconciled; const fallback = active.task.status === 'dead-letter' && chainFallbackEnabled(active.task) ? await advanceCinemaChainFallback(options.root, options.projectSlug, active.task.taskId) : undefined; return { disposition: 'pending', task: active.task, candidate: null, providerCalls: 1, generationCalls: 0, spendAuthorized: active.task.spendAuthorizationId !== null, ...(fallback ? { fallback } : {}), }; } const completedObservation = observation as CinemaCompatibilityWorkerObservation | null; const output = completedObservation?.outputs.find((entry) => entry.kind === 'video'); if (!output || !completedObservation?.actualCost) { await handoffCinemaTaskLease(options.root, options.projectSlug, { taskId: options.taskId, leaseId, now }); throw new Error(`Compatibility completed task ${options.taskId} lacks a video output or actual cost`); } await recordCinemaTaskDownloadStarted(options.root, options.projectSlug, { taskId: options.taskId, leaseId, temporaryPath: `${output.path}.provider`, now }); const completed = await completeCinemaTaskDownload(options.root, options.projectSlug, { taskId: options.taskId, leaseId, finalPath: output.path, evidenceIds: [reconciled.task.execution!.providerStatusHash], actualCost: completedObservation.actualCost, laneQueue, ...(options.spendAuthorization ? { spendAuthorization: options.spendAuthorization } : {}), now, }); const candidate = await bridgeCompletedTask(options.root, options.projectSlug, completed.task, payload); return { disposition: 'completed', task: completed.task, candidate, providerCalls: 1, generationCalls: 0, spendAuthorized: completed.task.spendAuthorizationId !== null }; }