import { execFileSync } from 'node:child_process'; import { createHash, randomUUID } from 'node:crypto'; import { hostname } from 'node:os'; import { LaneQueue, defaultLaneDbPath, type AcquireResult, type LaneId, type LaneStatus, type LaneStatusReceipt, type LaneTicket, type LaneCooldown, type LaneCooldownReason, } from './lane-queue.js'; export interface LaneCoordinatorPort { readonly coordinator: 'local-sqlite' | 'cloudflare-durable-object'; acquire(params: { lane: LaneId; project: string; scene?: number | null; promptHash?: string | null; limit: number | null; ttlSec?: number; holder?: string; }): AcquireResult; status(lane: LaneId, limit?: number | null, options?: { includeReceipts?: boolean }): LaneStatus; heartbeat(ticketId: string, ttlSec?: number): boolean; release(ticketId: string, lane?: LaneId, limit?: number | null, outcome?: LaneReleaseOutcome): boolean; /** `contentHashVerified`: the coordinator recomputed the payload hash and it matched; null from a coordinator too old to check. */ recordReceipt?(receipt: SharedLaneReceipt): { stored: boolean; deduped: boolean; contentHashVerified?: boolean | null }; /** The render log (remote coordinator only): done tickets newest first. */ history?(lane: LaneId, options?: { limit?: number; project?: string }): LaneHistory; /** Pause the lane for every caller (see LaneCooldown). Only an active lease may. */ cooldown?(ticketId: string, lane: LaneId, reasonCode: LaneCooldownReason, retryAfterSeconds: number): LaneCooldown; } /** What the driver saw when it gave the slot back — the render log's one word per * draw (2026-09-04). `expired` is the coordinator's own verdict, never sent. */ export const LANE_OUTCOMES = ['completed', 'failed', 'timeout', 'lease-lost', 'not-submitted', 'slot-busy', 'abandoned'] as const; export type LaneOutcome = (typeof LANE_OUTCOMES)[number]; export interface LaneReleaseOutcome { outcome: LaneOutcome; providerJobId?: string; note?: string; } export interface LaneHistory { lane: LaneId; tickets: LaneTicket[]; terminalRetained: number; coordinator: string; } export type SharedLaneReceiptPhase = | 'lease-acquired' | 'references-applied' | 'provider-submitted' | 'provider-terminal'; export interface SharedLaneReceipt { lane: string; ticketId: string; receiptId: string; contentHash: string; phase: SharedLaneReceiptPhase; payload: Record; } export interface RemoteLaneCoordinatorOptions { url: string; token: string; machineId: string; timeoutMs?: number; request?: RemoteLaneRequest; } export type RemoteLaneRequest = (input: { url: string; token: string; timeoutMs: number; body: Record; }) => { status: number; text: string }; export interface CreateLaneCoordinatorOptions { env?: NodeJS.ProcessEnv; dbPath?: string; } const PROTOCOL_VERSION = 1; const DEFAULT_TTL_SECONDS = 1800; const REMOTE_HELPER = String.raw` let raw = ''; process.stdin.setEncoding('utf8'); process.stdin.on('data', (chunk) => { raw += chunk; }); process.stdin.on('end', async () => { let input; try { input = JSON.parse(raw); } catch (error) { process.stderr.write('invalid coordinator helper input: ' + String(error)); process.exitCode = 2; return; } const controller = new AbortController(); const timer = setTimeout(() => controller.abort(), input.timeoutMs); try { const response = await fetch(input.url, { method: input.method, headers: { authorization: 'Bearer ' + input.token, 'content-type': 'application/json', accept: 'application/json', }, body: input.body === null ? undefined : JSON.stringify(input.body), signal: controller.signal, }); const text = await response.text(); process.stdout.write(JSON.stringify({ status: response.status, text })); } catch (error) { process.stderr.write(error instanceof Error ? error.message : String(error)); process.exitCode = 2; } finally { clearTimeout(timer); } }); `; function requiredText(value: string | undefined, name: string): string { const result = value?.trim(); if (!result) throw new Error(`${name} is required for the shared lane coordinator`); return result; } function validateRemoteUrl(value: string): string { let url: URL; try { url = new URL(value); } catch { throw new Error('VCLAW_SHARED_QUEUE_URL must be an absolute URL'); } const local = ['localhost', '127.0.0.1', '::1'].includes(url.hostname); if (url.protocol !== 'https:' && !(local && url.protocol === 'http:')) { throw new Error('VCLAW_SHARED_QUEUE_URL must use HTTPS (HTTP is allowed only for localhost)'); } url.pathname = url.pathname.replace(/\/$/, ''); url.search = ''; url.hash = ''; return url.toString().replace(/\/$/, ''); } function object(value: unknown, label: string): Record { if (!value || typeof value !== 'object' || Array.isArray(value)) { throw new Error(`shared lane coordinator returned malformed ${label}`); } return value as Record; } function defaultRemoteRequest(input: Parameters[0]): ReturnType { const raw = execFileSync(process.execPath, ['--input-type=module', '--eval', REMOTE_HELPER], { input: JSON.stringify({ ...input, method: 'POST' }), encoding: 'utf8', timeout: input.timeoutMs + 2_000, maxBuffer: 512 * 1024, windowsHide: true, }); return object(JSON.parse(raw) as unknown, 'HTTP envelope') as unknown as ReturnType; } /** * Canonical JSON for a receipt's content hash. Keys sort by UTF-16 code unit * (`<`), never `localeCompare`: the coordinator recomputes this hash in another * runtime, and a locale-aware order (`a` before `B`) is not the same string * everywhere. Kept byte-for-byte with `stableJson` in * services/shared-lane-coordinator/src/protocol.ts; one shared vector is pinned * in both test suites. */ function stableJson(value: unknown): string { if (Array.isArray(value)) return `[${value.map(stableJson).join(',')}]`; if (value && typeof value === 'object') { return `{${Object.entries(value as Record) .sort(([left], [right]) => (left < right ? -1 : left > right ? 1 : 0)) .map(([key, entry]) => `${JSON.stringify(key)}:${stableJson(entry)}`) .join(',')}}`; } return JSON.stringify(value); } /** One short attempt: see RemoteLaneCoordinator.recordReceipt. */ export const RECEIPT_TIMEOUT_MS = 3_000; export function sharedLaneReceiptContentHash(payload: Record): string { return `sha256:${createHash('sha256').update(stableJson(payload)).digest('hex')}`; } function numberOrNull(value: unknown): number | null { return value === null || value === undefined ? null : Number(value); } function coordinatorRetryDelay(milliseconds: number): void { Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0, milliseconds); } export class RemoteLaneCoordinator implements LaneCoordinatorPort { readonly coordinator = 'cloudflare-durable-object' as const; private readonly baseUrl: string; private readonly token: string; private readonly machineId: string; private readonly timeoutMs: number; private readonly request: RemoteLaneRequest; // One claim per DRIVER, not per process. lane_render spawns a fresh `vclaw video lane` // process for every call, so a per-process claim could never match its own ticket on // heartbeat; the driver exports VCLAW_LANE_CLAIM once and every call carries it. // Acquire always names a claim (the ledger stores one per ticket); heartbeat and // release name it ONLY when the driver set VCLAW_LANE_CLAIM. A hand // `vclaw video lane release --ticket X` from the same machine therefore still // unwedges a SIGKILLed driver's ticket under the machine rule — a throwaway UUID // there would leave the slot dead until the TTL. private readonly claimId = process.env.VCLAW_LANE_CLAIM?.trim() || randomUUID(); private readonly namedClaim: string | undefined = process.env.VCLAW_LANE_CLAIM?.trim() || undefined; constructor(options: RemoteLaneCoordinatorOptions) { this.baseUrl = validateRemoteUrl(requiredText(options.url, 'VCLAW_SHARED_QUEUE_URL')); this.token = requiredText(options.token, 'VCLAW_SHARED_QUEUE_TOKEN'); this.machineId = requiredText(options.machineId, 'VCLAW_MACHINE_ID'); this.timeoutMs = options.timeoutMs ?? 10_000; this.request = options.request ?? defaultRemoteRequest; if (!Number.isInteger(this.timeoutMs) || this.timeoutMs < 500 || this.timeoutMs > 60_000) { throw new Error('shared lane coordinator timeout must be 500-60000ms'); } } private call(path: string, payload: Record, limits: { attempts: number; timeoutMs: number } = { attempts: 3, timeoutMs: this.timeoutMs }): Record { let envelope: Record | null = null; const request = { url: `${this.baseUrl}${path}`, token: this.token, timeoutMs: limits.timeoutMs, body: { protocolVersion: PROTOCOL_VERSION, machineId: this.machineId, ...payload }, }; let lastError: unknown; for (let attempt = 0; attempt < limits.attempts; attempt += 1) { try { envelope = object(this.request(request), 'HTTP envelope'); const responseStatus = Number(envelope.status); if (![401, 429].includes(responseStatus) && responseStatus < 500) break; lastError = new Error(`HTTP ${responseStatus}`); } catch (error) { lastError = error; } // Queue operations are idempotent by stable request/ticket/receipt ID. // A bounded retry covers Worker cold start, secret propagation and edge // overload; it never falls back to SQLite or retries provider submission. if (attempt < limits.attempts - 1) coordinatorRetryDelay(attempt === 0 ? 250 : 750); } if (!envelope) { const detail = lastError instanceof Error ? lastError.message : String(lastError); throw new Error(`shared lane coordinator is unavailable; provider submission is blocked (${detail})`); } const status = Number(envelope.status); const responseText = String(envelope.text ?? ''); let response: Record; try { response = object(JSON.parse(responseText) as unknown, 'JSON response'); } catch { throw new Error(`shared lane coordinator returned HTTP ${status} with a malformed JSON response`); } if (!Number.isInteger(status) || status < 200 || status >= 300) { throw new Error(`shared lane coordinator rejected the request with HTTP ${status}: ${String(response.error ?? 'unknown error')}`); } if (response.coordinator !== this.coordinator) { throw new Error('shared lane coordinator response is missing its Cloudflare authority marker'); } return response; } acquire(params: Parameters[0]): AcquireResult { if (params.limit === null) { return { status: 'granted', ticketId: `unlimited-${Date.now()}`, position: 0, etaSeconds: 0, lane: params.lane, deduped: false, coordinator: this.coordinator, }; } const promptHash = params.promptHash?.trim(); if (!promptHash) throw new Error(`lane ${params.lane} requires a stable promptHash/request identity`); const response = this.call('/v1/lanes/acquire', { lane: params.lane, project: params.project, scene: params.scene ?? null, promptHash, limit: params.limit, ttlSeconds: params.ttlSec ?? DEFAULT_TTL_SECONDS, claimId: this.claimId, }); const status = response.status; if (status !== 'granted' && status !== 'queued') throw new Error('shared lane coordinator returned an invalid acquire status'); return { status, ticketId: requiredText(typeof response.ticketId === 'string' ? response.ticketId : undefined, 'ticketId'), position: Number(response.position), etaSeconds: numberOrNull(response.etaSeconds), lane: requiredText(typeof response.lane === 'string' ? response.lane : undefined, 'lane'), deduped: response.deduped === true, coordinator: this.coordinator, // A Worker deployed before cooldowns existed sends no field: read as none. cooldown: cooldownOrNull(response.cooldown), }; } status(lane: LaneId, limit: number | null = null, options: { includeReceipts?: boolean } = {}): LaneStatus { if (limit === null) { return { lane, limit, held: [], queued: [], medianJobSeconds: null, terminalRetained: 0, coordinator: this.coordinator }; } const response = this.call('/v1/lanes/status', { lane, limit }); if (!Array.isArray(response.held) || !Array.isArray(response.queued)) { throw new Error('shared lane coordinator returned malformed lane status'); } return { lane: requiredText(typeof response.lane === 'string' ? response.lane : undefined, 'lane'), limit: Number(response.limit), held: response.held as LaneStatus['held'], queued: response.queued as LaneStatus['queued'], medianJobSeconds: numberOrNull(response.medianJobSeconds), terminalRetained: Number(response.terminalRetained), coordinator: this.coordinator, cooldown: cooldownOrNull(response.cooldown), // The coordinator sends these either way and `call` has already parsed them; // carrying them is opt-in only so the answer a driver reads in a loop stays // short to print. `unreadableReceipts` rides along so a shortened trail is // never mistaken for a step that left no evidence. ...(options.includeReceipts ? receiptTrail(response.recentReceipts) : {}), }; } cooldown(ticketId: string, lane: LaneId, reasonCode: LaneCooldownReason, retryAfterSeconds: number): LaneCooldown { if (ticketId.startsWith('unlimited-')) throw new Error('an uncoordinated lane cannot publish a shared cooldown'); // The owner rule is the release rule: the machine, plus the claim when the // driver named one (two drivers on one machine, #425). const response = this.call('/v1/lanes/cooldown', { lane, ticketId, reasonCode, retryAfterSeconds, ...(this.namedClaim ? { claimId: this.namedClaim } : {}), }); const cooldown = cooldownOrNull(response.cooldown); if (!cooldown) throw new Error('shared lane coordinator returned no cooldown for a cooldown request'); return cooldown; } heartbeat(ticketId: string, ttlSec = DEFAULT_TTL_SECONDS): boolean { if (ticketId.startsWith('unlimited-')) return true; const lane = ticketId.slice(0, ticketId.lastIndexOf(':')); if (!lane) throw new Error('remote lane heartbeat requires a Cloudflare ticket containing its lane'); return this.call('/v1/lanes/heartbeat', { lane, ticketId, ttlSeconds: ttlSec, ...(this.namedClaim ? { claimId: this.namedClaim } : {}) }).refreshed === true; } release(ticketId: string, lane?: LaneId, limit: number | null = null, outcome?: LaneReleaseOutcome): boolean { if (ticketId.startsWith('unlimited-')) return true; const exactLane = lane ?? ticketId.slice(0, ticketId.lastIndexOf(':')); if (!exactLane) throw new Error('remote lane release requires the exact lane'); return this.call('/v1/lanes/release', { lane: exactLane, ticketId, ...(this.namedClaim ? { claimId: this.namedClaim } : {}), ...(limit === null ? {} : { limit }), ...(outcome ? { outcome: outcome.outcome } : {}), ...(outcome?.providerJobId ? { providerJobId: outcome.providerJobId.slice(0, 256) } : {}), ...(outcome?.note ? { note: outcome.note.slice(0, 512) } : {}), }).released === true; } history(lane: LaneId, options: { limit?: number; project?: string } = {}): LaneHistory { const response = this.call('/v1/lanes/history', { lane, ...(options.limit === undefined ? {} : { limit: options.limit }), ...(options.project === undefined ? {} : { project: options.project }), }); return { lane, tickets: (response.tickets ?? []) as LaneTicket[], terminalRetained: Number(response.terminalRetained ?? 0), coordinator: this.coordinator, }; } recordReceipt(receipt: SharedLaneReceipt): { stored: boolean; deduped: boolean; contentHashVerified: boolean | null } { // A receipt is evidence about a render, written from inside the render path // (one of them between the provider accepting a job and its id reaching // disk). The default three attempts at ten seconds could hold that path for // about 37 s per receipt; a receipt gets one short attempt and is dropped. const response = this.call('/v1/lanes/receipt', { ...receipt }, { attempts: 1, timeoutMs: Math.min(this.timeoutMs, RECEIPT_TIMEOUT_MS) }); return { stored: response.stored === true, deduped: response.deduped === true, contentHashVerified: typeof response.contentHashVerified === 'boolean' ? response.contentHashVerified : null, }; } verify(lane: string, limit: number): LaneStatus { return this.status(lane, limit); } } const RECEIPT_PHASES: ReadonlySet = new Set(['lease-acquired', 'references-applied', 'provider-submitted', 'provider-terminal']); /** * Map the coordinator's receipt rows, dropping anything malformed rather than * failing a status read — but COUNTING what was dropped. On a surface whose whole * job is evidence, a silently shorter list reads as "that step left no receipt" * when the truth is "this client could not read the row the coordinator sent". * The count is what lets a reader tell those apart. */ function receiptsOf(value: unknown): { receipts: LaneStatusReceipt[]; unreadable: number } { if (!Array.isArray(value)) return { receipts: [], unreadable: 0 }; const receipts: LaneStatusReceipt[] = []; let unreadable = 0; for (const entry of value) { if (!entry || typeof entry !== 'object' || Array.isArray(entry)) { unreadable += 1; continue; } const row = entry as Record; const phase = String(row.phase ?? ''); const payload = row.payload; // One rule: a row this client cannot represent faithfully is COUNTED, never // half-reported. Coercing an absent field would be a quiet lie — `''` reads as // "this receipt's hash is empty" rather than "the row was incomplete", and the // hash is what the whole trail turns on. `createdAt` must be a finite number // for the same reason and one more: `Number(undefined ?? 0)` is `0`, a // plausible-looking 1970 timestamp that would sail through unflagged, while // `Number('x')` is `NaN`, which `JSON.stringify` emits as `null` — violating // this field's own declared type. if (!RECEIPT_PHASES.has(phase) || typeof row.receiptId !== 'string' || typeof row.ticketId !== 'string' || typeof row.contentHash !== 'string' || typeof row.machineId !== 'string' || typeof row.createdAt !== 'number' || !Number.isFinite(row.createdAt)) { unreadable += 1; continue; } receipts.push({ receiptId: row.receiptId, ticketId: row.ticketId, phase: phase as LaneStatusReceipt['phase'], contentHash: row.contentHash, contentHashVerified: typeof row.contentHashVerified === 'boolean' ? row.contentHashVerified : null, machineId: row.machineId, createdAt: row.createdAt, payload: payload && typeof payload === 'object' && !Array.isArray(payload) ? payload as Record : {}, }); } return { receipts, unreadable }; } /** The two fields a verbose status answer carries: the trail, and what was lost from it. */ function receiptTrail(value: unknown): Pick { const { receipts, unreadable } = receiptsOf(value); return { recentReceipts: receipts, unreadableReceipts: unreadable }; } function cooldownOrNull(value: unknown): LaneCooldown | null { if (!value || typeof value !== 'object' || Array.isArray(value)) return null; const c = value as Record; if (c.reasonCode !== 'provider-human-check' || typeof c.until !== 'number' || typeof c.ticketId !== 'string') return null; return { reasonCode: 'provider-human-check', until: c.until, remainingSeconds: typeof c.remainingSeconds === 'number' ? c.remainingSeconds : Math.max(1, Math.ceil(c.until - Date.now() / 1000)), machineId: typeof c.machineId === 'string' ? c.machineId : null, ticketId: c.ticketId, reportedAt: typeof c.reportedAt === 'number' ? c.reportedAt : c.until, }; } export function sharedLaneCoordinatorRequired(env: NodeJS.ProcessEnv = process.env): boolean { return env.VCLAW_SHARED_QUEUE_REQUIRED === '1'; } export function createLaneCoordinator(options: CreateLaneCoordinatorOptions = {}): LaneCoordinatorPort { const env = options.env ?? process.env; const url = env.VCLAW_SHARED_QUEUE_URL?.trim(); const required = sharedLaneCoordinatorRequired(env); if (url || required) { if (!url) throw new Error('VCLAW_SHARED_QUEUE_REQUIRED=1 but VCLAW_SHARED_QUEUE_URL is missing'); const configuredMachine = env.VCLAW_MACHINE_ID?.trim(); if (required && !configuredMachine) { throw new Error('VCLAW_SHARED_QUEUE_REQUIRED=1 but VCLAW_MACHINE_ID is missing; assign a distinct stable ID on each computer'); } return new RemoteLaneCoordinator({ url, token: requiredText(env.VCLAW_SHARED_QUEUE_TOKEN, 'VCLAW_SHARED_QUEUE_TOKEN'), machineId: configuredMachine || hostname(), timeoutMs: Number(env.VCLAW_SHARED_QUEUE_TIMEOUT_MS ?? 10_000), }); } return new LaneQueue({ dbPath: options.dbPath ?? defaultLaneDbPath(env) }); }