import { createHash } from 'node:crypto'; import { spawn } from 'node:child_process'; import { readFile } from 'node:fs/promises'; import { sealCinemaCapabilitySnapshot, sealCinemaGenerationQuote, } from './cinema-provider-contracts.js'; import { assertCinemaProjectionQuoteParity, type CinemaProviderProjection } from './cinema-provider-projection.js'; import { writeCinemaCapabilitySnapshot, writeCinemaGenerationQuote, writeCinemaProviderProjectionBundle, } from './cinema-provider-store.js'; import type { CinemaProviderAccountClass, CinemaProviderCapabilitySnapshotArtifact, CinemaGenerationQuoteArtifact, } from './cinema-provider-types.js'; import type { CinemaQueueMoney, CinemaQueueTask } from './cinema-queue-types.js'; import type { CinemaSha256 } from './cinema-types.js'; import { buildCinemaCompatibilityExecutionPayload, resolveCinemaCompatibilityExecutionPayload, } from './cinema-compatibility-worker.js'; import { sha256Text, stableCinemaJson } from './cinema-shot-compiler.js'; import { executionAdapterCommandHash } from './execution-adapter.js'; const PAID_ACCOUNT_CLASSES = new Set(['provider-credits', 'api-credits']); export interface CinemaCompatibilityQuoteObservation { authoritative: true; routeId: string; transportId: string; accountClass: 'provider-credits' | 'api-credits'; observedAt: string; expiresAt: string; capabilityVersion: string; adapterCommandHash: string; maximumReferenceImages: number; maximumConcurrentJobs: number | null; parameterSchema: Record; models: string[]; workflows: string[]; unitCost: CinemaQueueMoney; maximumCost: CinemaQueueMoney; balance: { observed: CinemaQueueMoney; minimum: CinemaQueueMoney; maximum: CinemaQueueMoney; }; } export interface CinemaCompatibilityQuotePort { quote(input: { taskId: string; routeId: string; payload: ReturnType; providerArgumentsHash: CinemaSha256; sourceHashes: CinemaSha256[]; adapterCommandHash: CinemaSha256; }): Promise; } export interface CinemaCompatibilityQuoteResult { snapshot: CinemaProviderCapabilitySnapshotArtifact; snapshotPath: string; quote: ReturnType; quotePath: string; projection: CinemaProviderProjection; projectionBundlePath: string; providerCalls: 1; providerMutations: 0; generationCalls: 0; spendAuthorized: false; executable: false; } export interface CinemaCompatibilityQuoteBinding { payload: ReturnType; providerArguments: Record; providerArgumentsHash: CinemaSha256; sourceHashes: CinemaSha256[]; } function sha256Bytes(bytes: Uint8Array): CinemaSha256 { return `sha256:${createHash('sha256').update(bytes).digest('hex')}`; } function finiteMoney(value: CinemaQueueMoney, label: string): void { if (!value.currency.trim() || !Number.isFinite(value.amount) || value.amount < 0) { throw new Error(`${label} must contain a currency and finite non-negative amount`); } } function instant(value: string, label: string): number { const parsed = Date.parse(value); if (!Number.isFinite(parsed)) throw new Error(`${label} must be an ISO date-time`); return parsed; } function uniqueStrings(value: string[], label: string): void { if (!Array.isArray(value) || value.some((entry) => typeof entry !== 'string' || !entry.trim()) || new Set(value).size !== value.length) throw new Error(`${label} must contain unique non-empty strings`); } async function sourceHashes(task: CinemaQueueTask, payload: ReturnType): Promise { const hashes: CinemaSha256[] = []; for (const path of payload.tasks.flatMap((entry) => entry.referencePaths)) { if (/^[a-z][a-z0-9+.-]*:\/\//i.test(path)) { throw new Error(`Paid compatibility quotation requires local hashable references, not remote URL ${path}`); } try { hashes.push(sha256Bytes(await readFile(path))); } catch (error) { throw new Error(`Paid compatibility quotation cannot hash reference ${path}: ${error instanceof Error ? error.message : String(error)}`); } } const unique = [...new Set(hashes)]; if (task.payload.artifactIds.length > 0 && payload.tasks[0]?.referencePaths.length === 0) { // Artifact IDs remain bound by payloadHash/providerArgumentsHash. They are not // misrepresented as byte hashes when the compatibility source lacks a path. } return unique; } export async function buildCinemaCompatibilityQuoteBinding( root: string, projectSlug: string, task: CinemaQueueTask, ): Promise { const payload = await resolveCinemaCompatibilityExecutionPayload(root, projectSlug, task); const providerArguments = structuredClone(payload) as unknown as Record; return { payload, providerArguments, providerArgumentsHash: sha256Text(stableCinemaJson(providerArguments)), sourceHashes: await sourceHashes(task, payload), }; } function validateObservation(task: CinemaQueueTask, observation: CinemaCompatibilityQuoteObservation): void { if (observation.authoritative !== true) throw new Error('Compatibility quote must be authoritative'); if (observation.routeId !== task.routeId) throw new Error(`Compatibility quote route changed from ${task.routeId} to ${observation.routeId}`); if (!observation.transportId.trim() || !observation.capabilityVersion.trim()) throw new Error('Compatibility quote lacks transport or capability version'); if (!/^sha256:[a-f0-9]{64}$/.test(observation.adapterCommandHash)) throw new Error('Compatibility quote lacks a valid submit adapter command hash'); if (!PAID_ACCOUNT_CLASSES.has(observation.accountClass)) throw new Error('Compatibility quote must identify a paid credit account class'); const observedAt = instant(observation.observedAt, 'Compatibility quote observedAt'); const expiresAt = instant(observation.expiresAt, 'Compatibility quote expiresAt'); if (expiresAt <= observedAt) throw new Error('Compatibility quote expiry must be after observation'); if (!Number.isInteger(observation.maximumReferenceImages) || observation.maximumReferenceImages < 0) throw new Error('Compatibility quote maximumReferenceImages is invalid'); if (observation.maximumConcurrentJobs !== null && (!Number.isInteger(observation.maximumConcurrentJobs) || observation.maximumConcurrentJobs < 1)) throw new Error('Compatibility quote maximumConcurrentJobs is invalid'); if (!observation.parameterSchema || typeof observation.parameterSchema !== 'object' || Array.isArray(observation.parameterSchema)) throw new Error('Compatibility quote parameterSchema must be an object'); uniqueStrings(observation.models, 'Compatibility quote models'); uniqueStrings(observation.workflows, 'Compatibility quote workflows'); [observation.unitCost, observation.maximumCost, observation.balance.observed, observation.balance.minimum, observation.balance.maximum] .forEach((money, index) => finiteMoney(money, `Compatibility quote money[${index}]`)); const currency = observation.unitCost.currency; if ([observation.maximumCost, observation.balance.observed, observation.balance.minimum, observation.balance.maximum] .some((money) => money.currency !== currency)) throw new Error('Compatibility quote currencies must match'); if (observation.maximumCost.amount < observation.unitCost.amount) throw new Error('Compatibility maximum cost cannot be below unit cost'); if (observation.balance.minimum.amount > observation.balance.observed.amount || observation.balance.observed.amount > observation.balance.maximum.amount) throw new Error('Compatibility observed balance lies outside its quoted envelope'); if (observation.balance.minimum.amount < observation.maximumCost.amount) throw new Error('Compatibility quote balance envelope cannot cover maximum cost'); } function safeId(value: string): string { return value.replace(/[^A-Za-z0-9._-]/g, '-').replace(/^-+/, '') || 'compatibility'; } export async function quoteCinemaCompatibilityTask(options: { root: string; projectSlug: string; task: CinemaQueueTask; port: CinemaCompatibilityQuotePort; adapterEnv?: NodeJS.ProcessEnv; }): Promise { const { task } = options; if (task.authorizationRequirement !== 'exact-quote' || task.spendAuthorizationId) { throw new Error(`Compatibility task ${task.taskId} is not awaiting an exact quote`); } if (task.status !== 'awaiting-quote') throw new Error(`Compatibility task ${task.taskId} is ${task.status}, not awaiting-quote`); const binding = await buildCinemaCompatibilityQuoteBinding(options.root, options.projectSlug, task); const { payload, providerArguments, providerArgumentsHash, sourceHashes: hashes } = binding; const adapterCommandHash = executionAdapterCommandHash(payload.routeId, options.adapterEnv ?? process.env) as CinemaSha256; const observation = await options.port.quote({ taskId: task.taskId, routeId: task.routeId, payload, providerArgumentsHash, sourceHashes: hashes, adapterCommandHash, }); validateObservation(task, observation); if (observation.adapterCommandHash !== adapterCommandHash) throw new Error('Compatibility quote does not bind the resolved submit adapter command'); if (Object.hasOwn(observation.parameterSchema, 'x-videoclaw-submit-adapter-command-hash')) throw new Error('Compatibility provider parameter schema uses a reserved VideoClaw key'); if (payload.tasks[0].referencePaths.length > observation.maximumReferenceImages) { throw new Error(`Compatibility route accepts ${observation.maximumReferenceImages} references, task requires ${payload.tasks[0].referencePaths.length}`); } const timestamp = Date.parse(observation.observedAt); const suffix = `${safeId(task.routeId)}-${safeId(task.taskId)}-${timestamp}`; const snapshot = sealCinemaCapabilitySnapshot({ schemaVersion: 1, snapshotId: `capability-${suffix}`, projectSlug: options.projectSlug, routeId: task.routeId, transportId: observation.transportId, accountClass: observation.accountClass, discoveredAt: observation.observedAt, expiresAt: observation.expiresAt, discovery: { kind: 'transport-probe', version: observation.capabilityVersion, authenticated: true, sourceEvidenceIds: [providerArgumentsHash] }, capabilities: { mediaKinds: ['video'], operations: [payload.operationKind], models: [...observation.models], workflows: [...observation.workflows], maximumReferenceImages: observation.maximumReferenceImages, maximumConcurrentJobs: observation.maximumConcurrentJobs, parameterSchema: { ...structuredClone(observation.parameterSchema), 'x-videoclaw-submit-adapter-command-hash': adapterCommandHash, }, }, }); const projection: CinemaProviderProjection = { taskId: task.taskId, taskKind: task.kind, routeId: task.routeId, transportId: observation.transportId, accountClass: observation.accountClass, capabilitySnapshotId: snapshot.snapshotId, capabilitySnapshotHash: snapshot.contentHash, payloadHash: task.payloadHash, sourceHashes: hashes, providerArguments, providerArgumentsHash, fallbackPolicy: 'forbidden', }; const quote = sealCinemaGenerationQuote({ schemaVersion: 1, quoteId: `quote-${suffix}`, projectSlug: options.projectSlug, routeId: task.routeId, accountClass: observation.accountClass, capabilitySnapshotId: snapshot.snapshotId, capabilitySnapshotHash: snapshot.contentHash, createdAt: observation.observedAt, expiresAt: observation.expiresAt, jobs: [{ jobId: task.taskId, taskKind: task.kind, payloadHash: task.payloadHash, sourceHashes: hashes, providerArgumentsHash, quantity: 1, unitCost: structuredClone(observation.unitCost), maximumCost: structuredClone(observation.maximumCost), }], totalCost: structuredClone(observation.unitCost), maximumBalanceDebit: structuredClone(observation.maximumCost), balanceEnvelope: { observedBefore: structuredClone(observation.balance.observed), minimumBefore: structuredClone(observation.balance.minimum), maximumBefore: structuredClone(observation.balance.maximum), }, }); const snapshotPath = await writeCinemaCapabilitySnapshot(options.root, snapshot); const projectionBundle = await writeCinemaProviderProjectionBundle(options.root, { quote, snapshot, projections: [projection] }); const quotePath = await writeCinemaGenerationQuote(options.root, quote, snapshot); return { snapshot, snapshotPath, quote, quotePath, projection, projectionBundlePath: projectionBundle.path, providerCalls: 1, providerMutations: 0, generationCalls: 0, spendAuthorized: false, executable: false, }; } export async function revalidateCinemaCompatibilityQuote(options: { root: string; projectSlug: string; task: CinemaQueueTask; snapshot: CinemaProviderCapabilitySnapshotArtifact; quote: CinemaGenerationQuoteArtifact; projection: CinemaProviderProjection; port: CinemaCompatibilityQuotePort; adapterEnv?: NodeJS.ProcessEnv; }): Promise<{ observation: CinemaCompatibilityQuoteObservation; binding: CinemaCompatibilityQuoteBinding }> { const { task, snapshot, quote, projection } = options; if (task.authorizationRequirement !== 'exact-quote' || !task.spendAuthorizationId) { throw new Error(`Compatibility task ${task.taskId} does not carry an exact spend authorization`); } const job = quote.jobs.find((entry) => entry.jobId === task.taskId); if (!job) throw new Error(`Compatibility quote ${quote.quoteId} does not contain task ${task.taskId}`); if (quote.routeId !== task.routeId || snapshot.routeId !== task.routeId || quote.capabilitySnapshotId !== snapshot.snapshotId || quote.capabilitySnapshotHash !== snapshot.contentHash) { throw new Error('Compatibility quote, capability snapshot and queue route do not match'); } const binding = await buildCinemaCompatibilityQuoteBinding(options.root, options.projectSlug, task); const adapterCommandHash = executionAdapterCommandHash(binding.payload.routeId, options.adapterEnv ?? process.env) as CinemaSha256; if (snapshot.capabilities.parameterSchema['x-videoclaw-submit-adapter-command-hash'] !== adapterCommandHash) { throw new Error('Compatibility submit adapter command changed after authorization'); } assertCinemaProjectionQuoteParity(projection, { routeId: task.routeId, payloadHash: task.payloadHash, providerArgumentsHash: binding.providerArgumentsHash, sourceHashes: binding.sourceHashes, }); if (projection.taskId !== task.taskId || projection.transportId !== snapshot.transportId || projection.accountClass !== snapshot.accountClass) { throw new Error('Compatibility provider projection no longer matches its task or capability transport'); } const observation = await options.port.quote({ taskId: task.taskId, routeId: task.routeId, payload: binding.payload, providerArgumentsHash: binding.providerArgumentsHash, sourceHashes: binding.sourceHashes, adapterCommandHash, }); validateObservation(task, observation); if (observation.adapterCommandHash !== adapterCommandHash) throw new Error('Compatibility quote refresh does not bind the resolved submit adapter command'); if (observation.transportId !== snapshot.transportId || observation.accountClass !== snapshot.accountClass || observation.capabilityVersion !== snapshot.discovery.version) { throw new Error('Compatibility provider transport, account class or capability version changed after authorization'); } if (observation.maximumReferenceImages < binding.payload.tasks[0].referencePaths.length || !snapshot.capabilities.operations.includes(binding.payload.operationKind)) { throw new Error('Compatibility provider capability no longer supports the exact authorized payload'); } if (stableCinemaJson(observation.unitCost) !== stableCinemaJson(job.unitCost) || stableCinemaJson(observation.maximumCost) !== stableCinemaJson(job.maximumCost)) { throw new Error('Compatibility provider cost changed after authorization; obtain and approve a new quote'); } const envelope = quote.balanceEnvelope; if (observation.balance.observed.currency !== envelope.observedBefore.currency || observation.balance.observed.amount < envelope.minimumBefore.amount || observation.balance.observed.amount > envelope.maximumBefore.amount || observation.balance.observed.amount < job.maximumCost.amount) { throw new Error('Compatibility provider balance changed outside the authorized quote envelope'); } return { observation, binding }; } export function createCinemaCompatibilityQuoteCommandPort(executable: string, env: NodeJS.ProcessEnv = process.env): CinemaCompatibilityQuotePort { if (!executable.trim()) throw new Error('Compatibility quote adapter executable is required'); return { async quote(input) { const result = await new Promise<{ code: number | null; stdout: string; stderr: string }>((resolve, reject) => { const child = spawn(executable, [], { env, stdio: ['pipe', 'pipe', 'pipe'] }); let stdout = ''; let stderr = ''; child.stdout.on('data', (chunk) => { stdout += String(chunk); if (stdout.length > 1_000_000) child.kill(); }); child.stderr.on('data', (chunk) => { stderr += String(chunk); if (stderr.length > 1_000_000) child.kill(); }); child.on('error', reject); child.on('close', (code) => resolve({ code, stdout, stderr })); child.stdin.end(JSON.stringify({ action: 'quote', ...input })); }); if (result.code !== 0) throw new Error(`Compatibility quote adapter failed: ${result.stderr.trim() || `exit ${result.code}`}`); try { return JSON.parse(result.stdout) as CinemaCompatibilityQuoteObservation; } catch (error) { throw new Error(`Compatibility quote adapter returned invalid JSON: ${error instanceof Error ? error.message : String(error)}`); } }, }; }