// Remote-unit CLIENT side: delegate ONE backfill unit (attribution, // diarization, …) to a unit executor over HTTP. The clone of // remoteTranscribe.ts for the general seam: submit the kind's input files and // the primary's identity config, poll progress, pull back the kind's declared // output files, clean up the remote scratch. the backfill lane calls this when a // remote worker slot matching the kind is free. // // Failure classes reuse TranscribeError verbatim: a network/5xx/poll problem // is "transport" (the batch retries on another worker or falls back to a local // run); a remote job that finishes "failed" is "transcription" (the WORK // failed — an executor with the same code would fail the same way, so no // retry). The names say "transcribe" for history; the split is the same. import type { Worker } from "../lib/workers"; import type { StartWorkerUnitInput } from "./workerServer"; import { TranscribeError } from "./transcribeError"; const POLL_MS = 1000; // Stop polling if the remote goes silent this long (no successful poll), to // avoid a wedged executor parking a batch slot forever. const POLL_TIMEOUT_MS = 10 * 60 * 1000; export type RemoteUnitOptions = { worker: Worker; // kind === "remote" op: string; channelSlug: string; videoId: string; // Input files by name, raw bytes (the caller read them; base64 happens here). files: Record; target: unknown; config?: StartWorkerUnitInput["config"]; force?: boolean; onLog?: (line: string) => void; signal?: AbortSignal; // Injectable for tests; production callers leave the defaults. pollMs?: number; pollTimeoutMs?: number; }; export type RemoteUnitOutcome = { // The executor-side OperationRunOutcome ("done", "already-present", …). outcome: string; // The kind's declared output files as UTF-8 JSON text, to be applied on the // primary through kind.applyResult — never written raw. files: Record; }; type EventsResponse = { status: "queued" | "running" | "done" | "failed" | "cancelled" | "unknown"; log?: string; }; function authHeaders(worker: Worker): Record { const token = worker.remote?.token?.trim(); return token ? { Authorization: `Bearer ${token}` } : {}; } function base(worker: Worker): string { return (worker.remote?.baseUrl ?? "").replace(/\/+$/, ""); } export async function runUnitViaRemote( opts: RemoteUnitOptions, ): Promise { const { worker } = opts; const baseUrl = base(worker); if (!baseUrl) { throw new TranscribeError( `remote worker "${worker.name}" has no base URL`, "transport", ); } const log = opts.onLog ?? (() => {}); const pollMs = opts.pollMs ?? POLL_MS; const pollTimeoutMs = opts.pollTimeoutMs ?? POLL_TIMEOUT_MS; // 1. Submit the envelope → remoteJobId. let remoteJobId: string; try { const files: Record = {}; for (const [name, bytes] of Object.entries(opts.files)) { files[name] = bytes.toString("base64"); } const res = await fetch(`${baseUrl}/api/worker/unit`, { method: "POST", headers: { ...authHeaders(worker), "Content-Type": "application/json" }, body: JSON.stringify({ op: opts.op, channelSlug: opts.channelSlug, videoId: opts.videoId, files, target: opts.target, force: opts.force, config: opts.config, } satisfies StartWorkerUnitInput & { force?: boolean }), signal: opts.signal, }); if (!res.ok) { // A 4xx (bad envelope, kind refused, an executor too old to have the // endpoint) is a WORK-class refusal — another executor running the same // code would refuse the same way. 5xx/429 is the box's problem. throw new TranscribeError( `remote executor rejected the unit (${res.status} ${res.statusText})`, res.status >= 400 && res.status < 500 && res.status !== 429 ? "transcription" : "transport", ); } const body = (await res.json()) as { remoteJobId?: string }; if (!body.remoteJobId) { throw new TranscribeError("remote executor returned no job id", "transport"); } remoteJobId = body.remoteJobId; } catch (err) { if (err instanceof TranscribeError) throw err; if (opts.signal?.aborted) throw err; throw new TranscribeError( `could not reach remote executor "${worker.name}": ${(err as Error).message}`, "transport", ); } const jobUrl = `${baseUrl}/api/worker/unit/${remoteJobId}`; // 2. Poll events until terminal, forwarding the executor's log. try { let lastOk = Date.now(); for (;;) { if (opts.signal?.aborted) throw abortError(); await delay(pollMs, opts.signal); let ev: EventsResponse; try { const res = await fetch(`${jobUrl}/events`, { headers: authHeaders(worker), signal: opts.signal, }); if (!res.ok) throw new Error(`${res.status} ${res.statusText}`); ev = (await res.json()) as EventsResponse; lastOk = Date.now(); } catch (err) { if (opts.signal?.aborted) throw abortError(); if (Date.now() - lastOk > pollTimeoutMs) { throw new TranscribeError( `remote executor "${worker.name}" stopped responding`, "transport", ); } continue; // transient — keep polling until the timeout } if (ev.log) for (const line of ev.log.split(/\r?\n/)) if (line) log(line); if (ev.status === "done") break; if (ev.status === "failed") { throw new TranscribeError( `remote unit failed on "${worker.name}"`, "transcription", ); } if (ev.status === "cancelled") { throw new TranscribeError( `remote unit was cancelled on "${worker.name}"`, "transport", ); } if (ev.status === "unknown") { throw new TranscribeError( `remote executor "${worker.name}" forgot the unit`, "transport", ); } } // 3. Pull the outcome + output files. const res = await fetch(`${jobUrl}/result`, { headers: authHeaders(worker), signal: opts.signal, }); if (!res.ok) { throw new TranscribeError( `remote executor "${worker.name}" had no unit result (${res.status})`, "transport", ); } const body = (await res.json()) as { outcome?: string; files?: Record; }; if (typeof body.outcome !== "string") { throw new TranscribeError( `remote executor "${worker.name}" returned a malformed unit result`, "transport", ); } return { outcome: body.outcome, files: body.files ?? {} }; } finally { // 4. Best-effort cleanup of the remote scratch (also cancels if still // running — e.g. when we abort). Fire-and-forget; never block on it. void fetch(jobUrl, { method: "DELETE", headers: authHeaders(worker) }).catch( () => {}, ); } } function abortError(): Error { const err = new Error("remote unit aborted"); err.name = "AbortError"; return err; } function delay(ms: number, signal?: AbortSignal): Promise { return new Promise((resolve, reject) => { if (signal?.aborted) return reject(abortError()); const t = setTimeout(resolve, ms); signal?.addEventListener( "abort", () => { clearTimeout(t); reject(abortError()); }, { once: true }, ); }); }