import { join } from 'node:path'; import { beginCinemaTaskSubmission, claimCinemaTask, failCinemaTask, completeCinemaTaskDownload, handoffCinemaTaskLease, markCinemaTaskSubmissionUnknown, readCinemaProductionQueue, recordCinemaTaskDownloadStarted, recordCinemaTaskSubmitted, reconcileCinemaProviderTask, type CinemaLaneQueuePort, } from './cinema-production-queue.js'; import type { CinemaSpendAuthorizationPort } from './cinema-provider-types.js'; import { createLaneCoordinator } from './lane-coordinator.js'; import { buildCinemaImageSubmitPayload } from './cinema-image-payload.js'; import { ingestCinemaImageFromQueue, type CinemaImageResult } from './cinema-image-jobs.js'; import type { CinemaImageWorkerTransport } from './cinema-image-transport.js'; import type { CinemaQueueTask } from './cinema-queue-types.js'; export interface CinemaImageWorkerResult { disposition: 'submitted' | 'pending' | 'completed' | 'waiting-lane' | 'already-succeeded'; task: CinemaQueueTask; /** The pending-review image record, once the bytes are durably recorded. */ result: CinemaImageResult | null; providerCalls: number; generationCalls: number; spendAuthorized: boolean; } /** * Drain ONE `image-generate` task. * * This is a sibling of `runCinemaCompatibilityWorkerOnce`, not a fork of it. It * reuses every durable-state function that worker uses — claim, begin, record, * reconcile, download, complete, handoff — all of which are already * media-agnostic. What it cannot reuse is the video worker itself: its payload * builder refuses any kind but `video-generate` and any route outside the core * `ProviderRouteId` space, and it selects the completed output with a hardcoded * `kind === 'video'`. Those guards stay exactly as they are; they are the correct * refusal if an image task ever wanders in. * * A completed task bridges into `ingestCinemaImageFromQueue`, mirroring how a * completed video task bridges into a scene candidate. */ export async function runCinemaImageWorkerOnce(options: { root: string; projectSlug: string; taskId: string; workerId: string; transport: CinemaImageWorkerTransport; confirmProviderCall: boolean; laneQueue?: CinemaLaneQueuePort; spendAuthorization?: CinemaSpendAuthorizationPort; now?: string; }): Promise { const { root, projectSlug, taskId } = options; const now = options.now ?? new Date().toISOString(); const queue = await readCinemaProductionQueue(root, projectSlug); const before = queue?.tasks.find((entry) => entry.taskId === taskId); if (!before) throw new Error(`Cinema image task ${taskId} does not exist`); const payload = await buildCinemaImageSubmitPayload(root, projectSlug, before); if (before.status === 'succeeded') { return { disposition: 'already-succeeded', task: before, result: await bridgeRecordedImage(root, projectSlug, before, payload.jobId), providerCalls: 0, generationCalls: 0, spendAuthorized: before.spendAuthorizationId !== null, }; } // Diagnose the missing quote HERE. Falling through to the claim would fail with // "was not leased", which says nothing about the actual problem and sends the // reader looking at leases rather than at the ceremony they have not finished. if (before.status === 'awaiting-quote' || (before.authorizationRequirement === 'exact-quote' && !before.spendAuthorizationId)) { throw new Error(`Cinema image task ${taskId} has no bound spend authorization: take an exact quote and authorise it before any render`); } if (['ready', 'waiting-lane', 'retry-ready'].includes(before.status) && !options.confirmProviderCall) { throw new Error('Cinema image worker requires explicit provider-call confirmation before it renders'); } const laneQueue = options.laneQueue ?? createLaneCoordinator({ env: process.env, dbPath: join(root, 'lanes.db') }); const claimed = await claimCinemaTask(root, projectSlug, { 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, result: null, providerCalls: 0, generationCalls: 0, spendAuthorized: claimed.task.spendAuthorizationId !== null }; } if (!claimed.task?.lease) throw new Error(`Cinema image task ${taskId} was not leased`); const leaseId = claimed.task.lease.leaseId; if (claimed.task.status === 'leased') { // While still `leased`, a task can simply FAIL — which is an ending an // operator can act on. Once `beginCinemaTaskSubmission` has run, a failure is // an ambiguous submission and parks in `reconciling`, and nothing can take a // reconciling image task back (claimCinemaTask only selects ready / // waiting-lane). So ask first whether the render can be attempted at all. if (options.transport.preflight) { try { await options.transport.preflight(); } catch (error) { const failed = await failCinemaTask(root, projectSlug, { taskId, leaseId, laneQueue, now, failure: { code: 'image-transport-not-ready', message: `${error instanceof Error ? error.message : String(error)} Fix it and compile the job again.`, retryable: true, }, }); return { disposition: 'pending', task: failed.task, result: null, providerCalls: 0, generationCalls: 0, spendAuthorized: failed.task.spendAuthorizationId !== null }; } } await beginCinemaTaskSubmission(root, projectSlug, { taskId, leaseId, now }); let submitted; try { submitted = await options.transport.submit(payload); } catch (error) { // The call may have cost money before it failed. Park the task for // reconciliation; it is never resubmitted automatically. // Reaching here means the render was attempted, so the provider may have // charged. That is genuinely ambiguous and must park for reconciliation. // (An earlier attempt recorded a `reachedProvider` flag on the envelope for // a later lookup to read; it never survived, because `recoverExpiredLeases` // replaces `providerStatusEnvelope` wholesale. The preflight above is what // separates the two cases now, before the queue is told anything.) await markCinemaTaskSubmissionUnknown(root, projectSlug, { taskId, leaseId, laneQueue, reason: `Cinema image provider acceptance is unknown: ${error instanceof Error ? error.message : String(error)}`, now, }); throw error; } const recorded = await recordCinemaTaskSubmitted(root, projectSlug, { taskId, leaseId, providerJobId: submitted.providerJobId, providerStatus: submitted.providerStatus, envelope: submitted.envelope, now, }); const handedOff = await handoffCinemaTaskLease(root, projectSlug, { taskId, leaseId, now }); return { disposition: 'submitted', task: handedOff.task, result: null, providerCalls: 1, generationCalls: 1, spendAuthorized: recorded.task.spendAuthorizationId !== null }; } const reconciled = await reconcileCinemaProviderTask(root, projectSlug, { taskId, leaseId, laneQueue, now, port: { lookup: async ({ providerJobId }) => options.transport.lookup({ payload, providerJobId }), }, }); if (reconciled.task.status !== 'downloading') { const active = ['polling', 'reconciling'].includes(reconciled.task.status) ? await handoffCinemaTaskLease(root, projectSlug, { taskId, leaseId, now }) : reconciled; return { disposition: 'pending', task: active.task, result: null, providerCalls: 1, generationCalls: 0, spendAuthorized: active.task.spendAuthorizationId !== null }; } // The bytes are already on disk — submit wrote them. "Download" here is the // durable record of where they are and what they cost. const observation = await options.transport.lookup({ payload, providerJobId: reconciled.task.execution?.providerJobId ?? null }); const output = observation.outputs.find((entry) => entry.kind === 'image'); if (!output || !observation.actualCost) { await handoffCinemaTaskLease(root, projectSlug, { taskId, leaseId, now }); throw new Error(`Cinema image task ${taskId} completed without an image output or an actual cost`); } await recordCinemaTaskDownloadStarted(root, projectSlug, { taskId, leaseId, temporaryPath: `${output.path}.provider`, now }); const completed = await completeCinemaTaskDownload(root, projectSlug, { taskId, leaseId, finalPath: output.path, evidenceIds: [reconciled.task.execution!.providerStatusHash], actualCost: observation.actualCost, laneQueue, ...(options.spendAuthorization ? { spendAuthorization: options.spendAuthorization } : {}), now, }); const result = await bridgeRecordedImage(root, projectSlug, completed.task, payload.jobId); return { disposition: 'completed', task: completed.task, result, providerCalls: 1, generationCalls: 0, spendAuthorized: completed.task.spendAuthorizationId !== null }; } /** * Hand the rendered bytes to the image lane's existing review chain. The record * is write-once, so a repeated call on an already-succeeded task returns the same * record rather than failing. */ async function bridgeRecordedImage( root: string, projectSlug: string, task: CinemaQueueTask, jobId: string, ): Promise { const finalPath = task.execution?.download.finalPath; const providerJobId = task.execution?.providerJobId; if (!finalPath || !providerJobId) { throw new Error(`Cinema image task ${task.taskId} lacks a terminal download receipt`); } return ingestCinemaImageFromQueue(root, projectSlug, jobId, { file: finalPath, routeId: task.routeId, taskId: task.taskId, providerJobId, }); }