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 { cinemaArtifactContentHash } from './cinema-provider-contracts.js'; import { withCinemaFileLock } from './cinema-file-lock.js'; import { enqueueCinemaTask, readCinemaProductionQueue, type EnqueueCinemaTaskInput } from './cinema-production-queue.js'; import type { CinemaQueueMoney, CinemaQueueTaskKind, CinemaQueueTaskStatus } from './cinema-queue-types.js'; import { cinemaArtifactsDirFor } from './cinema-store.js'; import { sha256Text, stableCinemaJson } from './cinema-shot-compiler.js'; import type { CinemaSha256 } from './cinema-types.js'; export type CinemaQueueFrontDoor = 'batch' | 'pool' | 'chain' | 'mograph' | 'cinema'; export interface CinemaQueueCompatibilityTask { externalTaskId: string; kind: CinemaQueueTaskKind; payload: EnqueueCinemaTaskInput['payload']; /** Already-persisted queue tasks that must succeed before this plan task. */ existingDependencies?: string[]; /** Dependencies expressed as externalTaskId values inside this plan. */ dependencies: string[]; routeId: string; shotId?: string; attemptId?: string; lane?: EnqueueCinemaTaskInput['lane']; estimatedCost?: CinemaQueueMoney; spendAuthorizationId?: string; authorizationRequirement?: 'none' | 'exact-quote'; } export interface CinemaQueueCompatibilityReceiptArtifact { schemaVersion: 1; receiptId: string; projectSlug: string; sourceFrontDoor: CinemaQueueFrontDoor; sourceArtifactHash: CinemaSha256; generatedAt: string; mappings: Array<{ externalTaskId: string; queueTaskId: string; executionHash: CinemaSha256; payloadHash: CinemaSha256; dependencies: string[]; kind: CinemaQueueTaskKind; routeId: string; enqueued: boolean; }>; queueHashAtEnqueue: CinemaSha256; providerCalls: 0; spendAuthorized: false; contentHash: CinemaSha256; } export interface CinemaQueueCompatibilityStatus { receiptId: string; projectSlug: string; sourceFrontDoor: CinemaQueueFrontDoor; receiptHash: CinemaSha256; currentQueueHash: CinemaSha256; tasks: Array<{ externalTaskId: string; queueTaskId: string; status: CinemaQueueTaskStatus; evidenceIds: string[] }>; } 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 compatibilityDir(root: string, projectSlug: string): string { return join(cinemaArtifactsDirFor(root, projectSlug), 'queue-compatibility'); } export function cinemaQueueCompatibilityReceiptPathFor(root: string, projectSlug: string, receiptId: string): string { return join(compatibilityDir(root, projectSlug), `${safe(receiptId, 'receiptId')}.json`); } function queueHash(queue: NonNullable>>): CinemaSha256 { return sha256Text(stableCinemaJson(queue)); } function planHash(input: { projectSlug: string; sourceFrontDoor: CinemaQueueFrontDoor; tasks: CinemaQueueCompatibilityTask[] }): CinemaSha256 { return sha256Text(stableCinemaJson(input)); } function topological(tasks: CinemaQueueCompatibilityTask[]): CinemaQueueCompatibilityTask[] { const byId = new Map(); for (const task of tasks) { if (!task.externalTaskId.trim()) throw new Error('Cinema compatibility external task IDs must be non-empty'); if (byId.has(task.externalTaskId)) throw new Error(`Cinema compatibility task ID ${task.externalTaskId} is duplicated`); byId.set(task.externalTaskId, task); } for (const task of tasks) { if (new Set(task.existingDependencies ?? []).size !== (task.existingDependencies ?? []).length) { throw new Error(`Cinema compatibility task ${task.externalTaskId} repeats an existing queue dependency`); } const missing = task.dependencies.filter((dependency) => !byId.has(dependency)); if (missing.length) throw new Error(`Cinema compatibility task ${task.externalTaskId} has unknown dependencies: ${missing.join(', ')}`); if (new Set(task.dependencies).size !== task.dependencies.length) throw new Error(`Cinema compatibility task ${task.externalTaskId} repeats a dependency`); } const pending = new Map(byId); const ordered: CinemaQueueCompatibilityTask[] = []; const done = new Set(); while (pending.size) { const ready = [...pending.values()].filter((task) => task.dependencies.every((dependency) => done.has(dependency))); if (!ready.length) throw new Error(`Cinema compatibility plan contains a dependency cycle: ${[...pending.keys()].join(', ')}`); ready.sort((left, right) => left.externalTaskId.localeCompare(right.externalTaskId)); for (const task of ready) { ordered.push(task); done.add(task.externalTaskId); pending.delete(task.externalTaskId); } } return ordered; } function queueTaskId(frontDoor: CinemaQueueFrontDoor, externalTaskId: string, sourceHash: CinemaSha256): string { const normalized = externalTaskId.replace(/[^A-Za-z0-9._-]/g, '-').replace(/^-+/, '') || 'task'; return `compat-${frontDoor}-${normalized}-${sourceHash.slice(7, 15)}`; } async function writeReceipt(root: string, receipt: CinemaQueueCompatibilityReceiptArtifact): Promise { const path = cinemaQueueCompatibilityReceiptPathFor(root, receipt.projectSlug, receipt.receiptId); await mkdir(dirname(path), { recursive: true }); return withCinemaFileLock(path, async () => { if (existsSync(path)) { const existing = JSON.parse(await readFile(path, 'utf8')) as CinemaQueueCompatibilityReceiptArtifact; if (stableCinemaJson(existing) === stableCinemaJson(receipt)) return path; throw new Error(`Cinema queue compatibility receipt ${receipt.receiptId} is immutable`); } await writeTextFileAtomic(path, `${JSON.stringify(receipt, null, 2)}\n`); return path; }); } export async function enqueueCinemaCompatibilityPlan(root: string, input: { projectSlug: string; sourceFrontDoor: CinemaQueueFrontDoor; generatedAt: string; tasks: CinemaQueueCompatibilityTask[]; }): Promise<{ receipt: CinemaQueueCompatibilityReceiptArtifact; receiptPath: string }> { if (!input.tasks.length) throw new Error('Cinema queue compatibility plan requires at least one task'); const sourceArtifactHash = planHash({ projectSlug: input.projectSlug, sourceFrontDoor: input.sourceFrontDoor, tasks: input.tasks }); const mapped = new Map(); const mappings: CinemaQueueCompatibilityReceiptArtifact['mappings'] = []; for (const task of topological(input.tasks)) { const dependencies = [ ...(task.existingDependencies ?? []), ...task.dependencies.map((dependency) => mapped.get(dependency)!), ]; if (new Set(dependencies).size !== dependencies.length) { throw new Error(`Cinema compatibility task ${task.externalTaskId} repeats a resolved dependency`); } const result = await enqueueCinemaTask(root, input.projectSlug, { taskId: queueTaskId(input.sourceFrontDoor, task.externalTaskId, sourceArtifactHash), kind: task.kind, payload: task.payload, dependencies, routeId: task.routeId, ...(task.shotId ? { shotId: task.shotId } : {}), ...(task.attemptId ? { attemptId: task.attemptId } : {}), ...(task.lane ? { lane: task.lane } : {}), ...(task.estimatedCost ? { estimatedCost: task.estimatedCost } : {}), ...(task.spendAuthorizationId ? { spendAuthorizationId: task.spendAuthorizationId } : {}), ...(task.authorizationRequirement ? { authorizationRequirement: task.authorizationRequirement } : {}), createdAt: input.generatedAt, }); mapped.set(task.externalTaskId, result.task.taskId); mappings.push({ externalTaskId: task.externalTaskId, queueTaskId: result.task.taskId, executionHash: result.task.executionHash, payloadHash: result.task.payloadHash, dependencies: [...result.task.dependencies], kind: result.task.kind, routeId: result.task.routeId, enqueued: result.enqueued }); } const queue = await readCinemaProductionQueue(root, input.projectSlug); if (!queue) throw new Error('Cinema compatibility adapter did not persist the durable queue'); const receiptEvent = { schemaVersion: 1 as const, projectSlug: input.projectSlug, sourceFrontDoor: input.sourceFrontDoor, sourceArtifactHash, generatedAt: input.generatedAt, mappings, queueHashAtEnqueue: queueHash(queue), providerCalls: 0 as const, spendAuthorized: false as const, }; const draft = { ...receiptEvent, receiptId: `compat-${input.sourceFrontDoor}-${sha256Text(stableCinemaJson(receiptEvent)).slice(7, 23)}`, }; const receipt = { ...draft, contentHash: sha256Text('placeholder') }; receipt.contentHash = cinemaArtifactContentHash(receipt); const receiptPath = await writeReceipt(root, receipt); return { receipt, receiptPath }; } export async function readCinemaQueueCompatibilityReceipt(root: string, projectSlug: string, receiptId: string): Promise { const path = cinemaQueueCompatibilityReceiptPathFor(root, projectSlug, receiptId); if (!existsSync(path)) return null; const receipt = JSON.parse(await readFile(path, 'utf8')) as CinemaQueueCompatibilityReceiptArtifact; if (receipt.projectSlug !== projectSlug || receipt.receiptId !== receiptId || receipt.contentHash !== cinemaArtifactContentHash(receipt)) throw new Error(`Cinema queue compatibility receipt ${receiptId} failed identity or content-hash validation`); return receipt; } export async function projectCinemaCompatibilityStatus(root: string, receipt: CinemaQueueCompatibilityReceiptArtifact): Promise { const queue = await readCinemaProductionQueue(root, receipt.projectSlug); if (!queue) throw new Error(`Cinema durable queue for ${receipt.projectSlug} is missing`); const tasks = new Map(queue.tasks.map((task) => [task.taskId, task])); return { receiptId: receipt.receiptId, projectSlug: receipt.projectSlug, sourceFrontDoor: receipt.sourceFrontDoor, receiptHash: receipt.contentHash, currentQueueHash: queueHash(queue), tasks: receipt.mappings.map((mapping) => { const task = tasks.get(mapping.queueTaskId); if (!task) throw new Error(`Cinema compatibility queue task ${mapping.queueTaskId} is missing`); return { externalTaskId: mapping.externalTaskId, queueTaskId: task.taskId, status: task.status, evidenceIds: [...task.evidenceIds] }; }), }; }