import { cancelCinemaChainFallbackTail, readCinemaProductionQueue, } from './cinema-production-queue.js'; import type { CinemaQueueTask } from './cinema-queue-types.js'; import { enqueueCinemaCompatibilityPlan, projectCinemaCompatibilityStatus, type CinemaQueueCompatibilityReceiptArtifact, type CinemaQueueCompatibilityStatus, type CinemaQueueCompatibilityTask, } from './cinema-queue-compatibility.js'; import { sha256Text, stableCinemaJson } from './cinema-shot-compiler.js'; import { appendProjectEvent } from './events.js'; import { ensureProjectWorkspace } from './workspace.js'; interface ChainPolicy { kind: 'chain' | 'image-only'; sourceSceneIndex?: number; } interface ChainBinding extends Record { sourcePolicy: ChainPolicy[]; activeSourcePolicyIndex: number; failFast: boolean; } export interface CinemaChainFallbackTransition { transitionId: string; sourceTaskId: string; fromPolicyIndex: number; activePolicyIndex: number; fromPolicy: ChainPolicy; activePolicy: ChainPolicy; cancelledTaskIds: string[]; replacementTaskIds: string[]; receipt: CinemaQueueCompatibilityReceiptArtifact; receiptPath: string; status: CinemaQueueCompatibilityStatus; providerCalls: 0; generationCalls: 0; spendAuthorized: false; } 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 policy(value: unknown, label: string): ChainPolicy { const entry = record(value, label); if (entry.kind === 'image-only') return { kind: 'image-only' }; if (entry.kind === 'chain' && Number.isInteger(entry.sourceSceneIndex) && Number(entry.sourceSceneIndex) >= 0) { return { kind: 'chain', sourceSceneIndex: Number(entry.sourceSceneIndex) }; } throw new Error(`${label} is invalid`); } function bindingFor(task: CinemaQueueTask): ChainBinding { const binding = record(task.payload.parameters.chainBinding, `Cinema chain binding for ${task.taskId}`); if (!Array.isArray(binding.sourcePolicy) || binding.sourcePolicy.length === 0) { throw new Error(`Cinema chain task ${task.taskId} has no source policy`); } const sourcePolicy = binding.sourcePolicy.map((entry, index) => policy(entry, `Cinema chain policy ${index} for ${task.taskId}`)); const activeSourcePolicyIndex = binding.activeSourcePolicyIndex === undefined ? 0 : Number(binding.activeSourcePolicyIndex); if (!Number.isInteger(activeSourcePolicyIndex) || activeSourcePolicyIndex < 0 || activeSourcePolicyIndex >= sourcePolicy.length) { throw new Error(`Cinema chain task ${task.taskId} has an invalid active source-policy index`); } return { ...binding, sourcePolicy, activeSourcePolicyIndex, failFast: binding.failFast !== false, }; } function descendantTail(queueTasks: CinemaQueueTask[], sourceTaskId: string): CinemaQueueTask[] { const ids = new Set([sourceTaskId]); let changed = true; while (changed) { changed = false; for (const task of queueTasks) { if (ids.has(task.taskId) || !task.dependencies.some((dependency) => ids.has(dependency))) continue; ids.add(task.taskId); changed = true; } } return queueTasks.filter((task) => ids.has(task.taskId)); } function externalId(transitionId: string, task: CinemaQueueTask): string { return `fallback-${transitionId.slice(-16)}-${task.taskId}`; } /** * Advance one failed durable auto-chain task to exactly the next declared * source-policy rung. The failed attempt remains immutable. Any untouched * blocked descendants are cancelled and recreated against the replacement * dependency, so no tail can remain attached to the dead-letter task. * * This function performs no provider call and carries no prior authorization: * every replacement provider task must receive a fresh exact quote/approval. */ export async function advanceCinemaChainFallback( root: string, projectSlug: string, sourceTaskId: string, ): Promise { const queue = await readCinemaProductionQueue(root, projectSlug); const source = queue?.tasks.find((task) => task.taskId === sourceTaskId); if (!queue || !source) throw new Error(`Cinema chain fallback source ${sourceTaskId} does not exist`); if (source.status !== 'dead-letter' || !source.failure || source.failure.retryable) { throw new Error(`Cinema chain fallback requires an authoritative non-retryable dead-letter task (got ${source.status})`); } const sourceBinding = bindingFor(source); if (sourceBinding.failFast) throw new Error(`Cinema chain task ${sourceTaskId} is fail-fast and has no authorized fallback transition`); const nextPolicyIndex = sourceBinding.activeSourcePolicyIndex + 1; if (nextPolicyIndex >= sourceBinding.sourcePolicy.length) { throw new Error(`Cinema chain task ${sourceTaskId} exhausted its declared fallback policy`); } const tail = descendantTail(queue.tasks, sourceTaskId); const descendants = tail.filter((task) => task.taskId !== sourceTaskId); for (const task of descendants) { bindingFor(task); if (task.status === 'cancelled') continue; if (!['blocked', 'awaiting-quote'].includes(task.status) || task.execution || task.lease) { throw new Error(`Cinema chain fallback tail ${task.taskId} has already begun execution (status ${task.status})`); } } for (const dependencyId of source.dependencies) { if (queue.tasks.find((task) => task.taskId === dependencyId)?.status !== 'succeeded') { throw new Error(`Cinema chain fallback source dependency ${dependencyId} is not succeeded`); } } const transitionSeed = { projectSlug, sourceTaskId, sourceExecutionHash: source.executionHash, providerStatusHash: source.execution?.providerStatusHash ?? null, failure: source.failure, fromPolicyIndex: sourceBinding.activeSourcePolicyIndex, activePolicyIndex: nextPolicyIndex, }; const transitionId = `chain-fallback-${sha256Text(stableCinemaJson(transitionSeed)).slice(7, 31)}`; const tailIds = new Set(tail.map((task) => task.taskId)); const externalIds = new Map(tail.map((task) => [task.taskId, externalId(transitionId, task)])); const compatibilityTasks: CinemaQueueCompatibilityTask[] = tail.map((task) => { const oldBinding = bindingFor(task); const activeSourcePolicyIndex = task.taskId === sourceTaskId ? nextPolicyIndex : oldBinding.activeSourcePolicyIndex; const payload = structuredClone(task.payload); payload.parameters.chainBinding = { ...oldBinding, activeSourcePolicyIndex, fallbackState: { transitionId, previousTaskId: task.taskId, previousPolicyIndex: oldBinding.activeSourcePolicyIndex, activePolicyIndex: activeSourcePolicyIndex, triggerTaskId: sourceTaskId, triggerFailureCode: source.failure!.code, triggerProviderStatusHash: source.execution?.providerStatusHash ?? null, }, }; return { externalTaskId: externalIds.get(task.taskId)!, kind: task.kind, payload, existingDependencies: task.dependencies.filter((dependency) => !tailIds.has(dependency)), dependencies: task.dependencies.filter((dependency) => tailIds.has(dependency)).map((dependency) => externalIds.get(dependency)!), routeId: task.routeId, ...(task.shotId ? { shotId: task.shotId } : {}), ...(task.lane ? { lane: { laneId: task.lane.laneId, limit: task.lane.limit } } : {}), authorizationRequirement: task.authorizationRequirement, }; }); const compatibility = await enqueueCinemaCompatibilityPlan(root, { projectSlug, sourceFrontDoor: 'chain', generatedAt: source.updatedAt, tasks: compatibilityTasks, }); const replacementTaskIds = compatibility.receipt.mappings.map((mapping) => mapping.queueTaskId); await cancelCinemaChainFallbackTail(root, projectSlug, { taskIds: descendants.map((task) => task.taskId), transitionId, replacementTaskIds, evidenceId: compatibility.receipt.contentHash, now: source.updatedAt, }); if (compatibility.receipt.mappings.some((mapping) => mapping.enqueued)) { const workspace = await ensureProjectWorkspace(projectSlug, root); await appendProjectEvent(workspace, { type: 'cinema.chain.fallback.advanced', recordedAt: source.updatedAt, payload: { transitionId, sourceTaskId, fromPolicyIndex: sourceBinding.activeSourcePolicyIndex, activePolicyIndex: nextPolicyIndex, cancelledTaskIds: descendants.map((task) => task.taskId), replacementTaskIds, receiptId: compatibility.receipt.receiptId, receiptHash: compatibility.receipt.contentHash, }, }); } return { transitionId, sourceTaskId, fromPolicyIndex: sourceBinding.activeSourcePolicyIndex, activePolicyIndex: nextPolicyIndex, fromPolicy: sourceBinding.sourcePolicy[sourceBinding.activeSourcePolicyIndex], activePolicy: sourceBinding.sourcePolicy[nextPolicyIndex], cancelledTaskIds: descendants.map((task) => task.taskId), replacementTaskIds, receipt: compatibility.receipt, receiptPath: compatibility.receiptPath, status: await projectCinemaCompatibilityStatus(root, compatibility.receipt), providerCalls: 0, generationCalls: 0, spendAuthorized: false, }; }