import { createHash, randomUUID } from 'node:crypto'; import { existsSync } from 'node:fs'; import { mkdir, readFile } from 'node:fs/promises'; import { dirname, join } from 'node:path'; import { writeTextFileAtomic } from './atomic-write.js'; import { withCinemaFileLock } from './cinema-file-lock.js'; import { cinemaArtifactsDirFor } from './cinema-store.js'; import type { CinemaProductionQueueArtifact, CinemaQueueLaneOwnership, CinemaQueueMoney, CinemaQueuePayload, CinemaQueueTask, CinemaQueueTaskKind, } from './cinema-queue-types.js'; import { cinemaQueueExecutionHash, validateCinemaProductionQueue } from './cinema-queue-validation.js'; import { sha256Text, stableCinemaJson } from './cinema-shot-compiler.js'; import type { CinemaSpendAuthorizationPort } from './cinema-provider-types.js'; import type { CinemaProviderContractBundle } from './cinema-provider-types.js'; import { validateCinemaGenerationAuthorization } from './cinema-provider-contracts.js'; import type { AcquireResult, LaneStatus } from './lane-queue.js'; export interface CinemaLaneQueuePort { acquire(params: { lane: string; project: string; scene?: number | null; promptHash?: string | null; limit: number | null; ttlSec?: number; holder?: string; }): AcquireResult; status(lane: string, limit?: number | null): LaneStatus; heartbeat(ticketId: string, ttlSec?: number): boolean; release(ticketId: string, lane?: string, limit?: number | null): void; } export interface EnqueueCinemaTaskInput { taskId: string; kind: CinemaQueueTaskKind; payload: CinemaQueuePayload; dependencies?: string[]; routeId: string; shotId?: string; attemptId?: string; lane?: { laneId: string; limit: number | null } | null; estimatedCost?: CinemaQueueMoney; authorizationRequirement?: 'none' | 'exact-quote'; spendAuthorizationId?: string | null; createdAt?: string; } export interface CinemaQueueMutationResult { queue: CinemaProductionQueueArtifact; task: CinemaQueueTask; } export interface CinemaQueueEnqueueResult extends CinemaQueueMutationResult { enqueued: boolean; } export interface CinemaQueueClaimResult { queue: CinemaProductionQueueArtifact; task: CinemaQueueTask | null; disposition: 'leased' | 'waiting-lane' | 'empty'; } export interface CinemaQueueAuthorizationBindingResult { queue: CinemaProductionQueueArtifact; tasks: CinemaQueueTask[]; } export interface CinemaProviderReconciliationPort { lookup(input: { routeId: string; submissionKey: string; providerJobId: string | null; }): Promise<{ state: 'not-found' | 'accepted' | 'processing' | 'completed' | 'failed' | 'unknown'; authoritative: boolean; providerJobId: string | null; providerStatus: string; envelope: Record; failure?: { code: string; message: string; retryable: boolean }; }>; } function verifyPaidTaskAuthorization( projectSlug: string, task: CinemaQueueTask, queue: CinemaProductionQueueArtifact, spendAuthorization: CinemaSpendAuthorizationPort | undefined, now: string, cost = task.estimatedCost, ): void { if (!task.spendAuthorizationId) { throw new Error(`Cinema paid task ${task.taskId} cannot be claimed without exact spend authorization`); } if (!spendAuthorization) { throw new Error(`Cinema paid task ${task.taskId} requires cryptographic spend-authorization verification`); } spendAuthorization.verifyPaidTask({ projectSlug, taskId: task.taskId, taskKind: task.kind, routeId: task.routeId, payloadHash: task.payloadHash, estimatedCost: cost, spendAuthorizationId: task.spendAuthorizationId, now, queueAuthorizationUses: queue.tasks .filter((entry) => entry.spendAuthorizationId === task.spendAuthorizationId) .map((entry) => ({ taskId: entry.taskId, payloadHash: entry.payloadHash, status: entry.status, estimatedCost: entry.estimatedCost, actualCost: entry.actualCost, })), }); } export function cinemaProductionQueuePathFor(root: string, projectSlug: string): string { return join(cinemaArtifactsDirFor(root, projectSlug), 'production-queue.json'); } function emptyQueue(projectSlug: string, generatedAt: string): CinemaProductionQueueArtifact { return { schemaVersion: 1, projectSlug, generatedAt, tasks: [] }; } function assertQueue(queue: CinemaProductionQueueArtifact): void { const validation = validateCinemaProductionQueue(queue); if (!validation.ok) { const detail = validation.issues.slice(0, 8).map((issue) => `${issue.path}: ${issue.message}`).join('; '); throw new Error(`Cinema production queue failed validation: ${detail}`); } } export async function readCinemaProductionQueue( root: string, projectSlug: string, ): Promise { const path = cinemaProductionQueuePathFor(root, projectSlug); if (!existsSync(path)) return null; let queue: CinemaProductionQueueArtifact; try { queue = JSON.parse(await readFile(path, 'utf8')) as CinemaProductionQueueArtifact; } catch (error) { throw new Error(`Cinema production queue is not valid JSON at ${path}: ${error instanceof Error ? error.message : String(error)}`); } if (queue.projectSlug !== projectSlug) throw new Error(`Cinema production queue belongs to ${queue.projectSlug}, expected ${projectSlug}`); assertQueue(queue); return queue; } function dependenciesSucceeded(task: CinemaQueueTask, tasks: Map): boolean { return task.dependencies.every((taskId) => tasks.get(taskId)?.status === 'succeeded'); } function refreshReadiness(queue: CinemaProductionQueueArtifact, updatedAt: string): void { const tasks = new Map(queue.tasks.map((task) => [task.taskId, task])); queue.tasks.forEach((task) => { if (!['awaiting-quote', 'blocked', 'ready'].includes(task.status)) return; const dependenciesReady = dependenciesSucceeded(task, tasks); const status = !dependenciesReady ? 'blocked' : task.authorizationRequirement === 'exact-quote' && !task.spendAuthorizationId ? 'awaiting-quote' : 'ready'; if (task.status !== status) { task.status = status; task.updatedAt = updatedAt; } }); } async function updateQueue( root: string, projectSlug: string, now: string, operation: (queue: CinemaProductionQueueArtifact) => T, ): Promise { const path = cinemaProductionQueuePathFor(root, projectSlug); return withCinemaFileLock(path, async () => { const queue = await readCinemaProductionQueue(root, projectSlug) ?? emptyQueue(projectSlug, now); const result = operation(queue); queue.generatedAt = now; refreshReadiness(queue, now); assertQueue(queue); await mkdir(dirname(path), { recursive: true }); await writeTextFileAtomic(path, `${JSON.stringify(queue, null, 2)}\n`); return result; }); } function noTicketLane(input: EnqueueCinemaTaskInput['lane']): CinemaQueueLaneOwnership | null { if (!input) return null; return { laneId: input.laneId, limit: input.limit, ticketId: null, ticketState: 'none', position: null, etaSeconds: null, }; } export async function enqueueCinemaTask( root: string, projectSlug: string, input: EnqueueCinemaTaskInput, ): Promise { const now = input.createdAt ?? new Date().toISOString(); let response!: CinemaQueueEnqueueResult; await updateQueue(root, projectSlug, now, (queue) => { const dependencies = input.dependencies ?? []; const authorizationRequirement = input.authorizationRequirement ?? ((input.estimatedCost?.amount ?? 0) > 0 ? 'exact-quote' : 'none'); const payload = structuredClone(input.payload); const payloadHash = sha256Text(stableCinemaJson(payload)); const executionHash = cinemaQueueExecutionHash({ kind: input.kind, payloadHash, dependencies, routeId: input.routeId, ...(input.shotId ? { shotId: input.shotId } : {}), ...(input.attemptId ? { attemptId: input.attemptId } : {}), lane: input.lane ?? null, authorizationRequirement, }); const duplicate = queue.tasks.find((task) => task.executionHash === executionHash); if (duplicate) { response = { enqueued: false, queue, task: duplicate }; return; } const idCollision = queue.tasks.find((task) => task.taskId === input.taskId); if (idCollision) throw new Error(`Cinema queue task id ${input.taskId} already exists with a different payload`); const taskIds = new Set(queue.tasks.map((task) => task.taskId)); const missing = dependencies.filter((taskId) => !taskIds.has(taskId)); if (missing.length > 0) throw new Error(`Cinema queue task ${input.taskId} has missing dependencies: ${missing.join(', ')}`); const task: CinemaQueueTask = { taskId: input.taskId, kind: input.kind, executionHash, payloadHash, payload, status: dependencies.length > 0 ? 'blocked' : authorizationRequirement === 'exact-quote' && !input.spendAuthorizationId ? 'awaiting-quote' : 'ready', dependencies: [...dependencies], routeId: input.routeId, ...(input.shotId ? { shotId: input.shotId } : {}), ...(input.attemptId ? { attemptId: input.attemptId } : {}), lane: noTicketLane(input.lane), lease: null, execution: null, createdAt: now, updatedAt: now, estimatedCost: input.estimatedCost ?? { currency: 'USD', amount: 0 }, actualCost: { currency: input.estimatedCost?.currency ?? 'USD', amount: 0 }, authorizationRequirement, spendAuthorizationId: input.spendAuthorizationId ?? null, evidenceIds: [], failure: null, }; queue.tasks.push(task); response = { enqueued: true, queue, task }; }); return response; } export async function bindCinemaQuotedAuthorization( root: string, projectSlug: string, input: { bundle: CinemaProviderContractBundle; now?: string }, ): Promise { const now = input.now ?? new Date().toISOString(); const { snapshot, quote, authorization } = input.bundle; if (projectSlug !== quote.projectSlug) throw new Error(`Cinema quote belongs to ${quote.projectSlug}, expected ${projectSlug}`); const validation = validateCinemaGenerationAuthorization(authorization, quote, snapshot, now); if (!validation.ok) { throw new Error(`Cinema generation authorization cannot bind queue tasks: ${validation.issues.map((issue) => `${issue.code} (${issue.path})`).join(', ')}`); } let response!: CinemaQueueAuthorizationBindingResult; await updateQueue(root, projectSlug, now, (queue) => { const bound: CinemaQueueTask[] = []; for (const job of quote.jobs) { const matches = queue.tasks.filter((task) => task.routeId === quote.routeId && task.kind === job.taskKind && task.payloadHash === job.payloadHash); if (matches.length !== 1) { throw new Error(`Cinema quoted job ${job.jobId} resolves to ${matches.length} durable queue tasks; exactly one is required`); } const task = matches[0]; if (task.authorizationRequirement !== 'exact-quote') { throw new Error(`Cinema quoted job ${job.jobId} resolved to provider-free task ${task.taskId}`); } if (!['awaiting-quote', 'blocked', 'ready'].includes(task.status)) { throw new Error(`Cinema queue task ${task.taskId} cannot bind authorization from status ${task.status}`); } if (task.spendAuthorizationId && task.spendAuthorizationId !== authorization.authorizationId) { throw new Error(`Cinema queue task ${task.taskId} already binds a different authorization`); } task.estimatedCost = structuredClone(job.maximumCost); // The enqueue placeholder is `{ USD, 0 }`; a quote in another currency // (Flow quotes in credits) retags the untouched placeholder so the // queue's own currency invariant holds. Recorded spend never changes // currency underneath an authorization. if (task.actualCost.currency !== job.maximumCost.currency) { if (task.actualCost.amount !== 0) { throw new Error(`Cinema queue task ${task.taskId} already recorded ${task.actualCost.amount} ${task.actualCost.currency}; cannot bind a ${job.maximumCost.currency} quote`); } task.actualCost = { currency: job.maximumCost.currency, amount: 0 }; } task.spendAuthorizationId = authorization.authorizationId; task.evidenceIds = [...new Set([ ...task.evidenceIds, snapshot.contentHash, quote.contentHash, authorization.contentHash, ...authorization.complianceEvidenceIds, ])]; task.updatedAt = now; bound.push(task); } const quotedHashes = new Set(quote.jobs.map((job) => job.payloadHash)); const foreignUses = queue.tasks.filter((task) => task.spendAuthorizationId === authorization.authorizationId && !quotedHashes.has(task.payloadHash)); if (foreignUses.length > 0) { throw new Error(`Cinema authorization ${authorization.authorizationId} is already attached outside its exact quote`); } response = { queue, tasks: bound }; }); return response; } function clearLaneTicket(task: CinemaQueueTask): void { if (!task.lane) return; task.lane.ticketId = null; task.lane.ticketState = 'none'; task.lane.position = null; task.lane.etaSeconds = null; } function releaseLane(task: CinemaQueueTask, laneQueue: CinemaLaneQueuePort | undefined): void { if (task.lane?.ticketId && !laneQueue) { throw new Error(`Cinema task ${task.taskId} cannot release LaneQueue ticket ${task.lane.ticketId} without a LaneQueue`); } if (task.lane?.ticketId && laneQueue) laneQueue.release(task.lane.ticketId, task.lane.laneId, task.lane.limit); clearLaneTicket(task); } function recoverExpiredLeases( queue: CinemaProductionQueueArtifact, now: string, laneQueue: CinemaLaneQueuePort | undefined, workerId: string, leaseSeconds: number, taskId?: string, ): CinemaQueueTask | null { const nowMs = Date.parse(now); const active = ['leased', 'submitting', 'reconciling', 'submitted', 'polling', 'downloading']; const task = queue.tasks.find((entry) => (!taskId || entry.taskId === taskId) && active.includes(entry.status) && entry.lease && Date.parse(entry.lease.expiresAt) <= nowMs); if (!task) return null; if (task.execution) { const previousStatus = task.status; task.status = 'reconciling'; task.execution.reconciliation = { required: true, reason: `worker lease expired during ${previousStatus}; provider state must reconcile before any resubmission`, attempts: task.execution.reconciliation.attempts, }; setProviderObservation(task, now, 'worker-lease-expired', { phase: 'worker-lease-expired', previousTaskStatus: previousStatus, providerJobId: task.execution.providerJobId, }); leaseTask(task, workerId, now, leaseSeconds, true); return task; } if (task.status === 'leased') { releaseLane(task, laneQueue); task.lease = null; task.status = 'ready'; task.updatedAt = now; } return null; } function reconcileWaitingLane(task: CinemaQueueTask, laneQueue: CinemaLaneQueuePort): void { if (!task.lane?.ticketId) return; const status = laneQueue.status(task.lane.laneId, task.lane.limit); const held = status.held.find((ticket) => ticket.id === task.lane?.ticketId); if (held) { task.lane.ticketState = 'held'; task.lane.position = 0; task.lane.etaSeconds = 0; return; } const queued = status.queued.find((ticket) => ticket.id === task.lane?.ticketId); if (queued) { task.lane.ticketState = 'queued'; task.lane.position = queued.position; task.lane.etaSeconds = status.medianJobSeconds === null ? null : Math.round(queued.position * status.medianJobSeconds); return; } clearLaneTicket(task); task.status = 'ready'; } function leaseTask( task: CinemaQueueTask, workerId: string, now: string, leaseSeconds: number, preserveStatus = false, ): void { if (!preserveStatus) task.status = 'leased'; task.lease = { leaseId: randomUUID(), workerId, acquiredAt: now, expiresAt: new Date(Date.parse(now) + leaseSeconds * 1000).toISOString(), }; task.updatedAt = now; } function setProviderObservation( task: CinemaQueueTask, now: string, providerStatus: string, envelope: Record, ): void { if (!task.execution) throw new Error(`Cinema task ${task.taskId} has no provider execution receipt`); task.execution.lastObservedAt = now; task.execution.providerStatus = providerStatus; task.execution.providerStatusEnvelope = structuredClone(envelope); task.execution.providerStatusHash = sha256Text(stableCinemaJson(envelope)); task.updatedAt = now; } export async function beginCinemaTaskSubmission( root: string, projectSlug: string, input: { taskId: string; leaseId: string; now?: string }, ): Promise { const now = input.now ?? new Date().toISOString(); let response!: CinemaQueueMutationResult; await updateQueue(root, projectSlug, now, (queue) => { const task = queue.tasks.find((entry) => entry.taskId === input.taskId); if (!task || task.status !== 'leased' || task.lease?.leaseId !== input.leaseId) { throw new Error(`Cinema task ${input.taskId} is not owned by lease ${input.leaseId}`); } const envelope = { phase: 'submitting', submissionKey: task.payloadHash }; task.status = 'submitting'; task.execution = { submissionKey: task.payloadHash, providerJobId: null, submittedAt: null, lastObservedAt: now, providerStatus: 'submitting', providerStatusEnvelope: envelope, providerStatusHash: sha256Text(stableCinemaJson(envelope)), reconciliation: { required: false, reason: null, attempts: 0 }, download: { temporaryPath: null, finalPath: null, contentHash: null, byteLength: null, completedAt: null }, }; task.updatedAt = now; response = { queue, task }; }); return response; } export async function recordCinemaTaskSubmitted( root: string, projectSlug: string, input: { taskId: string; leaseId: string; providerJobId: string; providerStatus: string; envelope: Record; now?: string; }, ): Promise { const now = input.now ?? new Date().toISOString(); let response!: CinemaQueueMutationResult; await updateQueue(root, projectSlug, now, (queue) => { const task = queue.tasks.find((entry) => entry.taskId === input.taskId); if (!task || task.status !== 'submitting' || task.lease?.leaseId !== input.leaseId || !task.execution) { throw new Error(`Cinema task ${input.taskId} is not in an owned submitting state`); } if (!input.providerJobId.trim()) throw new Error('Cinema provider submission requires a durable provider job id'); task.status = 'submitted'; task.execution.providerJobId = input.providerJobId; task.execution.submittedAt = now; task.execution.reconciliation = { required: false, reason: null, attempts: task.execution.reconciliation.attempts }; setProviderObservation(task, now, input.providerStatus, input.envelope); response = { queue, task }; }); return response; } export async function markCinemaTaskSubmissionUnknown( root: string, projectSlug: string, input: { taskId: string; leaseId: string; reason: string; envelope?: Record; now?: string }, ): Promise { const now = input.now ?? new Date().toISOString(); let response!: CinemaQueueMutationResult; await updateQueue(root, projectSlug, now, (queue) => { const task = queue.tasks.find((entry) => entry.taskId === input.taskId); if (!task || task.status !== 'submitting' || task.lease?.leaseId !== input.leaseId || !task.execution) { throw new Error(`Cinema task ${input.taskId} is not in an owned submitting state`); } if (!input.reason.trim()) throw new Error('Cinema unknown submission requires a reconciliation reason'); task.status = 'reconciling'; task.execution.reconciliation = { required: true, reason: input.reason, attempts: task.execution.reconciliation.attempts }; setProviderObservation(task, now, 'submission-unknown', input.envelope ?? { phase: 'submission-unknown', reason: input.reason }); response = { queue, task }; }); return response; } export async function reconcileCinemaProviderTask( root: string, projectSlug: string, input: { taskId: string; leaseId: string; port: CinemaProviderReconciliationPort; laneQueue?: CinemaLaneQueuePort; now?: string; }, ): Promise { const now = input.now ?? new Date().toISOString(); const before = await readCinemaProductionQueue(root, projectSlug); const snapshot = before?.tasks.find((entry) => entry.taskId === input.taskId); if (!snapshot?.execution || snapshot.lease?.leaseId !== input.leaseId || !['reconciling', 'submitted', 'polling', 'downloading'].includes(snapshot.status)) { throw new Error(`Cinema task ${input.taskId} is not in an owned reconcilable state`); } const expectedStatusHash = snapshot.execution.providerStatusHash; const observed = await input.port.lookup({ routeId: snapshot.routeId, submissionKey: snapshot.execution.submissionKey, providerJobId: snapshot.execution.providerJobId, }); let response!: CinemaQueueMutationResult; await updateQueue(root, projectSlug, now, (queue) => { const task = queue.tasks.find((entry) => entry.taskId === input.taskId); if (!task?.execution || task.lease?.leaseId !== input.leaseId || task.execution.providerStatusHash !== expectedStatusHash) { throw new Error(`Cinema task ${input.taskId} changed while provider reconciliation was in flight`); } task.execution.reconciliation.attempts += 1; setProviderObservation(task, now, observed.providerStatus, observed.envelope); if (observed.state === 'not-found') { if (!observed.authoritative) { task.status = 'reconciling'; task.execution.reconciliation = { required: true, reason: 'provider lookup was not authoritative', attempts: task.execution.reconciliation.attempts }; } else { releaseLane(task, input.laneQueue); task.status = 'retry-ready'; task.lease = null; task.execution.reconciliation = { required: false, reason: null, attempts: task.execution.reconciliation.attempts }; task.failure = { code: 'provider-job-not-found', message: 'Authoritative lookup found no accepted provider job; explicit retry is safe', retryable: true }; } } else if (observed.state === 'unknown') { task.status = 'reconciling'; task.execution.reconciliation = { required: true, reason: 'provider state remains unknown', attempts: task.execution.reconciliation.attempts }; } else if (observed.state === 'failed') { releaseLane(task, input.laneQueue); task.status = observed.failure?.retryable ? 'retry-ready' : 'dead-letter'; task.lease = null; task.failure = observed.failure ?? { code: 'provider-failed', message: observed.providerStatus, retryable: false }; task.execution.reconciliation = { required: false, reason: null, attempts: task.execution.reconciliation.attempts }; } else { if (!observed.providerJobId?.trim()) throw new Error(`Cinema provider ${observed.state} observation requires a provider job id`); task.execution.providerJobId = observed.providerJobId; task.execution.submittedAt ??= now; task.execution.reconciliation = { required: false, reason: null, attempts: task.execution.reconciliation.attempts }; task.status = observed.state === 'completed' ? 'downloading' : 'polling'; } response = { queue, task }; }); return response; } export async function recordCinemaTaskDownloadStarted( root: string, projectSlug: string, input: { taskId: string; leaseId: string; temporaryPath: string; now?: string }, ): Promise { const now = input.now ?? new Date().toISOString(); let response!: CinemaQueueMutationResult; await updateQueue(root, projectSlug, now, (queue) => { const task = queue.tasks.find((entry) => entry.taskId === input.taskId); if (!task?.execution || task.lease?.leaseId !== input.leaseId || !['submitted', 'polling', 'downloading'].includes(task.status) || !task.execution.providerJobId) { throw new Error(`Cinema task ${input.taskId} is not an owned downloadable provider job`); } if (!input.temporaryPath.trim()) throw new Error('Cinema download requires a temporary path'); task.status = 'downloading'; task.execution.download.temporaryPath = input.temporaryPath; task.updatedAt = now; response = { queue, task }; }); return response; } export async function completeCinemaTaskDownload( root: string, projectSlug: string, input: { taskId: string; leaseId: string; finalPath: string; evidenceIds: string[]; actualCost?: CinemaQueueMoney; spendAuthorization?: CinemaSpendAuthorizationPort; laneQueue?: CinemaLaneQueuePort; now?: string; }, ): Promise { const now = input.now ?? new Date().toISOString(); if (!input.finalPath.trim()) throw new Error('Cinema completed download requires a final path'); const bytes = await readFile(input.finalPath); const contentHash = `sha256:${createHash('sha256').update(bytes).digest('hex')}` as const; let response!: CinemaQueueMutationResult; await updateQueue(root, projectSlug, now, (queue) => { const task = queue.tasks.find((entry) => entry.taskId === input.taskId); if (!task?.execution || task.status !== 'downloading' || task.lease?.leaseId !== input.leaseId) { throw new Error(`Cinema task ${input.taskId} is not an owned in-progress download`); } if (input.evidenceIds.length === 0) throw new Error(`Cinema task ${input.taskId} download completion requires evidence`); const actualCost = input.actualCost ?? { currency: task.estimatedCost.currency, amount: 0 }; if (actualCost.amount < 0 || actualCost.currency !== task.estimatedCost.currency) { throw new Error(`Cinema task ${input.taskId} actual cost must be non-negative and use ${task.estimatedCost.currency}`); } if (task.authorizationRequirement === 'exact-quote') { verifyPaidTaskAuthorization(projectSlug, task, queue, input.spendAuthorization, now, actualCost); } else if (actualCost.amount > 0) { throw new Error(`Cinema provider-free task ${input.taskId} cannot record spend`); } task.execution.download = { temporaryPath: task.execution.download.temporaryPath, finalPath: input.finalPath, contentHash, byteLength: bytes.byteLength, completedAt: now, }; setProviderObservation(task, now, 'downloaded', { status: 'downloaded', providerJobId: task.execution.providerJobId, contentHash, byteLength: bytes.byteLength, }); releaseLane(task, input.laneQueue); task.status = 'succeeded'; task.lease = null; task.evidenceIds = [...new Set([...task.evidenceIds, ...input.evidenceIds])]; task.actualCost = actualCost; task.failure = null; response = { queue, task }; }); return response; } export async function claimCinemaTask( root: string, projectSlug: string, options: { workerId: string; laneQueue?: CinemaLaneQueuePort; spendAuthorization?: CinemaSpendAuthorizationPort; now?: string; leaseSeconds?: number; taskId?: string; }, ): Promise { const now = options.now ?? new Date().toISOString(); const leaseSeconds = options.leaseSeconds ?? 300; let response!: CinemaQueueClaimResult; await updateQueue(root, projectSlug, now, (queue) => { const recoveredProviderTask = recoverExpiredLeases( queue, now, options.laneQueue, options.workerId, leaseSeconds, options.taskId, ); if (recoveredProviderTask) { response = { queue, task: recoveredProviderTask, disposition: 'leased' }; return; } queue.tasks.filter((task) => task.status === 'waiting-lane').forEach((task) => { if (!options.laneQueue) throw new Error(`Cinema task ${task.taskId} requires a LaneQueue to reconcile ticket ${task.lane?.ticketId}`); reconcileWaitingLane(task, options.laneQueue); }); refreshReadiness(queue, now); const selected = (task: CinemaQueueTask) => !options.taskId || task.taskId === options.taskId; const held = queue.tasks.find((task) => selected(task) && task.status === 'waiting-lane' && task.lane?.ticketState === 'held'); const task = held ?? queue.tasks.find((entry) => selected(entry) && entry.status === 'ready'); if (!task) { response = { queue, task: null, disposition: 'empty' }; return; } if (task.authorizationRequirement === 'exact-quote') { verifyPaidTaskAuthorization(projectSlug, task, queue, options.spendAuthorization, now); } if (task.lane?.ticketState === 'held') { leaseTask(task, options.workerId, now, leaseSeconds); response = { queue, task, disposition: 'leased' }; return; } if (task.lane) { if (!options.laneQueue) throw new Error(`Cinema task ${task.taskId} requires LaneQueue ${task.lane.laneId}`); const ticket = options.laneQueue.acquire({ lane: task.lane.laneId, project: projectSlug, scene: task.shotId ? Number(task.shotId.match(/\d+$/)?.[0] ?? 0) : null, promptHash: task.payloadHash, limit: task.lane.limit, holder: options.workerId, }); task.lane.ticketId = ticket.ticketId; task.lane.ticketState = ticket.status === 'granted' ? 'held' : 'queued'; task.lane.position = ticket.position; task.lane.etaSeconds = ticket.etaSeconds; task.updatedAt = now; if (ticket.status === 'queued') { task.status = 'waiting-lane'; response = { queue, task, disposition: 'waiting-lane' }; return; } } leaseTask(task, options.workerId, now, leaseSeconds); response = { queue, task, disposition: 'leased' }; }); return response; } export async function completeCinemaTask( root: string, projectSlug: string, input: { taskId: string; leaseId: string; evidenceIds: string[]; actualCost?: CinemaQueueMoney; laneQueue?: CinemaLaneQueuePort; spendAuthorization?: CinemaSpendAuthorizationPort; now?: string; }, ): Promise { const now = input.now ?? new Date().toISOString(); let response!: CinemaQueueMutationResult; await updateQueue(root, projectSlug, now, (queue) => { const task = queue.tasks.find((entry) => entry.taskId === input.taskId); if (!task) throw new Error(`Cinema queue task not found: ${input.taskId}`); if (task.status !== 'leased' || task.lease?.leaseId !== input.leaseId) throw new Error(`Cinema task ${input.taskId} is not owned by lease ${input.leaseId}`); if (input.evidenceIds.length === 0) throw new Error(`Cinema task ${input.taskId} completion requires evidence`); const actualCost = input.actualCost ?? { currency: task.estimatedCost.currency, amount: 0 }; if (actualCost.amount < 0 || actualCost.currency !== task.estimatedCost.currency) { throw new Error(`Cinema task ${input.taskId} actual cost must be non-negative and use ${task.estimatedCost.currency}`); } if (task.authorizationRequirement === 'exact-quote') { verifyPaidTaskAuthorization(projectSlug, task, queue, input.spendAuthorization, now, actualCost); } else if (actualCost.amount > 0) { throw new Error(`Cinema provider-free task ${input.taskId} cannot record spend`); } releaseLane(task, input.laneQueue); task.status = 'succeeded'; task.lease = null; task.evidenceIds = [...input.evidenceIds]; task.actualCost = actualCost; task.updatedAt = now; response = { queue, task }; }); return response; } export async function failCinemaTask( root: string, projectSlug: string, input: { taskId: string; leaseId: string; failure: NonNullable; evidenceIds?: string[]; laneQueue?: CinemaLaneQueuePort; now?: string; }, ): Promise { const now = input.now ?? new Date().toISOString(); let response!: CinemaQueueMutationResult; await updateQueue(root, projectSlug, now, (queue) => { const task = queue.tasks.find((entry) => entry.taskId === input.taskId); if (!task) throw new Error(`Cinema queue task not found: ${input.taskId}`); if (task.status !== 'leased' || task.lease?.leaseId !== input.leaseId) throw new Error(`Cinema task ${input.taskId} is not owned by lease ${input.leaseId}`); if (!input.failure.code.trim() || !input.failure.message.trim()) { throw new Error(`Cinema task ${input.taskId} failure requires a code and message`); } releaseLane(task, input.laneQueue); task.status = 'failed'; task.lease = null; task.failure = { ...input.failure }; task.evidenceIds = [...(input.evidenceIds ?? [])]; task.updatedAt = now; response = { queue, task }; }); return response; } /** * Cancel only untouched descendants of a dead-letter chain task after their * immutable replacement tail has been enqueued. Idempotent for the same * transition; refuses any task that has begun provider or local execution. */ export async function cancelCinemaChainFallbackTail( root: string, projectSlug: string, input: { taskIds: string[]; transitionId: string; replacementTaskIds: string[]; evidenceId: string; now?: string; }, ): Promise { const now = input.now ?? new Date().toISOString(); if (!input.transitionId.trim() || !input.evidenceId.trim() || input.replacementTaskIds.length === 0) { throw new Error('Cinema chain fallback cancellation requires transition, replacement and evidence identities'); } let response!: CinemaProductionQueueArtifact; await updateQueue(root, projectSlug, now, (queue) => { for (const taskId of input.taskIds) { const task = queue.tasks.find((entry) => entry.taskId === taskId); if (!task) throw new Error(`Cinema chain fallback descendant ${taskId} is missing`); if (task.status === 'cancelled') { if (!task.evidenceIds.includes(input.transitionId)) { throw new Error(`Cinema chain fallback descendant ${taskId} was cancelled by a different transition`); } continue; } if (!['blocked', 'awaiting-quote'].includes(task.status) || task.execution || task.lease) { throw new Error(`Cinema chain fallback cannot supersede ${taskId} after execution began (status ${task.status})`); } task.status = 'cancelled'; task.failure = null; task.evidenceIds = [...new Set([...task.evidenceIds, input.transitionId, input.evidenceId])]; task.updatedAt = now; } response = queue; }); return response; } export async function heartbeatCinemaTask( root: string, projectSlug: string, input: { taskId: string; leaseId: string; laneQueue?: CinemaLaneQueuePort; now?: string; leaseSeconds?: number; }, ): Promise { const now = input.now ?? new Date().toISOString(); const leaseSeconds = input.leaseSeconds ?? 300; let response!: CinemaQueueMutationResult; await updateQueue(root, projectSlug, now, (queue) => { const task = queue.tasks.find((entry) => entry.taskId === input.taskId); if (!task || task.status !== 'leased' || task.lease?.leaseId !== input.leaseId) throw new Error(`Cinema task ${input.taskId} is not owned by lease ${input.leaseId}`); task.lease.expiresAt = new Date(Date.parse(now) + leaseSeconds * 1000).toISOString(); task.updatedAt = now; if (task.lane?.ticketId && input.laneQueue && !input.laneQueue.heartbeat(task.lane.ticketId, leaseSeconds)) { throw new Error(`Cinema task ${task.taskId} lost LaneQueue ticket ${task.lane.ticketId}`); } response = { queue, task }; }); return response; } export async function handoffCinemaTaskLease( root: string, projectSlug: string, input: { taskId: string; leaseId: string; now?: string }, ): Promise { const now = input.now ?? new Date().toISOString(); let response!: CinemaQueueMutationResult; await updateQueue(root, projectSlug, now, (queue) => { const task = queue.tasks.find((entry) => entry.taskId === input.taskId); if (!task?.lease || task.lease.leaseId !== input.leaseId || !['submitted', 'polling', 'downloading', 'reconciling'].includes(task.status)) { throw new Error(`Cinema task ${input.taskId} has no owned provider lease to hand off`); } task.lease.expiresAt = now; task.updatedAt = now; response = { queue, task }; }); return response; }