// Remote-worker CLIENT side: delegate one transcription to another instance of // this app over HTTP. Uploads the audio, polls progress (forwarding the remote's // log + parsed {fraction, detail}), pulls back the transcript.json bytes, and // cleans up the remote scratch. transcribeOne calls this for kind:"remote". // // Failure classes (see TranscribeError): a network/5xx/poll problem is // "transport" (the batch retries on another worker); a remote job that finishes // "failed" is "transcription" (the audio itself failed — no point retrying). import fs from "fs-extra"; import type { Worker } from "../lib/workers"; import { TranscribeError } from "./transcribeError"; const POLL_MS = 1000; // Stop polling if the remote goes silent this long (no successful poll), to // avoid a wedged remote parking a batch slot forever. const POLL_TIMEOUT_MS = 10 * 60 * 1000; export type RemoteTranscribeOptions = { worker: Worker; // kind === "remote" audioPath: string; audioName: string; videoId: string; // Required for the shared-fs fast path (worker.remote.sharedFs): the remote // transcribes channelsDir//data/ in place on the shared mount. channelSlug?: string; onLog?: (line: string) => void; onProgress?: (patch: { fraction?: number; detail?: string }) => void; signal?: AbortSignal; }; type EventsResponse = { status: "queued" | "running" | "done" | "failed" | "cancelled" | "unknown"; fraction?: number; detail?: string; 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(/\/+$/, ""); } // Probe a remote worker's /api/worker/health. Returns true only on a 200 with a // healthy body. Used to confirm a remote is actually down (vs a transient job // error) so the scheduler can degrade it immediately instead of burning retries. export async function pingRemoteHealth( worker: Worker, timeoutMs = 3000, ): Promise { const baseUrl = base(worker); if (!baseUrl) return false; try { const res = await fetch(`${baseUrl}/api/worker/health`, { headers: authHeaders(worker), signal: AbortSignal.timeout(timeoutMs), }); if (!res.ok) return false; const body = (await res.json()) as { ok?: boolean }; return body.ok === true; } catch { return false; } } // Delegate a transcription to the remote and poll to completion. Returns the // transcript.json bytes for the upload path, or null for the shared-fs path // (the remote wrote the transcript directly onto the shared mount, so the caller // already has it in place and only needs to normalize). export async function transcribeViaRemote( opts: RemoteTranscribeOptions, ): 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 sharedFs = worker.remote?.sharedFs === true && !!opts.channelSlug; // 1. Submit the job → remoteJobId. Shared-fs sends a JSON path reference; // otherwise the raw audio bytes are uploaded. let remoteJobId: string; try { const res = sharedFs ? await fetch(`${baseUrl}/api/worker/transcribe`, { method: "POST", headers: { ...authHeaders(worker), "Content-Type": "application/json" }, body: JSON.stringify({ channelSlug: opts.channelSlug, videoId: opts.videoId, audioFilename: opts.audioName, }), signal: opts.signal, }) : await fetch(`${baseUrl}/api/worker/transcribe`, { method: "POST", headers: { ...authHeaders(worker), "Content-Type": "application/octet-stream", "X-Audio-Name": opts.audioName, "X-Job-Label": opts.videoId, }, body: await fs.readFile(opts.audioPath), signal: opts.signal, }); if (!res.ok) { throw new TranscribeError( `remote worker rejected the job (${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 worker 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 worker "${worker.name}": ${(err as Error).message}`, "transport", ); } const jobUrl = `${baseUrl}/api/worker/transcribe/${remoteJobId}`; // 2. Poll events until terminal, forwarding log + progress. try { let lastOk = Date.now(); for (;;) { if (opts.signal?.aborted) throw abortError(); await delay(POLL_MS, 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 > POLL_TIMEOUT_MS) { throw new TranscribeError( `remote worker "${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.fraction !== undefined || ev.detail !== undefined) { opts.onProgress?.({ fraction: ev.fraction, detail: ev.detail }); } if (ev.status === "done") break; if (ev.status === "failed") { throw new TranscribeError( `remote transcription failed on "${worker.name}"`, "transcription", ); } if (ev.status === "cancelled") { throw new TranscribeError( `remote transcription was cancelled on "${worker.name}"`, "transport", ); } if (ev.status === "unknown") { throw new TranscribeError( `remote worker "${worker.name}" forgot the job`, "transport", ); } } // 3. Shared-fs: the transcript is already on the shared mount, nothing to // pull. Otherwise fetch the produced bytes. if (sharedFs) return null; const res = await fetch(`${jobUrl}/result`, { headers: authHeaders(worker), signal: opts.signal, }); if (!res.ok) { throw new TranscribeError( `remote worker "${worker.name}" had no result (${res.status})`, "transport", ); } return Buffer.from(await res.arrayBuffer()); } 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 transcription 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 }, ); }); }