commit a79b3c4e8cf8be94549edc127af6572d4b335573
parent 5e4d31824f4dbabbc19322d419bbd69d5c1f16b3
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Mon, 24 Aug 2026 10:16:50 -0400
backfill units: the lane ships work to tagged executors and applies it guarded
common/controller/remoteUnit.ts is the transcribe client's general sibling —
submit envelope / poll events / pull outcome+outputs / fire-and-forget DELETE,
with the TranscribeError transport-vs-work split verbatim (5xx/unreachable/
silence → transport; a 4xx refusal or a failed job → work-class, final).
backfillBatch's runOne becomes a dispatch ladder: runBackfillUnit (tryAcquire
[kind.id, contendsFor] on TAGGED remote workers only — tagging is the opt-in
that keeps an untagged pre-protocol remote from being shipped envelopes it
cannot serve), then the llm fan-out, then the local run unchanged. Unit
results land through kind.applyResult (downgrade rule re-checked on current
disk; "refused" maps to already-present); transport failures mark/degrade the
worker via pingRemoteHealth and retry elsewhere or locally, incrementing
neither failed nor attempted. The limit gains a unit term of
unitActive + min-across-kinds free tagged slots — MIN so a mixed run cannot
dispatch one kind locally on a slot justified by another kind's remote
capacity. Envelope config resolved once per run: attribution settings,
appConfig with baseUrl stripped, injected channel context, diarization
settings. e2e: worker-unit.spec.ts proves auth, the non-kind door guard, and
a full attribution round trip against a scratch corpus + in-test ollama stub
(6/6 green with worker-remote.spec.ts).
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Diffstat:
8 files changed, 863 insertions(+), 49 deletions(-)
diff --git a/RUNNING_IN_DOCKER.md b/RUNNING_IN_DOCKER.md
@@ -268,8 +268,10 @@ corpus at all. Two shapes, by cost:
sidecars back and applies them through its own guarded writers. The
executor's own settings are never consulted for the work's identity — the
primary injects its model/prompt configuration into every unit, and a unit
- that would fall back to local defaults fails loudly instead. Tag the worker
- (e.g. `cpu, diarization`) to say what it should take.
+ that would fall back to local defaults fails loudly instead. **Tag the
+ worker** (e.g. `cpu, diarization`) — units are shipped only to *tagged*
+ remote workers, so an untagged remote from before this protocol keeps doing
+ transcription only.
The executor never opens the primary's LMDB, never downloads media, and never
arms a sweep — the primary stays the sole scheduler.
diff --git a/common/controller/backfillBatch.ts b/common/controller/backfillBatch.ts
@@ -31,7 +31,7 @@
// it a share. Nothing new in the scheduler.
import path from "node:path";
-import { readdir } from "node:fs/promises";
+import { readFile, readdir } from "node:fs/promises";
import type { Paths } from "../lib/paths";
import { getSettings } from "../lib/settings";
import { runPool } from "../jobs/concurrentRunner";
@@ -52,6 +52,12 @@ import { transcriptionActivity } from "./digestYield";
import { acquireLlmSlot, freeLlmSlots } from "./llmWorkers";
import { getWorkerPool } from "../jobs/workerPool";
import { resolveAttributionTarget } from "./attributionTarget";
+import { pingRemoteHealth } from "./remoteTranscribe";
+import { runUnitViaRemote } from "./remoteUnit";
+import { TranscribeError } from "./transcribeError";
+import { readDigestContext } from "../lib/digestContext-server";
+import type { DigestAppConfig } from "../lib/digest";
+import type { StartWorkerUnitInput } from "./workerServer";
import { reacquireMediaFor, type ReacquireOutcome } from "./backfillReacquire";
import type { AutoQueueOrder } from "../jobs/autoQueuePolicy";
import { batchRecencyComparator } from "./batchRecency";
@@ -230,6 +236,14 @@ type Candidate = { id: string; kind: BackfillKind; target: unknown };
// pass, which only a unit executor (a full instance of this app) can take.
const LLM_BACKFILL_OPS = new Set(["attribution-text", "attribution-diarized"]);
+// The executor localises its own endpoint (its OLLAMA_URL); everything else in
+// the config IS identity and travels verbatim.
+function stripBaseUrl(config: DigestAppConfig): DigestAppConfig {
+ const { baseUrl: _baseUrl, ...rest } = config;
+ void _baseUrl;
+ return rest;
+}
+
export async function runBackfillBatch(
opts: BackfillBatchOptions,
): Promise<BackfillBatchResult> {
@@ -292,12 +306,17 @@ export async function runBackfillBatch(
// zero churn — and only an HTTP engine with a baseUrl field (ollama) has an
// endpoint to swap at all.
const llmOps = kinds.map((k) => k.id).filter((id) => LLM_BACKFILL_OPS.has(id));
- const llmFanout = (() => {
- if (llmOps.length === 0) return null;
- const resolved = resolveAttributionTarget("text-only", settings.attribution);
- if (!resolved.app.fields.baseUrl) return null;
- return { config: resolved.config, modelRequested: resolved.modelRequested };
- })();
+ const attrResolved =
+ llmOps.length > 0
+ ? resolveAttributionTarget("text-only", settings.attribution)
+ : null;
+ const llmFanout =
+ attrResolved && attrResolved.app.fields.baseUrl
+ ? {
+ config: attrResolved.config,
+ modelRequested: attrResolved.modelRequested,
+ }
+ : null;
// The remote term is sound only when EVERY kind in this run can take an llm
// lease. runPool's limit is pool-wide: on a mixed run, a slot justified by
// remote capacity could dispatch a non-LLM kind (diarization) onto this
@@ -308,6 +327,126 @@ export async function runBackfillBatch(
const remoteEligible = llmFanout !== null && llmOps.length === kinds.length;
let llmActive = 0;
+ // UNIT-EXECUTOR fan-out: a TAGGED remote worker matching [kind.id,
+ // contendsFor] takes whole units. Tagged-only is the opt-in gate — an
+ // untagged remote predates the unit protocol (and may be an older build with
+ // no /api/worker/unit), and shipping it units it cannot serve would poison
+ // the sweep with refusals. The envelope's identity config is resolved ONCE
+ // per run, appConfig with baseUrl STRIPPED so the executor localises its own
+ // endpoint, and the channel context injected because a scratch corpus has no
+ // digest-context.md to read.
+ const unitConfig: StartWorkerUnitInput["config"] = {
+ ...(attrResolved
+ ? {
+ attribution: settings.attribution,
+ appConfig: stripBaseUrl(attrResolved.config),
+ context: await readDigestContext(opts.paths, opts.channelSlug),
+ }
+ : {}),
+ ...(kinds.some((k) => k.id === "diarization")
+ ? { diarization: settings.diarization }
+ : {}),
+ };
+ let unitActive = 0;
+ const unitRequires = (kind: BackfillKind): string[] => [
+ kind.id,
+ (kind.laneFor?.(settings) ?? kind.lane).contendsFor,
+ ];
+
+ // Ship one candidate to a unit executor; apply its result through the kind's
+ // guarded writers (trap: never a raw file copy — the unit ran against a
+ // snapshot minutes old, and the primary's disk may have moved meanwhile).
+ // Returns null when no tagged remote slot is free or every attempt hit a
+ // transport failure — the caller then runs locally, exactly as today. The
+ // retry shape copies transcribeOne's: transport → degrade-or-mark and try
+ // another worker; work-class → final (the same code would fail the same way
+ // anywhere).
+ const MAX_UNIT_ATTEMPTS = 3;
+ const runBackfillUnit = async (
+ candidate: Candidate,
+ videoDir: string,
+ itemLog: (m: string) => void,
+ runSignal: AbortSignal,
+ ): Promise<BackfillRunOutcome | null> => {
+ const pool = getWorkerPool();
+ for (let attempt = 0; attempt < MAX_UNIT_ATTEMPTS; attempt++) {
+ const lease = pool.tryAcquire(unitRequires(candidate.kind), {
+ kind: "remote",
+ taggedOnly: true,
+ });
+ if (!lease) return null;
+ unitActive++;
+ try {
+ const listing = await readVideoFiles(videoDir, {
+ checkUntranscribable: true,
+ });
+ const files: Record<string, Buffer> = {};
+ for (const name of candidate.kind.inputs(listing)) {
+ files[name] = await readFile(path.join(videoDir, name));
+ }
+ const result = await runUnitViaRemote({
+ worker: lease.worker,
+ op: candidate.kind.id,
+ channelSlug: opts.channelSlug,
+ videoId: candidate.id,
+ files,
+ target: candidate.target,
+ config: unitConfig,
+ force: opts.force,
+ onLog: itemLog,
+ signal: runSignal,
+ });
+ pool.markSuccess(lease.worker.id);
+ if (result.outcome === "done") {
+ const applied = await candidate.kind.applyResult(
+ videoDir,
+ result.files,
+ );
+ if (applied === "applied") return "done";
+ if (applied === "refused") {
+ // A guard said no (the downgrade rule) — the same answer the
+ // local runner reports as outranked/already-present.
+ return "already-present";
+ }
+ itemLog(
+ `Remote unit ${candidate.kind.id} ${candidate.id}: unusable result payload.`,
+ );
+ return "failed";
+ }
+ if (
+ result.outcome === "already-present" ||
+ result.outcome === "missing-input" ||
+ result.outcome === "skipped"
+ ) {
+ return result.outcome;
+ }
+ return "failed";
+ } catch (err) {
+ if (runSignal.aborted || opts.signal?.aborted) throw err;
+ const failureClass =
+ err instanceof TranscribeError ? err.failureClass : "transcription";
+ if (failureClass === "transport") {
+ // A transport failure must NEVER count as work failure — the item
+ // retries on another worker or locally, and only the WORKER pays.
+ if (!(await pingRemoteHealth(lease.worker))) {
+ pool.markDegraded(lease.worker.id);
+ } else {
+ pool.markFailure(lease.worker.id);
+ }
+ itemLog(
+ `Remote unit ${candidate.kind.id} ${candidate.id} transport failure on ${lease.worker.id}: ${String(err)} — retrying elsewhere.`,
+ );
+ continue;
+ }
+ return "failed";
+ } finally {
+ unitActive--;
+ lease.release();
+ }
+ }
+ return null;
+ };
+
// Resolved once per run, not per video: the identity is a settings read plus
// some string work, and deriving it per item is how a counter and a runner end
// up disagreeing about what is stale.
@@ -473,43 +612,51 @@ export async function runBackfillBatch(
if (reacquired.status === "fetched") result.reacquired++;
}
- // A lease on a verified llm endpoint, or null → run on the local/default
- // endpoint exactly as today. Non-parking on purpose: a parked acquire
- // inside a runPool slot would deadlock the batch.
- const llm =
- llmFanout && LLM_BACKFILL_OPS.has(candidate.kind.id)
- ? await acquireLlmSlot(
- candidate.kind.id,
- llmFanout.modelRequested,
- itemLog,
- )
- : null;
- if (llm) llmActive++;
- let outcome: BackfillRunOutcome;
- try {
- outcome = await candidate.kind.run({
- paths: opts.paths,
- videoDir,
- videoId: candidate.id,
- channelSlug: opts.channelSlug,
- target: candidate.target,
- force: opts.force,
- ...(llm
- ? { appConfig: { ...llmFanout!.config, baseUrl: llm.baseUrl } }
- : {}),
- onLog: itemLog,
- signal: runSignal,
- });
- if (llm) getWorkerPool().markSuccess(llm.workerId);
- } catch (err) {
- if (llm && !runSignal.aborted && !opts.signal?.aborted) {
- getWorkerPool().markFailure(llm.workerId);
- }
- throw err;
- } finally {
- if (llm) {
- llmActive--;
- llm.lease.release();
+ // Dispatch ladder: a whole unit to a tagged remote executor when one is
+ // free; else the llm fan-out for the call-bound kinds; else the local
+ // run exactly as today. All claims are non-parking on purpose: a parked
+ // acquire inside a runPool slot would deadlock the batch.
+ let outcome: BackfillRunOutcome | null = await runBackfillUnit(
+ candidate,
+ videoDir,
+ itemLog,
+ runSignal,
+ );
+ if (outcome === null) {
+ const llm =
+ llmFanout && LLM_BACKFILL_OPS.has(candidate.kind.id)
+ ? await acquireLlmSlot(
+ candidate.kind.id,
+ llmFanout.modelRequested,
+ itemLog,
+ )
+ : null;
+ if (llm) llmActive++;
+ try {
+ outcome = await candidate.kind.run({
+ paths: opts.paths,
+ videoDir,
+ videoId: candidate.id,
+ channelSlug: opts.channelSlug,
+ target: candidate.target,
+ force: opts.force,
+ ...(llm
+ ? { appConfig: { ...llmFanout!.config, baseUrl: llm.baseUrl } }
+ : {}),
+ onLog: itemLog,
+ signal: runSignal,
+ });
+ if (llm) getWorkerPool().markSuccess(llm.workerId);
+ } catch (err) {
+ if (llm && !runSignal.aborted && !opts.signal?.aborted) {
+ getWorkerPool().markFailure(llm.workerId);
+ }
+ throw err;
+ } finally {
+ if (llm) {
+ llmActive--;
+ llm.lease.release();
+ }
}
}
if (outcome === "done") {
@@ -624,7 +771,21 @@ export async function runBackfillBatch(
const remoteTerm = remoteEligible
? llmActive + freeLlmSlots(llmOps)
: 0;
- const total = limit + remoteTerm;
+ // The UNIT term: tagged remote slots every kind in this run could take,
+ // plus the leases it already holds. The MIN across kinds keeps it sound
+ // on a mixed run — a slot justified by one kind's remote capacity must
+ // not dispatch another kind onto this box while the local term holds.
+ const unitTerm =
+ unitActive +
+ Math.min(
+ ...kinds.map((k) =>
+ getWorkerPool().freeSlots(unitRequires(k), {
+ kind: "remote",
+ taggedOnly: true,
+ }),
+ ),
+ );
+ const total = limit + remoteTerm + unitTerm;
if (total === 0) {
if (!yielding) {
yielding = true;
diff --git a/common/controller/remoteUnit.test.ts b/common/controller/remoteUnit.test.ts
@@ -0,0 +1,213 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import http from "node:http";
+import type { AddressInfo } from "node:net";
+import type { Worker } from "../lib/workers";
+import { TranscribeError } from "./transcribeError";
+import { runUnitViaRemote } from "./remoteUnit";
+
+// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test controller/remoteUnit.test.ts
+//
+// The unit client against a real node:http stub — no network, no fetch mocks.
+// What is pinned: the submit → poll → result → DELETE round trip; the failure
+// CLASSES (5xx/unreachable/silence → transport, so the batch retries the item
+// elsewhere instead of counting it failed; a 4xx refusal or a "failed" job →
+// work-class, final); and that the scratch DELETE fires even on success.
+
+function worker(baseUrl: string): Worker {
+ return {
+ id: "exec",
+ name: "exec",
+ kind: "remote",
+ enabled: true,
+ priority: 0,
+ tags: ["cpu"],
+ remote: { baseUrl, token: "sekrit" },
+ };
+}
+
+type Stub = {
+ baseUrl: string;
+ requests: Array<{ method: string; url: string; body: string }>;
+ close: () => Promise<void>;
+};
+
+async function startStub(handler: http.RequestListener): Promise<Stub> {
+ const requests: Stub["requests"] = [];
+ const server = http.createServer((req, res) => {
+ let body = "";
+ req.on("data", (c) => (body += c));
+ req.on("end", () => {
+ requests.push({ method: req.method ?? "", url: req.url ?? "", body });
+ handler(req, res);
+ });
+ });
+ await new Promise<void>((resolve) => server.listen(0, "127.0.0.1", resolve));
+ const { port } = server.address() as AddressInfo;
+ return {
+ baseUrl: `http://127.0.0.1:${port}`,
+ requests,
+ close: () => new Promise((resolve) => server.close(() => resolve())),
+ };
+}
+
+function json(res: http.ServerResponse, status: number, body: unknown): void {
+ res.statusCode = status;
+ res.setHeader("content-type", "application/json");
+ res.end(JSON.stringify(body));
+}
+
+test("submit → poll → result → DELETE round trip", async () => {
+ let polls = 0;
+ const stub = await startStub((req, res) => {
+ if (req.method === "POST" && req.url === "/api/worker/unit") {
+ return json(res, 202, { remoteJobId: "job1" });
+ }
+ if (req.url === "/api/worker/unit/job1/events") {
+ polls++;
+ return json(res, 200, { status: polls < 2 ? "running" : "done" });
+ }
+ if (req.url === "/api/worker/unit/job1/result") {
+ return json(res, 200, {
+ outcome: "done",
+ files: { "attribution.json": "{\"videoId\":\"v\"}" },
+ });
+ }
+ if (req.method === "DELETE") return json(res, 200, { ok: true });
+ return json(res, 404, {});
+ });
+ try {
+ const result = await runUnitViaRemote({
+ worker: worker(stub.baseUrl),
+ op: "attribution-text",
+ channelSlug: "chan",
+ videoId: "vid",
+ files: { "transcript.cues.json": Buffer.from("{}") },
+ target: {},
+ pollMs: 10,
+ pollTimeoutMs: 5000,
+ });
+ assert.equal(result.outcome, "done");
+ assert.equal(result.files["attribution.json"], '{"videoId":"v"}');
+ const submit = stub.requests.find((r) => r.method === "POST")!;
+ const body = JSON.parse(submit.body) as {
+ op: string;
+ files: Record<string, string>;
+ };
+ assert.equal(body.op, "attribution-text");
+ assert.equal(
+ Buffer.from(body.files["transcript.cues.json"], "base64").toString(),
+ "{}",
+ );
+ // The fire-and-forget cleanup reaches the stub shortly after.
+ await new Promise((r) => setTimeout(r, 50));
+ assert.ok(stub.requests.some((r) => r.method === "DELETE"));
+ } finally {
+ await stub.close();
+ }
+});
+
+test("a 5xx submit is a TRANSPORT failure (retry elsewhere)", async () => {
+ const stub = await startStub((_req, res) => json(res, 500, { error: "boom" }));
+ try {
+ await assert.rejects(
+ runUnitViaRemote({
+ worker: worker(stub.baseUrl),
+ op: "diarization",
+ channelSlug: "chan",
+ videoId: "vid",
+ files: {},
+ target: {},
+ pollMs: 10,
+ pollTimeoutMs: 200,
+ }),
+ (err: unknown) =>
+ err instanceof TranscribeError && err.failureClass === "transport",
+ );
+ } finally {
+ await stub.close();
+ }
+});
+
+test("a 4xx refusal is a WORK-class failure (no retry)", async () => {
+ const stub = await startStub((_req, res) =>
+ json(res, 400, { error: "not a backfill kind" }),
+ );
+ try {
+ await assert.rejects(
+ runUnitViaRemote({
+ worker: worker(stub.baseUrl),
+ op: "download",
+ channelSlug: "chan",
+ videoId: "vid",
+ files: {},
+ target: {},
+ pollMs: 10,
+ pollTimeoutMs: 200,
+ }),
+ (err: unknown) =>
+ err instanceof TranscribeError && err.failureClass === "transcription",
+ );
+ } finally {
+ await stub.close();
+ }
+});
+
+test("a remote job that finishes failed is a WORK-class failure", async () => {
+ const stub = await startStub((req, res) => {
+ if (req.method === "POST" && req.url === "/api/worker/unit") {
+ return json(res, 202, { remoteJobId: "job1" });
+ }
+ if (req.url?.endsWith("/events")) return json(res, 200, { status: "failed" });
+ return json(res, 200, { ok: true });
+ });
+ try {
+ await assert.rejects(
+ runUnitViaRemote({
+ worker: worker(stub.baseUrl),
+ op: "attribution-text",
+ channelSlug: "chan",
+ videoId: "vid",
+ files: {},
+ target: {},
+ pollMs: 10,
+ pollTimeoutMs: 5000,
+ }),
+ (err: unknown) =>
+ err instanceof TranscribeError && err.failureClass === "transcription",
+ );
+ } finally {
+ await stub.close();
+ }
+});
+
+test("a server that goes silent after submit times out as TRANSPORT", async () => {
+ const stub = await startStub((req, res) => {
+ if (req.method === "POST" && req.url === "/api/worker/unit") {
+ return json(res, 202, { remoteJobId: "job1" });
+ }
+ if (req.method === "DELETE") return json(res, 200, { ok: true });
+ // Every poll errors — the client keeps polling until its silence timeout.
+ return json(res, 500, {});
+ });
+ try {
+ await assert.rejects(
+ runUnitViaRemote({
+ worker: worker(stub.baseUrl),
+ op: "attribution-text",
+ channelSlug: "chan",
+ videoId: "vid",
+ files: {},
+ target: {},
+ pollMs: 10,
+ pollTimeoutMs: 100,
+ }),
+ (err: unknown) =>
+ err instanceof TranscribeError &&
+ err.failureClass === "transport" &&
+ /stopped responding/.test(err.message),
+ );
+ } finally {
+ await stub.close();
+ }
+});
diff --git a/common/controller/remoteUnit.ts b/common/controller/remoteUnit.ts
@@ -0,0 +1,222 @@
+// 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. backfillBatch 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<string, Buffer>;
+ 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 BackfillRunOutcome ("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<string, string>;
+};
+
+type EventsResponse = {
+ status: "queued" | "running" | "done" | "failed" | "cancelled" | "unknown";
+ log?: string;
+};
+
+function authHeaders(worker: Worker): Record<string, string> {
+ 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<RemoteUnitOutcome> {
+ 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<string, string> = {};
+ 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<string, string>;
+ };
+ 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<void> {
+ 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 },
+ );
+ });
+}
diff --git a/common/jobs/workerPool.test.ts b/common/jobs/workerPool.test.ts
@@ -391,6 +391,42 @@ test("an llm worker never satisfies a requirement-less claim", () => {
lease!.release();
});
+test("taggedOnly claims skip untagged workers (the unit-dispatch opt-in gate)", () => {
+ const pool = new WorkerPool();
+ pool.reconfigure(
+ [remoteWorker("plain-remote", 1, 0), worker("cpu-box", 1, ["cpu"])],
+ { applyEnabled: true },
+ );
+ // The untagged remote matches everything under the universal rule, but an
+ // untagged remote predates the unit protocol — units go only to workers the
+ // operator explicitly tagged.
+ assert.equal(
+ pool.tryAcquire(["diarization", "cpu"], {
+ kind: "remote",
+ taggedOnly: true,
+ }),
+ null,
+ );
+ assert.equal(
+ pool.freeSlots(["diarization", "cpu"], { kind: "remote", taggedOnly: true }),
+ 0,
+ );
+ // Tag the remote and the same claim lands.
+ pool.reconfigure(
+ [
+ { ...remoteWorker("plain-remote", 1, 0), tags: ["cpu"] },
+ worker("cpu-box", 1, ["cpu"]),
+ ],
+ { applyEnabled: true },
+ );
+ const lease = pool.tryAcquire(["diarization", "cpu"], {
+ kind: "remote",
+ taggedOnly: true,
+ });
+ assert.equal(lease?.worker.id, "plain-remote");
+ lease?.release();
+});
+
test("an llm worker's slots expand like a remote's", async () => {
const pool = new WorkerPool();
pool.reconfigure([llmWorker("mac", 2)], { applyEnabled: true });
diff --git a/common/jobs/workerPool.ts b/common/jobs/workerPool.ts
@@ -470,16 +470,24 @@ export class WorkerPool {
// REMOTE capacity only ("llm" endpoints, "remote" unit executors) and must
// never grab a local slot out from under the transcription scheduler.
//
+ // opts.taggedOnly restricts the claim to EXPLICITLY TAGGED workers. The unit
+ // dispatcher passes it: an untagged remote worker predates the unit protocol
+ // and matches everything under the untagged-is-universal rule, so without
+ // this gate an existing install's transcription-only remote (possibly an
+ // older build with no /api/worker/unit at all) would be shipped units it
+ // cannot serve. Tagging a worker is the operator's opt-in to unit work.
+ //
// Never starves a parked waiter: if any parked waiter could take the slot this
// would claim, the claim yields (returns null) and pump() serves the waiter.
tryAcquire(
requires: readonly string[],
- opts?: { kind?: WorkerKind },
+ opts?: { kind?: WorkerKind; taggedOnly?: boolean },
): Lease | null {
this.ensureInit();
let best: PoolEntry | null = null;
for (const entry of this.entries.values()) {
if (opts?.kind && entry.config.kind !== opts.kind) continue;
+ if (opts?.taggedOnly && !entry.config.tags?.length) continue;
if (!this.eligible(entry, requires)) continue;
if (!best || entry.config.priority < best.config.priority) best = entry;
}
@@ -496,12 +504,13 @@ export class WorkerPool {
// leasing keeps limit() pure and re-readable every poll.
freeSlots(
requires: readonly string[],
- opts?: { kind?: WorkerKind },
+ opts?: { kind?: WorkerKind; taggedOnly?: boolean },
): number {
this.ensureInit();
let n = 0;
for (const entry of this.entries.values()) {
if (opts?.kind && entry.config.kind !== opts.kind) continue;
+ if (opts?.taggedOnly && !entry.config.tags?.length) continue;
if (this.eligible(entry, requires)) n++;
}
return n;
diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md
@@ -1,6 +1,7 @@
# Changelog
## [Unreleased]
+- **The backfill lane now dispatches units to those executors, and the lane's ceiling rises to match.** For every candidate the runner first tries a free slot on a *tagged* remote worker whose tags cover the operation (tagging is the opt-in: an untagged remote from before this protocol keeps doing transcription only, rather than being shipped envelopes an older build answers with errors). A shipped unit comes back as records, applied through the same guarded writers a local run uses — so a unit's text-only speaker record still loses to a diarized one that landed while it was in flight. A network failure is charged to the *worker*, never the work: the item retries on another executor or falls back to running locally, and a machine that stops answering is health-checked and benched, exactly as remote transcription already does. The lane's concurrency becomes local + remote slots — with the remote share counted at the *minimum* across the run's operations, so on a mixed run a slot justified by one operation's remote capacity can never push a different operation onto this box while it is meant to be standing aside. Digest deliberately does not ship as units: its distribution is the LLM-endpoint fan-out, which reaches the same ollama with less machinery.
- **A machine with no corpus can now run whole backfill units for this one.** The remote-transcription protocol grew a general sibling: `/api/worker/unit` accepts one unit of any *backfill kind* — speaker attribution, diarization — as a small envelope of input files, runs it against a throwaway scratch corpus, and hands the produced sidecar back. Each kind now declares its own contract: what a unit needs (attribution ships the cue sidecar, metadata and raw transcript — the freshness gate compares their mtimes, so the executor writes the cues file *last* or the unit would silently do nothing), what it produces, and how the result lands back on the primary — always through the guarded writers, never a raw copy, so a unit's text-only record still cannot overwrite a diarized one and applying one digest section still preserves the other. Two refusals are load-bearing: the endpoint takes only backfill kinds (downloads and transcription are refused at the door, so download politeness stays one machine's promise), and a unit that reports "disabled" or "not configured" on the executor is an *error* — on a bare box that answer means the primary's injected model/prompt identity was dropped and the executor's default settings leaked in, which is exactly how a second machine writes permanently-stale records. Deploying an executor needs no new software: this app, `ARCHILYZER_IDLE_BOOT=1`, `WORKER_TOKEN`, and an empty transcripts dir — see RUNNING_IN_DOCKER.md.
- **A second machine can now carry the AI backlog by running nothing but `ollama serve`.** The digest and speaker-attribution sweeps — about 194,000 and 78,000 model calls on this archive — bottom out in exactly one HTTP call per chunk; everything around that call (chunking, prompts, parsing, the guarded sidecar writes) is cheap and stays on this box. So the new **LLM endpoint** worker kind is just a URL: no repo, no editor, no token, no copy of the corpus on the other machine. The digest and attribution runners fan their calls across every free endpoint slot alongside the local one, and the lanes' limits rise to match — including while the digest lane is yielding the GPU to transcription, when the *local* term goes to zero and remote endpoints keep the lane moving (an operator pause and the metered spend cap still stop everything; intent and money are global). With no endpoints configured, nothing changes at all.
- **An endpoint that can't serve the exact model is refused, loudly, instead of quietly poisoning the corpus.** What makes a record "current" pins the configured model string — and for digests the context size too — so a second machine running *almost* the right model would write records that this box marks stale and re-does forever, with no error anywhere. Only the endpoint URL may vary per call. Before an endpoint's first use the scheduler asks it (ollama's own `/api/tags`) whether the primary's exact tag is present, and one that lacks it — or is unreachable — is degraded with a log line naming the fix, at the cost of one probe rather than a failure per video. A separate tripwire logs when any engine reports having run a different model than was requested, because freshness deliberately compares the requested string and would never notice on its own. An LLM endpoint can never be handed a transcription, and a worker list of nothing but endpoints won't validate — they cannot transcribe, and letting them count as "enabled workers" would park every transcription forever.
diff --git a/editor/e2e/worker-unit.spec.ts b/editor/e2e/worker-unit.spec.ts
@@ -0,0 +1,170 @@
+// Unit-executor protocol (/api/worker/unit) — the generalisation of the
+// remote-transcription protocol to backfill kinds. The test server runs with
+// WORKER_TOKEN set (see package.json dev:test), so the endpoints are live.
+// Three angles:
+// 1. Auth + the door guard (only backfill KINDS are accepted — download and
+// transcription are refused, which is what keeps download politeness
+// single-machine).
+// 2. A full round trip: an attribution-text unit whose model calls land on an
+// ollama STUB started inside this test — proving the scratch-corpus
+// materialization (cues written last passes the mtime freshness gate), the
+// config injection, and the result pull, with no real model anywhere.
+// 3. Cleanup: DELETE removes the scratch and the result 404s.
+
+import http from "node:http";
+import type { AddressInfo } from "node:net";
+import { test, expect } from "@playwright/test";
+import { resetData } from "./helpers";
+import { baseUrl } from "./baseUrl";
+
+const TOKEN = "test-worker-token";
+const AUTH = { authorization: `Bearer ${TOKEN}` };
+
+test.beforeEach(async () => {
+ await resetData("empty");
+});
+
+test("unit endpoint enforces the bearer token and refuses non-kinds", async ({
+ request,
+}) => {
+ const noAuth = await request.post(`${baseUrl}/api/worker/unit`, {
+ data: { op: "attribution-text", channelSlug: "c", videoId: "v", files: {} },
+ });
+ expect(noAuth.status()).toBe(401);
+
+ // download/transcription are ExternalOperations, not backfill kinds — the
+ // executor refuses them at the door.
+ for (const op of ["download", "transcription", "nonsense"]) {
+ const refused = await request.post(`${baseUrl}/api/worker/unit`, {
+ headers: AUTH,
+ data: { op, channelSlug: "c", videoId: "v", files: {}, target: {} },
+ });
+ expect(refused.status(), op).toBe(400);
+ }
+});
+
+test("an attribution unit round-trips against a scratch corpus and a stub ollama", async ({
+ request,
+}) => {
+ test.setTimeout(60_000);
+ // A fake ollama the EXECUTOR's injected appConfig.baseUrl points at. The
+ // /api/chat reply names one speaker, in the schema the turn prompt pins.
+ const stub = http.createServer((req, res) => {
+ res.setHeader("content-type", "application/json");
+ if (req.url?.startsWith("/api/tags")) {
+ res.end(JSON.stringify({ models: [{ name: "stub-model" }] }));
+ return;
+ }
+ // Drain the request, then answer as ollama would.
+ req.resume();
+ req.on("end", () => {
+ res.end(
+ JSON.stringify({
+ model: "stub-model",
+ message: {
+ content: JSON.stringify({
+ turns: [{ start: "00:00:01", speaker: "Host" }],
+ }),
+ },
+ }),
+ );
+ });
+ });
+ await new Promise<void>((resolve) => stub.listen(0, "127.0.0.1", resolve));
+ const stubUrl = `http://127.0.0.1:${(stub.address() as AddressInfo).port}`;
+
+ try {
+ const cues = {
+ version: 1,
+ id: "unitvid1",
+ title: "Unit test video",
+ channel: "unit-chan",
+ duration: 9,
+ cues: [
+ { start: 0, end: 4, text: "hello there" },
+ { start: 4, end: 9, text: "general kenobi" },
+ ],
+ };
+ const b64 = (s: string) => Buffer.from(s).toString("base64");
+ const post = await request.post(`${baseUrl}/api/worker/unit`, {
+ headers: AUTH,
+ data: {
+ op: "attribution-text",
+ channelSlug: "unit-chan",
+ videoId: "unitvid1",
+ files: {
+ "metadata.info.json": b64(
+ JSON.stringify({ id: "unitvid1", title: "Unit test video", duration: 9 }),
+ ),
+ "transcript.json": b64(JSON.stringify({ transcription: [] })),
+ // Materialized LAST by the executor whatever this map's order is —
+ // the mtime freshness gate depends on it.
+ "transcript.cues.json": b64(JSON.stringify(cues)),
+ },
+ target: {},
+ config: {
+ // The primary's identity, injected. Without this the executor's
+ // default settings (attribution disabled) would fail the job loudly.
+ attribution: {
+ enabled: true,
+ appId: "ollama-direct",
+ model: "stub-model",
+ diarizedEnabled: false,
+ textOnlyEnabled: true,
+ promptVersion: 2,
+ },
+ appConfig: { model: "stub-model", baseUrl: stubUrl, numCtx: 8192 },
+ context: { hash: "none" },
+ },
+ },
+ });
+ expect(post.status()).toBe(202);
+ const { remoteJobId } = await post.json();
+ expect(remoteJobId).toBeTruthy();
+
+ await expect
+ .poll(
+ async () => {
+ const r = await request.get(
+ `${baseUrl}/api/worker/unit/${remoteJobId}/events`,
+ { headers: AUTH },
+ );
+ return ((await r.json()) as { status: string }).status;
+ },
+ { timeout: 30_000 },
+ )
+ .toBe("done");
+
+ const result = await request.get(
+ `${baseUrl}/api/worker/unit/${remoteJobId}/result`,
+ { headers: AUTH },
+ );
+ expect(result.status()).toBe(200);
+ const body = (await result.json()) as {
+ outcome: string;
+ files: Record<string, string>;
+ };
+ expect(body.outcome).toBe("done");
+ const record = JSON.parse(body.files["attribution.json"]) as {
+ speakers: Array<{ label: string }>;
+ provenance: { method: string; model: string };
+ };
+ expect(record.speakers[0]?.label).toBe("Host");
+ expect(record.provenance.method).toBe("text-only");
+ expect(record.provenance.model).toBe("stub-model");
+
+ // Cleanup removes the scratch corpus; the result then 404s.
+ const del = await request.delete(
+ `${baseUrl}/api/worker/unit/${remoteJobId}`,
+ { headers: AUTH },
+ );
+ expect(del.status()).toBe(200);
+ const gone = await request.get(
+ `${baseUrl}/api/worker/unit/${remoteJobId}/result`,
+ { headers: AUTH },
+ );
+ expect(gone.status()).toBe(404);
+ } finally {
+ await new Promise<void>((resolve) => stub.close(() => resolve()));
+ }
+});