import { createHash } from 'node:crypto'; import { constants } from 'node:fs'; import { access, mkdir, readFile, rename, rm, writeFile } from 'node:fs/promises'; import { dirname, join } from 'node:path'; import { claimCinemaTask, completeCinemaTask, failCinemaTask, readCinemaProductionQueue, } from './cinema-production-queue.js'; import type { CinemaQueueTask } from './cinema-queue-types.js'; import { appendProjectEvent } from './events.js'; import { runFfmpeg } from './assemble/ffmpeg.js'; import { ensureProjectWorkspace, resolveProjectWorkspace } from './workspace.js'; export interface CinemaCompatibilityLocalWorkerResult { disposition: 'completed' | 'already-succeeded'; task: CinemaQueueTask; artifactPaths: string[]; providerCalls: 0; generationCalls: 0; spendAuthorized: false; } type FfmpegRunner = (args: string[]) => Promise<{ exitCode: number }>; 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 safeId(value: unknown, label: string): string { if (typeof value !== 'string' || !/^[A-Za-z0-9][A-Za-z0-9._-]*$/.test(value)) throw new Error(`${label} must be a path-safe identifier`); return value; } function stringArray(value: unknown, label: string): string[] { if (!Array.isArray(value) || value.length === 0 || value.some((entry) => typeof entry !== 'string' || !entry.trim())) { throw new Error(`${label} must contain non-empty strings`); } return [...value] as string[]; } function sha256(bytes: Uint8Array): `sha256:${string}` { return `sha256:${createHash('sha256').update(bytes).digest('hex')}`; } async function exists(path: string): Promise { try { await access(path, constants.F_OK); return true; } catch { return false; } } async function writeBytesIdempotent(path: string, bytes: Uint8Array): Promise<`sha256:${string}`> { await mkdir(dirname(path), { recursive: true }); if (await exists(path)) { const present = await readFile(path); if (!Buffer.from(present).equals(Buffer.from(bytes))) throw new Error(`Local compatibility artifact already exists with different bytes: ${path}`); return sha256(present); } try { await writeFile(path, bytes, { flag: 'wx' }); } catch (error) { if ((error as NodeJS.ErrnoException).code !== 'EEXIST') throw error; const present = await readFile(path); if (!Buffer.from(present).equals(Buffer.from(bytes))) throw new Error(`Concurrent local compatibility artifact differs: ${path}`); } return sha256(bytes); } function mographPostprocessBinding(task: CinemaQueueTask): { blockId: string; renderGroupId: string; sidecar: Record; } { const binding = record(task.payload.parameters.mographPostprocess, 'Mograph post-process binding'); const blockId = safeId(binding.blockId, 'Mograph blockId'); const renderGroupId = safeId(binding.renderGroupId, 'Mograph renderGroupId'); const sidecar = record(binding.sidecar, 'Mograph sidecar'); if (sidecar.blockId !== blockId || sidecar.schemaVersion !== 1) throw new Error('Mograph sidecar does not match its block binding'); return { blockId, renderGroupId, sidecar: structuredClone(sidecar) }; } async function materializeMographBlock(root: string, projectSlug: string, task: CinemaQueueTask, queue: NonNullable>>): Promise<{ paths: string[]; evidenceIds: string[] }> { const binding = mographPostprocessBinding(task); if (task.dependencies.length !== 1) throw new Error(`Mograph post-process ${task.taskId} requires exactly one render dependency`); const source = queue.tasks.find((entry) => entry.taskId === task.dependencies[0]); if (source?.status !== 'succeeded' || !source.execution?.download.finalPath || !source.execution.download.contentHash) { throw new Error(`Mograph post-process ${task.taskId} lacks a succeeded render dependency`); } const sourceBytes = await readFile(source.execution.download.finalPath); if (sha256(sourceBytes) !== source.execution.download.contentHash) throw new Error(`Mograph render dependency ${source.taskId} changed on disk`); const projectDir = resolveProjectWorkspace(projectSlug, root).projectDir; const clipPath = join(projectDir, 'mograph', 'clips', `${binding.blockId}.mp4`); const sidecarPath = join(projectDir, 'artifacts', 'mograph-sidecars', `${binding.blockId}.json`); const clipHash = await writeBytesIdempotent(clipPath, sourceBytes); const sidecarHash = await writeBytesIdempotent(sidecarPath, Buffer.from(`${JSON.stringify(binding.sidecar, null, 2)}\n`, 'utf8')); return { paths: [clipPath, sidecarPath], evidenceIds: [source.execution.download.contentHash, clipHash, sidecarHash] }; } function concatLine(path: string): string { return `file '${path.replaceAll("'", "'\\''")}'`; } async function assembleMographMaster( root: string, projectSlug: string, task: CinemaQueueTask, queue: NonNullable>>, ffmpeg: FfmpegRunner, ): Promise<{ paths: string[]; evidenceIds: string[] }> { const binding = record(task.payload.parameters.mographAssemble, 'Mograph assembly binding'); const renderGroupId = safeId(binding.renderGroupId, 'Mograph renderGroupId'); const blockIds = stringArray(binding.blockIds, 'Mograph assembly blockIds').map((value) => safeId(value, 'Mograph blockId')); if (task.dependencies.length !== blockIds.length || task.dependencies.some((id) => queue.tasks.find((entry) => entry.taskId === id)?.status !== 'succeeded')) { throw new Error(`Mograph assembly ${task.taskId} requires every post-process dependency to succeed`); } const projectDir = resolveProjectWorkspace(projectSlug, root).projectDir; const clips = blockIds.map((blockId) => join(projectDir, 'mograph', 'clips', `${blockId}.mp4`)); const clipHashes: string[] = []; for (const clip of clips) clipHashes.push(sha256(await readFile(clip))); const masterPath = join(projectDir, 'final', 'videos', 'mograph-master.mp4'); const receiptPath = join(projectDir, 'artifacts', 'mograph-sidecars', `assembly-${renderGroupId}.json`); const temporaryMasterPath = join(projectDir, 'final', 'videos', `mograph-master-${safeId(task.taskId, 'Mograph assembly taskId')}.tmp.mp4`); if (await exists(receiptPath)) { const receipt = JSON.parse(await readFile(receiptPath, 'utf8')) as Record; if (receipt.renderGroupId !== renderGroupId || JSON.stringify(receipt.clipHashes) !== JSON.stringify(clipHashes) || typeof receipt.masterHash !== 'string') { throw new Error('Existing mograph assembly receipt does not match this task'); } if (await exists(masterPath)) { const masterHash = sha256(await readFile(masterPath)); if (receipt.masterHash !== masterHash) throw new Error('Existing mograph master does not match its assembly receipt'); return { paths: [masterPath, receiptPath], evidenceIds: [...clipHashes, masterHash, sha256(await readFile(receiptPath))] }; } if (await exists(temporaryMasterPath)) { const masterHash = sha256(await readFile(temporaryMasterPath)); if (receipt.masterHash !== masterHash) throw new Error('Recoverable mograph temporary master does not match its assembly receipt'); await mkdir(dirname(masterPath), { recursive: true }); await rename(temporaryMasterPath, masterPath); return { paths: [masterPath, receiptPath], evidenceIds: [...clipHashes, masterHash, sha256(await readFile(receiptPath))] }; } throw new Error('Mograph assembly receipt exists but neither its final nor temporary master exists'); } if (await exists(masterPath)) throw new Error(`Mograph master already exists without an assembly receipt: ${masterPath}`); await rm(temporaryMasterPath, { force: true }); const listPath = join(projectDir, 'mograph', 'clips', `concat-${renderGroupId}.txt`); await writeBytesIdempotent(listPath, Buffer.from(`${clips.map(concatLine).join('\n')}\n`, 'utf8')); await mkdir(dirname(masterPath), { recursive: true }); const result = await ffmpeg(['-f', 'concat', '-safe', '0', '-i', listPath, '-c:v', 'libx264', '-preset', 'medium', '-crf', '18', '-pix_fmt', 'yuv420p', '-c:a', 'aac', '-ar', '44100', temporaryMasterPath]); if (result.exitCode !== 0 || !(await exists(temporaryMasterPath))) throw new Error('Mograph assembly did not produce a temporary master'); const masterHash = sha256(await readFile(temporaryMasterPath)); const receipt = { schemaVersion: 1, renderGroupId, blockIds, clipHashes, masterPath, masterHash }; const receiptHash = await writeBytesIdempotent(receiptPath, Buffer.from(`${JSON.stringify(receipt, null, 2)}\n`, 'utf8')); await rename(temporaryMasterPath, masterPath); return { paths: [masterPath, receiptPath], evidenceIds: [...clipHashes, masterHash, receiptHash] }; } export async function runCinemaCompatibilityLocalWorkerOnce(options: { root: string; projectSlug: string; taskId: string; workerId: string; now?: string; ffmpeg?: FfmpegRunner; }): Promise { const now = options.now ?? new Date().toISOString(); const beforeQueue = await readCinemaProductionQueue(options.root, options.projectSlug); const before = beforeQueue?.tasks.find((entry) => entry.taskId === options.taskId); if (!before) throw new Error(`Compatibility local task ${options.taskId} does not exist`); if (!['post-process', 'assemble'].includes(before.kind) || before.routeId !== 'local-mograph') { throw new Error(`Compatibility local worker cannot execute ${before.kind} on ${before.routeId}`); } if (before.status === 'succeeded') { return { disposition: 'already-succeeded', task: before, artifactPaths: [], providerCalls: 0, generationCalls: 0, spendAuthorized: false }; } const claimed = await claimCinemaTask(options.root, options.projectSlug, { taskId: options.taskId, workerId: options.workerId, now }); if (!claimed.task?.lease || claimed.task.status !== 'leased') throw new Error(`Compatibility local task ${options.taskId} is not ready to run`); const leaseId = claimed.task.lease.leaseId; try { const queue = await readCinemaProductionQueue(options.root, options.projectSlug); if (!queue) throw new Error(`Compatibility queue ${options.projectSlug} disappeared`); const materialized = claimed.task.kind === 'post-process' ? await materializeMographBlock(options.root, options.projectSlug, claimed.task, queue) : await assembleMographMaster(options.root, options.projectSlug, claimed.task, queue, options.ffmpeg ?? (async (args) => runFfmpeg(args))); const completed = await completeCinemaTask(options.root, options.projectSlug, { taskId: options.taskId, leaseId, evidenceIds: materialized.evidenceIds, actualCost: { currency: claimed.task.estimatedCost.currency, amount: 0 }, now, }); const workspace = await ensureProjectWorkspace(options.projectSlug, options.root); await appendProjectEvent(workspace, { type: 'cinema.compatibility.local.completed', recordedAt: now, payload: { taskId: completed.task.taskId, kind: completed.task.kind, artifactPaths: materialized.paths, evidenceIds: materialized.evidenceIds }, }); return { disposition: 'completed', task: completed.task, artifactPaths: materialized.paths, providerCalls: 0, generationCalls: 0, spendAuthorized: false }; } catch (error) { await failCinemaTask(options.root, options.projectSlug, { taskId: options.taskId, leaseId, failure: { code: 'local-post-process-failed', message: error instanceof Error ? error.message : String(error), retryable: false }, now, }); throw error; } }