import { createHash } from 'node:crypto'; import { existsSync } from 'node:fs'; import { readFile, rename, unlink } from 'node:fs/promises'; import sharp from 'sharp'; import { generateGenImage } from './gen-image.js'; import type { CinemaImageSubmitPayload } from './cinema-image-payload.js'; import type { CinemaCompatibilityWorkerObservation } from './cinema-compatibility-worker.js'; import type { CinemaQueueMoney } from './cinema-queue-types.js'; /** * What each Cinema image route can actually DO today — declared once, so the * compile gate, the worker's transport choice and the docs cannot disagree. * * `transmitsReferences` is the load-bearing one. OpenAI's * `/v1/images/generations` is a pure text-to-image endpoint that accepts no * image input, so a plan carrying verified `sources` cannot be honoured there. * Sending it anyway would render from prose alone while the dry run showed * references — the contract diverging from the submission, which is the single * thing this lane exists to prevent. */ export const CINEMA_IMAGE_TRANSPORTS: Record = { 'openai-images': { transmitsReferences: false }, }; /** Null when the route has no transport at all, so nothing can render it yet. */ export function cinemaImageTransportFor(routeId: string): { transmitsReferences: boolean } | null { return CINEMA_IMAGE_TRANSPORTS[routeId] ?? null; } /** * An image transport. `lookup` reuses the media-neutral * `CinemaCompatibilityWorkerObservation` — its `outputs[].kind` already includes * `'image'` and it already carries `authoritative`, so there is nothing to widen. */ export interface CinemaImageWorkerTransport { /** * Throw when the render cannot be ATTEMPTED at all — no credential, no * reachable endpoint. Called while the task is still `leased`, before the queue * is told a submission is in flight, because "we never tried" and "we tried and * do not know" need different endings: the first can simply fail and be * retried, the second must park for reconciliation in case it charged. */ preflight?(): Promise; submit(payload: CinemaImageSubmitPayload): Promise<{ providerJobId: string; providerStatus: string; envelope: Record }>; lookup(input: { payload: CinemaImageSubmitPayload; providerJobId: string | null; }): Promise; } function bytesHash(bytes: Buffer): string { return `sha256:${createHash('sha256').update(bytes).digest('hex')}`; } /** * Report on bytes already on disk at the payload's deterministic output path. * * This is what makes a crash survivable. Image generation is synchronous: there * is no provider job to poll, so "did my call land?" can only be answered by * looking for the bytes. Because `outputPath` derives from the task id alone, a * worker that died after writing finds them on the next run and reports the job * complete — instead of asking a human, or paying a second time. */ /** * Bytes at the path are not evidence of a finished render — a crash or a full * disk mid-write leaves a truncated file. Claiming that as `completed` records * the full spend against a task that then fails to ingest, leaving it terminal, * paid, and with no artifact and no way out. Decode before believing it. * * This is a full sharp decode, deliberately STRICTER than the ingest's * magic-bytes sniff: the ingest decodes too, but only after the spend has already * been recorded, so this is the last point where refusing is still free. */ async function isCompleteImage(path: string): Promise { try { await sharp(await readFile(path), { failOn: 'error' }).raw().toBuffer(); return true; } catch { return false; } } async function observeWrittenBytes( payload: CinemaImageSubmitPayload, actualCost: CinemaQueueMoney, authoritativeWhenAbsent: boolean, ): Promise { // A crash between the provider write and the rename leaves COMPLETE, PAID bytes // at the temp path that nothing would otherwise look for. Recover them: finding // a render that was already bought is the whole reason this lookup reads the // disk instead of asking a human. const partial = `${payload.outputPath}.provider`; if (!existsSync(payload.outputPath) && existsSync(partial) && await isCompleteImage(partial)) { await rename(partial, payload.outputPath); } if (!existsSync(payload.outputPath) || !(await isCompleteImage(payload.outputPath))) { return { // Absence is genuinely ambiguous: the provider may have charged and the // write may have failed. Saying so is what stops an already-paid job being // resubmitted — `reconcileCinemaProviderTask` trusts this flag, and a false // `true` here is exactly how money gets spent twice. state: 'not-found', authoritative: authoritativeWhenAbsent, providerJobId: null, providerStatus: 'no-output', envelope: { outputPath: payload.outputPath }, outputs: [], actualCost: null, }; } const bytes = await readFile(payload.outputPath); const contentHash = bytesHash(bytes); return { state: 'completed', authoritative: true, providerJobId: contentHash, providerStatus: 'completed', envelope: { outputPath: payload.outputPath, contentHash, byteLength: bytes.byteLength }, outputs: [{ id: payload.taskId, kind: 'image', path: payload.outputPath, backend: payload.routeId }], actualCost, }; } export interface CinemaImageOpenAiTransportOptions { /** * What the exact quote said this render costs, recorded as the actual spend. * * It is the QUOTE, not a provider-reported charge: the OpenAI images endpoint * reports none, so there is nothing to measure. That is an honest constraint * rather than a measurement — a price change between quote and render records * the quoted figure. The authorization check is still real (the bytes are * re-read and re-hashed, and this value is what it verifies against), but it * verifies the quote against itself. */ unitCost: CinemaQueueMoney; /** Falls back to OPENAI_API_KEY. Explicit so a test never depends on ambient env. */ apiKey?: string; env?: NodeJS.ProcessEnv; /** Injected in tests so the suite never reaches the network or needs a real key. */ fetcher?: typeof fetch; model?: string; size?: string; } /** * `openai-images`. The OpenAI Images API answers in one call, so `submit` IS the * render: it writes the bytes and returns their hash as the provider job id. * * Its lookup can never be authoritative about an absent file — the API exposes no * job history to ask — so a missing output stays unknown and the task parks for * reconciliation rather than being retried. */ export function createCinemaImageOpenAiTransport(options: CinemaImageOpenAiTransportOptions): CinemaImageWorkerTransport { const authoritativeWhenAbsent = false; return { async preflight() { // The likeliest first failure by far, and the one that costs nothing: no // key means no call, so this must never become an ambiguous submission. const key = options.apiKey ?? (options.env ?? process.env).OPENAI_API_KEY; if (!key?.trim()) { throw new Error('openai-images needs an OpenAI API key: set OPENAI_API_KEY or pass one. Nothing was submitted and nothing was billed.'); } }, async submit(payload) { // Count the HTTP calls the render actually makes, by wrapping the fetcher. // A throw before the first one — a missing key, an unreadable reference — // cost nothing, and saying so is what lets a parked task be retried instead // of sitting in reconciliation forever. let providerCalls = 0; const base = options.fetcher ?? fetch; const counting: typeof fetch = async (...args: Parameters) => { providerCalls += 1; return base(...args); }; // Write beside the final path, prove the bytes decode, THEN rename. The // queue already names a `.provider` temp path in its download receipt; this // makes that real instead of fiction, so a crash mid-write can never be // mistaken for a finished render. const partial = `${payload.outputPath}.provider`; const result = await generateGenImage({ prompt: payload.prompt, // The role-driven prompt from the plan is authoritative. An identity // reference is not a diegetic prop, so the per-kind directive stays off; // `kind` still selects the canvas size. kind: 'prop', directive: 'none', backend: 'openai', outputPath: partial, ...(options.apiKey !== undefined ? { apiKey: options.apiKey } : {}), ...(options.model !== undefined ? { model: options.model } : {}), ...(options.size !== undefined ? { size: options.size } : {}), fetcher: counting, }).catch((error: unknown) => { // Pin it to the error: the worker records this so a later lookup can tell // "nothing was billed" from "it may have charged". if (error instanceof Error) Object.assign(error, { vclawReachedProvider: providerCalls > 0 }); throw error; }); if (!(await isCompleteImage(result.path))) { await unlink(result.path).catch(() => undefined); throw new Error('Cinema image render produced bytes that do not decode as an image; nothing was kept'); } await rename(result.path, payload.outputPath); const bytes = await readFile(payload.outputPath); return { providerJobId: bytesHash(bytes), providerStatus: 'completed', envelope: { path: payload.outputPath, model: result.model, sizeBytes: result.sizeBytes, ...(result.size ? { size: result.size } : {}) }, }; }, async lookup({ payload }) { // Absent bytes stay ambiguous here BY DESIGN: anything that reached this // point attempted the render, so the provider may have charged, and // `reconcileCinemaProviderTask` trusts `authoritative`. The case that cost // nothing — no credential — is caught by `preflight` before the queue is // ever told a submission is in flight. return observeWrittenBytes(payload, options.unitCost, authoritativeWhenAbsent); }, }; }