import { createHash } 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 { cinemaArtifactContentHash } from './cinema-provider-contracts.js'; import { readCinemaProductionQueue } from './cinema-production-queue.js'; import { appendCinemaAttemptOutcome, readCinemaShotManifest, cinemaArtifactsDirFor } from './cinema-store.js'; import { sha256Text, stableCinemaJson } from './cinema-shot-compiler.js'; import type { CinemaQueueTask } from './cinema-queue-types.js'; import type { CinemaSha256 } from './cinema-types.js'; import { appendProjectEvent } from './events.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 type { SceneCandidate } from './types.js'; import { ensureProjectWorkspace } from './workspace.js'; export interface CinemaCandidateBridgeReceiptArtifact { schemaVersion: 1; receiptId: string; projectSlug: string; queueTaskId: string; queueExecutionHash: CinemaSha256; queuePayloadHash: CinemaSha256; attemptId: string; outcomeId: string; outcomeHash: CinemaSha256; shotId: string; shotManifestHash: CinemaSha256; /** Unique cut/assembly coordinate derived from the shot ordinal. */ sceneIndex: number; /** Narrative scene coordinate retained for story/continuity projection. */ storySceneIndex: number; candidateId: string; providerJobId: string; media: { kind: 'video'; path: string; contentHash: CinemaSha256; byteLength: number; durationSeconds: number; }; routeId: string; actualCost: { currency: string; amount: number }; evidenceIds: string[]; ingestedAt: string; providerCalls: 0; generationCalls: 0; spendAuthorized: false; contentHash: CinemaSha256; } export interface CinemaCandidateBridgeResult { appended: boolean; receipt: CinemaCandidateBridgeReceiptArtifact; receiptPath: string; candidate: SceneCandidate; } function safe(value: string, label: string): string { if (!/^[A-Za-z0-9][A-Za-z0-9._-]*$/.test(value)) throw new Error(`${label} must be a path-safe identifier`); return value; } function bytesHash(bytes: Uint8Array): CinemaSha256 { return `sha256:${createHash('sha256').update(bytes).digest('hex')}`; } export function cinemaCandidateBridgeReceiptPathFor(root: string, projectSlug: string, taskId: string): string { return join(cinemaArtifactsDirFor(root, projectSlug), 'candidate-bridges', `${safe(taskId, 'taskId')}.json`); } async function readReceipt(path: string): Promise { if (!existsSync(path)) return null; const receipt = JSON.parse(await readFile(path, 'utf8')) as CinemaCandidateBridgeReceiptArtifact; if (receipt.contentHash !== cinemaArtifactContentHash(receipt)) throw new Error(`Cinema candidate bridge receipt failed content-hash validation: ${path}`); return receipt; } export async function readCinemaCandidateBridgeReceipt( root: string, projectSlug: string, taskId: string, ): Promise { const receipt = await readReceipt(cinemaCandidateBridgeReceiptPathFor(root, projectSlug, taskId)); if (!receipt) return null; if (receipt.projectSlug !== projectSlug || receipt.queueTaskId !== taskId) throw new Error(`Cinema candidate bridge receipt does not match ${projectSlug}/${taskId}`); return receipt; } function assertSucceededVideoTask(task: CinemaQueueTask): asserts task is CinemaQueueTask & { shotId: string; attemptId: string; execution: NonNullable & { providerJobId: string; download: { finalPath: string; contentHash: CinemaSha256; byteLength: number; completedAt: string }; }; } { if (task.kind !== 'video-generate') throw new Error(`Cinema task ${task.taskId} is ${task.kind}, not a video generation result`); if (task.status !== 'succeeded') throw new Error(`Cinema task ${task.taskId} is ${task.status}, not succeeded`); if (!task.shotId || !task.attemptId) throw new Error(`Cinema task ${task.taskId} lacks immutable shot or attempt lineage`); if (!task.execution?.providerJobId || !task.execution.download.finalPath || !task.execution.download.contentHash || task.execution.download.byteLength === null || !task.execution.download.completedAt) { throw new Error(`Cinema task ${task.taskId} lacks a complete provider/download receipt`); } } function findProviderCandidate( scenes: Awaited>, providerJobId: string, ): { sceneIndex: number; candidate: SceneCandidate } | null { for (const scene of scenes.scenes) { const candidate = scene.candidates.find((entry) => entry.source.externalJobId === providerJobId); if (candidate) return { sceneIndex: scene.sceneIndex, candidate }; } return null; } export async function bridgeCinemaTaskToCandidate( root: string, projectSlug: string, taskId: string, ingestedAt = new Date().toISOString(), ): Promise { const queue = await readCinemaProductionQueue(root, projectSlug); const task = queue?.tasks.find((entry) => entry.taskId === taskId); if (!task) throw new Error(`Cinema queue task ${taskId} does not exist`); assertSucceededVideoTask(task); const manifest = await readCinemaShotManifest(root, projectSlug, task.shotId); if (!manifest) throw new Error(`Cinema shot manifest ${task.shotId} does not exist`); if (manifest.sceneIndex < 0 || !Number.isInteger(manifest.sceneIndex)) throw new Error(`Cinema shot ${task.shotId} has an invalid scene index`); if (manifest.ordinal < 1 || !Number.isInteger(manifest.ordinal)) throw new Error(`Cinema shot ${task.shotId} has an invalid editorial ordinal`); const editorialIndex = manifest.ordinal - 1; const mediaBytes = await readFile(task.execution.download.finalPath); const observedHash = bytesHash(mediaBytes); if (observedHash !== task.execution.download.contentHash || mediaBytes.byteLength !== task.execution.download.byteLength) { throw new Error(`Cinema result bytes changed after queue completion for task ${taskId}`); } const outcomeInput = { outcomeId: `outcome-${taskId}`, attemptId: task.attemptId, shotId: task.shotId, queueTaskId: taskId, executionHash: task.executionHash, payloadHash: task.payloadHash, status: 'completed' as const, observedAt: task.execution.download.completedAt, providerJobId: task.execution.providerJobId, media: { path: task.execution.download.finalPath, contentHash: observedHash, byteLength: mediaBytes.byteLength }, actualCost: task.actualCost, evidenceIds: [...task.evidenceIds], }; const outcome = { ...outcomeInput, contentHash: sha256Text(stableCinemaJson(outcomeInput)) }; await appendCinemaAttemptOutcome(root, projectSlug, outcome, ingestedAt); const receiptPath = cinemaCandidateBridgeReceiptPathFor(root, projectSlug, taskId); return withCinemaFileLock(receiptPath, async () => { const existingReceipt = await readReceipt(receiptPath); if (existingReceipt) { if (existingReceipt.projectSlug !== projectSlug || existingReceipt.queueTaskId !== taskId || existingReceipt.queueExecutionHash !== task.executionHash || existingReceipt.queuePayloadHash !== task.payloadHash || existingReceipt.media.contentHash !== observedHash || existingReceipt.outcomeHash !== outcome.contentHash || existingReceipt.sceneIndex !== editorialIndex || existingReceipt.storySceneIndex !== manifest.sceneIndex) { throw new Error(`Cinema candidate bridge receipt ${existingReceipt.receiptId} no longer matches its queue task`); } const candidates = await readSceneCandidatesArtifact(root, projectSlug); const found = findProviderCandidate(candidates, task.execution.providerJobId); if (!found || found.candidate.id !== existingReceipt.candidateId || found.sceneIndex !== existingReceipt.sceneIndex) { throw new Error(`Cinema candidate bridge receipt ${existingReceipt.receiptId} points to a missing or changed candidate`); } return { appended: false, receipt: existingReceipt, receiptPath, candidate: found.candidate }; } let candidate!: SceneCandidate; let appended = false; await withSceneArtifactsLock(root, projectSlug, async () => { const candidates = await readSceneCandidatesArtifact(root, projectSlug); const selection = await readSceneSelectionArtifact(root, projectSlug); const found = findProviderCandidate(candidates, task.execution.providerJobId); if (found) { if (found.sceneIndex !== editorialIndex || found.candidate.outputs.length !== 1 || found.candidate.outputs[0].kind !== 'video' || found.candidate.outputs[0].path !== task.execution.download.finalPath) { throw new Error(`Provider job ${task.execution.providerJobId} is already bound to a different candidate result`); } candidate = found.candidate; const pending = markPending(selection, editorialIndex, [candidate.id]); await writeSceneSelectionArtifact(root, projectSlug, pending); return; } const generationRound = maxRoundForScene(candidates, editorialIndex) + 1; candidate = { id: nextCandidateId(candidates, editorialIndex), generationRound, prompt: manifest.compiledPrompt.text, route: task.routeId, submittedAt: task.execution.submittedAt ?? task.createdAt, completedAt: task.execution.download.completedAt, status: 'completed', outputs: [{ kind: 'video', path: task.execution.download.finalPath, durationSec: manifest.durationSeconds }], source: { executionRound: generationRound, adapter: 'native', externalJobId: task.execution.providerJobId, chainedFromCandidateId: null, }, }; await writeSceneCandidatesArtifact(root, projectSlug, appendCandidate(candidates, editorialIndex, candidate)); await writeSceneSelectionArtifact(root, projectSlug, markPending(selection, editorialIndex, [candidate.id])); appended = true; }); const receiptInput: Omit = { schemaVersion: 1, receiptId: `candidate-bridge-${taskId}`, projectSlug, queueTaskId: taskId, queueExecutionHash: task.executionHash, queuePayloadHash: task.payloadHash, attemptId: task.attemptId, outcomeId: outcome.outcomeId, outcomeHash: outcome.contentHash, shotId: task.shotId, shotManifestHash: sha256Text(stableCinemaJson(manifest)), sceneIndex: editorialIndex, storySceneIndex: manifest.sceneIndex, candidateId: candidate.id, providerJobId: task.execution.providerJobId, media: { kind: 'video', path: task.execution.download.finalPath, contentHash: observedHash, byteLength: mediaBytes.byteLength, durationSeconds: manifest.durationSeconds, }, routeId: task.routeId, actualCost: task.actualCost, evidenceIds: [...new Set([...task.evidenceIds, outcome.contentHash])], ingestedAt, providerCalls: 0, generationCalls: 0, spendAuthorized: false, }; const receipt = { ...receiptInput, contentHash: cinemaArtifactContentHash({ ...receiptInput, contentHash: sha256Text('placeholder') }) }; await mkdir(dirname(receiptPath), { recursive: true }); await writeTextFileAtomic(receiptPath, `${JSON.stringify(receipt, null, 2)}\n`); if (appended) { const workspace = await ensureProjectWorkspace(projectSlug, root); await appendProjectEvent(workspace, { type: 'cinema.candidate.ingested', recordedAt: ingestedAt, payload: { taskId, attemptId: task.attemptId, shotId: task.shotId, editorialIndex, storySceneIndex: manifest.sceneIndex, candidateId: candidate.id, receiptPath, contentHash: receipt.contentHash }, }); } return { appended, receipt, receiptPath, candidate }; }); }