import { sha256Text, stableCinemaJson } from './cinema-shot-compiler.js'; import type { CinemaProductionQueueArtifact, CinemaQueueTask } from './cinema-queue-types.js'; import type { CinemaContractIssue, CinemaContractValidation } from './cinema-types.js'; const SHA256_PATTERN = /^sha256:[a-f0-9]{64}$/; function add(issues: CinemaContractIssue[], code: string, path: string, message: string): void { issues.push({ code, path, message }); } function requireText(issues: CinemaContractIssue[], value: string, path: string): void { if (value.trim().length === 0) add(issues, 'required-text', path, 'must be a non-empty string'); } function duplicates(values: string[]): string[] { const seen = new Set(); return [...new Set(values.filter((value) => seen.has(value) || !seen.add(value)))]; } function dependenciesSucceeded(task: CinemaQueueTask, tasks: Map): boolean { return task.dependencies.every((taskId) => tasks.get(taskId)?.status === 'succeeded'); } export function cinemaQueueExecutionHash(input: Pick & { lane: { laneId: string; limit: number | null } | null }): `sha256:${string}` { return sha256Text(stableCinemaJson({ kind: input.kind, payloadHash: input.payloadHash, dependencies: input.dependencies, routeId: input.routeId, shotId: input.shotId ?? null, attemptId: input.attemptId ?? null, lane: input.lane ? { laneId: input.lane.laneId, limit: input.lane.limit } : null, authorizationRequirement: input.authorizationRequirement, })); } function validateTask( task: CinemaQueueTask, index: number, tasks: Map, issues: CinemaContractIssue[], ): void { const path = `tasks[${index}]`; [task.taskId, task.routeId].forEach((value, fieldIndex) => requireText(issues, value, `${path}.${['taskId', 'routeId'][fieldIndex]}`)); if (!SHA256_PATTERN.test(task.payloadHash)) { add(issues, 'invalid-content-hash', `${path}.payloadHash`, 'must be a lowercase sha256 content hash'); } else if (task.payloadHash !== sha256Text(stableCinemaJson(task.payload))) { add(issues, 'queue-payload-hash-mismatch', `${path}.payloadHash`, 'must match the canonical immutable payload'); } if (!SHA256_PATTERN.test(task.executionHash)) { add(issues, 'invalid-content-hash', `${path}.executionHash`, 'must be a lowercase sha256 execution hash'); } else if (task.executionHash !== cinemaQueueExecutionHash(task)) { add(issues, 'queue-execution-hash-mismatch', `${path}.executionHash`, 'must bind payload, route, dependencies, lane and authorization semantics'); } const repeatedDependencies = duplicates(task.dependencies); if (repeatedDependencies.length > 0) { add(issues, 'duplicate-dependency', `${path}.dependencies`, `contains duplicate task ids: ${repeatedDependencies.join(', ')}`); } task.dependencies.forEach((dependency, dependencyIndex) => { if (dependency === task.taskId) add(issues, 'self-dependency', `${path}.dependencies[${dependencyIndex}]`, 'a task cannot depend on itself'); else if (!tasks.has(dependency)) add(issues, 'missing-dependency', `${path}.dependencies[${dependencyIndex}]`, 'must reference a task in this queue'); }); const ready = dependenciesSucceeded(task, tasks); if (task.status === 'ready' && !ready) add(issues, 'task-ready-too-early', `${path}.status`, 'ready tasks require every dependency to have succeeded'); if (task.status === 'blocked' && ready) add(issues, 'task-blocked-without-dependency', `${path}.status`, 'blocked tasks must have an incomplete dependency'); if (task.status === 'awaiting-quote') { if (!ready) add(issues, 'quote-gate-too-early', `${path}.status`, 'awaiting-quote requires every dependency to have succeeded'); if (task.authorizationRequirement !== 'exact-quote' || task.spendAuthorizationId !== null) add(issues, 'invalid-quote-gate', `${path}.authorizationRequirement`, 'awaiting-quote requires an unresolved exact-quote authorization'); } if (task.authorizationRequirement === 'none' && task.spendAuthorizationId !== null) add(issues, 'unexpected-spend-authorization', `${path}.spendAuthorizationId`, 'provider-free tasks cannot carry spend authorization'); if (task.authorizationRequirement === 'none' && (task.estimatedCost.amount > 0 || task.actualCost.amount > 0)) add(issues, 'provider-free-spend', `${path}.authorizationRequirement`, 'provider-free tasks cannot estimate or record spend'); if (task.authorizationRequirement === 'exact-quote' && task.status === 'ready' && !task.spendAuthorizationId) add(issues, 'missing-ready-authorization', `${path}.spendAuthorizationId`, 'exact-quote tasks cannot become ready without authorization'); if (task.status === 'waiting-lane') { if (!task.lane || task.lane.ticketState !== 'queued' || !task.lane.ticketId) { add(issues, 'waiting-lane-ticket', `${path}.lane`, 'waiting-lane requires one exact queued LaneQueue ticket'); } if (task.lease !== null) add(issues, 'waiting-lane-lease', `${path}.lease`, 'waiting-lane cannot also own a worker lease'); } const providerActive = ['leased', 'submitting', 'reconciling', 'submitted', 'polling', 'downloading']; if (providerActive.includes(task.status)) { if (!task.lease) add(issues, 'missing-worker-lease', `${path}.lease`, 'leased task requires a worker lease'); if (task.lane && (task.lane.ticketState !== 'held' || !task.lane.ticketId)) { add(issues, 'missing-held-lane-ticket', `${path}.lane`, 'lane-bound leased task requires its exact held ticket'); } } else if (task.lease !== null) { add(issues, 'unexpected-worker-lease', `${path}.lease`, `${task.status} task cannot retain a worker lease`); } const executionRequired = ['submitting', 'reconciling', 'submitted', 'polling', 'downloading']; if (executionRequired.includes(task.status) && !task.execution) { add(issues, 'missing-provider-execution', `${path}.execution`, `${task.status} task requires a durable provider execution receipt`); } if (task.execution) { const execution = task.execution; if (execution.submissionKey !== task.payloadHash) { add(issues, 'submission-key-mismatch', `${path}.execution.submissionKey`, 'must equal the immutable queue payload hash'); } requireText(issues, execution.providerStatus, `${path}.execution.providerStatus`); if (execution.providerStatusHash !== sha256Text(stableCinemaJson(execution.providerStatusEnvelope))) { add(issues, 'provider-status-hash-mismatch', `${path}.execution.providerStatusHash`, 'must hash the exact provider status envelope'); } if (['submitted', 'polling', 'downloading'].includes(task.status) && !execution.providerJobId) { add(issues, 'missing-provider-job-id', `${path}.execution.providerJobId`, `${task.status} task requires a durable provider job id`); } if (task.status === 'reconciling' && !execution.reconciliation.required) { add(issues, 'reconciliation-not-required', `${path}.execution.reconciliation.required`, 'reconciling status must preserve the unresolved submission flag'); } if (execution.reconciliation.required && !execution.reconciliation.reason) { add(issues, 'missing-reconciliation-reason', `${path}.execution.reconciliation.reason`, 'required reconciliation needs a reason'); } const download = execution.download; const completedValues = [download.finalPath, download.contentHash, download.byteLength, download.completedAt]; if (completedValues.some((value) => value !== null) && completedValues.some((value) => value === null)) { add(issues, 'partial-download-receipt', `${path}.execution.download`, 'final path, content hash, byte length and completion time must become durable together'); } } if (task.lane) { requireText(issues, task.lane.laneId, `${path}.lane.laneId`); if (task.lane.limit !== null && (!Number.isInteger(task.lane.limit) || task.lane.limit < 1)) { add(issues, 'invalid-lane-limit', `${path}.lane.limit`, 'must be null or a positive integer'); } if (task.lane.ticketState === 'none' && task.lane.ticketId !== null) { add(issues, 'orphan-lane-ticket', `${path}.lane.ticketId`, 'must be null when ticketState is none'); } if (task.lane.ticketState !== 'none' && !task.lane.ticketId) { add(issues, 'missing-lane-ticket', `${path}.lane.ticketId`, 'must identify the exact queued or held ticket'); } } if (task.estimatedCost.amount < 0 || task.actualCost.amount < 0) { add(issues, 'negative-cost', path, 'estimated and actual costs cannot be negative'); } if (task.estimatedCost.currency !== task.actualCost.currency) { add(issues, 'cost-currency-mismatch', path, 'estimated and actual cost currencies must match'); } if ((task.estimatedCost.amount > 0 || task.actualCost.amount > 0) && [...providerActive, 'succeeded'].includes(task.status) && !task.spendAuthorizationId) { add(issues, 'paid-task-unauthorized', `${path}.spendAuthorizationId`, 'paid work cannot be leased or completed without an exact authorization id'); } if (['failed', 'dead-letter'].includes(task.status) && task.failure === null) { add(issues, 'missing-task-failure', `${path}.failure`, 'failed task must preserve a structured failure'); } if (!['failed', 'dead-letter', 'retry-ready'].includes(task.status) && task.failure !== null) { add(issues, 'unexpected-task-failure', `${path}.failure`, 'only failure or retry states can carry a structured failure'); } if (task.status === 'succeeded' && task.evidenceIds.length === 0) { add(issues, 'missing-task-evidence', `${path}.evidenceIds`, 'succeeded task must preserve output or QC evidence'); } } function detectCycles(tasks: Map, issues: CinemaContractIssue[]): void { const visiting = new Set(); const visited = new Set(); const walk = (taskId: string, trail: string[]): void => { if (visiting.has(taskId)) { add(issues, 'dependency-cycle', `tasks.${taskId}.dependencies`, `dependency cycle: ${[...trail, taskId].join(' -> ')}`); return; } if (visited.has(taskId)) return; visiting.add(taskId); const task = tasks.get(taskId); task?.dependencies.filter((dependency) => tasks.has(dependency)).forEach((dependency) => walk(dependency, [...trail, taskId])); visiting.delete(taskId); visited.add(taskId); }; tasks.forEach((_task, taskId) => walk(taskId, [])); } export function validateCinemaProductionQueue( queue: CinemaProductionQueueArtifact, ): CinemaContractValidation { const issues: CinemaContractIssue[] = []; requireText(issues, queue.projectSlug, 'projectSlug'); const repeatedIds = duplicates(queue.tasks.map((task) => task.taskId)); if (repeatedIds.length > 0) add(issues, 'duplicate-task-id', 'tasks.taskId', `contains duplicate ids: ${repeatedIds.join(', ')}`); const repeatedExecutions = duplicates(queue.tasks.map((task) => task.executionHash)); if (repeatedExecutions.length > 0) add(issues, 'duplicate-active-payload', 'tasks.executionHash', `equivalent execution contracts must deduplicate: ${repeatedExecutions.join(', ')}`); const tasks = new Map(queue.tasks.map((task) => [task.taskId, task])); queue.tasks.forEach((task, index) => validateTask(task, index, tasks, issues)); detectCycles(tasks, issues); return { ok: issues.length === 0, issues }; }