commit d6d99ff8c4fd42920ae526065ba3b8ae3ae52b78
parent ed66c5404e188347aea39d174ccce7c988092491
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Mon, 24 Aug 2026 09:59:38 -0400
llm fan-out: a WorkerKind for bare ollama endpoints, leased per call
The ~99%-LLM backlog distributes by fanning CALLS, not corpora: WorkerKind
gains "llm" ({baseUrl, slots}) — a machine contributing nothing but `ollama
serve`. The digest and backfill(attribution) batches tryAcquire a slot per
item and pass the primary's resolved engine config with ONLY baseUrl swapped
(baseUrl is outside the freshness identity, so endpoint choice causes zero
churn); their limits become localTerm + free llm slots, with the digest
yield-to-transcription hold zeroing only the local term (pause and spend cap
stay global), and backfill adding the remote term only on a pure-attribution
run so a mixed run can never dispatch diarization onto a "slot" justified by
remote capacity. Endpoints are verified against /api/tags for the primary's
exact model tag before first use and degraded when it is missing or the box
is dead (controller/llmWorkers.ts); a suspicious engine-reported model
resolution now logs in both runners, because freshness compares the requested
string on purpose and would never notice. An llm worker never matches a
requirement-less lease (transcription can't land on it), and validateWorkers
refuses a worker list whose only enabled entries are endpoints.
appConfig threads BackfillRunOptions → attributeOneVideo →
resolveAttributionTarget, which stops reading getSettings() unconditionally
when an override is supplied — the seam Stage 4's executor injection needs.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Diffstat:
17 files changed, 711 insertions(+), 31 deletions(-)
diff --git a/common/controller/attributeOne.ts b/common/controller/attributeOne.ts
@@ -70,6 +70,7 @@ import {
isCuesJsonFresh,
readNormalizedTranscript,
} from "./normalizeTranscript";
+import { isModelResolutionSuspicious } from "../lib/digest";
export type AttributeOneOptions = {
paths: Paths;
@@ -80,6 +81,13 @@ export type AttributeOneOptions = {
// Overrides settings.attribution. The backfill passes one resolved copy rather
// than re-reading settings per video.
settings?: AttributionSettings;
+ // Overrides the digest app's config block (settings.digest.apps[appId]).
+ // Two callers: the LLM fan-out, which passes the primary's resolved config
+ // with only `baseUrl` swapped to the leased endpoint (baseUrl is not part of
+ // the freshness identity, so endpoint choice causes zero churn); and the
+ // unit executor, which passes the primary's injected identity config so a
+ // bare box's default settings can never leak into provenance.
+ appConfig?: DigestAppConfig;
// Redo even when the recorded identity matches. Never overrides the DOWNGRADE
// rule — forcing a regeneration is not the same as asking for a worse record,
// and nothing in the UI should be able to request the second by accident.
@@ -144,6 +152,7 @@ export async function attributeOneVideo(
const { target, app, config, modelRequested } = resolveAttributionTarget(
method,
cfg,
+ opts.appConfig,
);
// The per-video half of the diarized identity. See
// AttributionProvenance.diarizationGeneratedAt: a cluster index means nothing
@@ -229,6 +238,15 @@ export async function attributeOneVideo(
return "failed";
}
+ // Same tripwire as digestVideo's: freshness compares the REQUESTED model
+ // string, so an endpoint serving the wrong weights would never invalidate
+ // anything — the log line is the only thing that can catch it.
+ if (isModelResolutionSuspicious(modelRequested, reportedModel)) {
+ log(
+ `Attribute ${opts.videoId}: engine reported model "${reportedModel}" for requested "${modelRequested}" — check the serving endpoint has the right weights.`,
+ );
+ }
+
if (speakers.length === 0) {
// Do NOT write. An empty record carrying the current identity would read as
// fresh, and this video would never be retried — the same trap digestVideo
diff --git a/common/controller/attributionTarget.ts b/common/controller/attributionTarget.ts
@@ -47,17 +47,24 @@ export type ResolvedAttributionTarget = {
// diarization.json it names clusters from, and that is a disk read. Callers that
// have the value pass it in; the registry's state() adds it only inside the one
// branch that has already paid for the read.
+// `appConfig` overrides the digest app's config block wholesale — the fan-out
+// passes the primary's resolved config with only baseUrl swapped, and the unit
+// executor passes the primary's injected identity config. When it is supplied,
+// getSettings() is NOT consulted for the config: on a bare executor the
+// settings file is defaults, and defaults leaking into the identity here is
+// exactly how a remote machine writes permanently-stale records.
export function resolveAttributionTarget(
method: AttributionMethod,
cfg?: AttributionSettings,
+ appConfig?: DigestAppConfig,
): ResolvedAttributionTarget {
- const settings = getSettings();
- const attribution = cfg ?? settings.attribution;
+ const attribution = cfg ?? getSettings().attribution;
const app = getDigestApp(attribution.appId);
// The digest app's OWN config block — the ollama URL, context size and
// timeout. Attribution is a digest-app workload; a second copy of that config
// would be one more thing to keep in step for no benefit.
- const config: DigestAppConfig = settings.digest.apps[app.id] ?? {};
+ const config: DigestAppConfig =
+ appConfig ?? getSettings().digest.apps[app.id] ?? {};
const modelRequested =
attribution.model.trim() || config.model?.trim() || app.defaultModel();
return {
diff --git a/common/controller/backfillBatch.ts b/common/controller/backfillBatch.ts
@@ -49,6 +49,9 @@ import {
type BackfillRunOutcome,
} from "../lib/backfillKinds";
import { transcriptionActivity } from "./digestYield";
+import { acquireLlmSlot, freeLlmSlots } from "./llmWorkers";
+import { getWorkerPool } from "../jobs/workerPool";
+import { resolveAttributionTarget } from "./attributionTarget";
import { reacquireMediaFor, type ReacquireOutcome } from "./backfillReacquire";
import type { AutoQueueOrder } from "../jobs/autoQueuePolicy";
import { batchRecencyComparator } from "./batchRecency";
@@ -222,6 +225,11 @@ export function candidateAction(
type Candidate = { id: string; kind: BackfillKind; target: unknown };
+// The lane kinds whose work bottoms out in model calls an "llm" endpoint
+// worker can serve. Diarization is deliberately absent — its work is an audio
+// pass, which only a unit executor (a full instance of this app) can take.
+const LLM_BACKFILL_OPS = new Set(["attribution-text", "attribution-diarized"]);
+
export async function runBackfillBatch(
opts: BackfillBatchOptions,
): Promise<BackfillBatchResult> {
@@ -278,6 +286,28 @@ export async function runBackfillBatch(
);
const wanted = byRecency ? [...wantedListed].sort(byRecency) : wantedListed;
+ // LLM fan-out, resolved ONCE per run: the primary's engine config and model
+ // for the attribution kinds. Only `baseUrl` will vary per leased endpoint —
+ // baseUrl is not part of the freshness identity, so endpoint choice causes
+ // 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 };
+ })();
+ // 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
+ // box's own CPU while the local term says stand aside — breaking the
+ // idle-only default. A pure attribution run (the scoped sweep) gets the full
+ // fan-out; a mixed run keeps today's local limit and still fans out whatever
+ // its local slots dispatch.
+ const remoteEligible = llmFanout !== null && llmOps.length === kinds.length;
+ let llmActive = 0;
+
// 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.
@@ -443,16 +473,45 @@ export async function runBackfillBatch(
if (reacquired.status === "fetched") result.reacquired++;
}
- const outcome: BackfillRunOutcome = await candidate.kind.run({
- paths: opts.paths,
- videoDir,
- videoId: candidate.id,
- channelSlug: opts.channelSlug,
- target: candidate.target,
- force: opts.force,
- onLog: itemLog,
- signal: runSignal,
- });
+ // 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();
+ }
+ }
if (outcome === "done") {
result.attempted++;
result.succeeded++;
@@ -557,7 +616,16 @@ export async function runBackfillBatch(
slots: live.concurrency,
primaryBusy: activity.busy,
});
- if (limit === 0) {
+ // The REMOTE term: llm-endpoint slots serving this run's kinds, plus the
+ // leases it already holds. transcriptionActivity zeroes only the LOCAL
+ // term — an endpoint contends for nothing on this box — while the
+ // lane-disabled hold above still zeroes everything (operator intent is
+ // global). See remoteEligible for why a mixed-kind run gets no term.
+ const remoteTerm = remoteEligible
+ ? llmActive + freeLlmSlots(llmOps)
+ : 0;
+ const total = limit + remoteTerm;
+ if (total === 0) {
if (!yielding) {
yielding = true;
log(
@@ -568,9 +636,13 @@ export async function runBackfillBatch(
}
if (yielding) {
yielding = false;
- log("Resuming the backfill lane.");
+ log(
+ limit === 0
+ ? "Resuming the backfill lane on remote LLM endpoints only."
+ : "Resuming the backfill lane.",
+ );
}
- return limit;
+ return total;
},
signal: opts.signal ?? NEVER,
drainSignal: opts.drainSignal ?? NEVER,
diff --git a/common/controller/digestBatch.ts b/common/controller/digestBatch.ts
@@ -40,6 +40,8 @@ import {
digestPreflight,
laneSharesDuplicates,
} from "./laneGuards";
+import { acquireLlmSlot, freeLlmSlots } from "./llmWorkers";
+import { getWorkerPool } from "../jobs/workerPool";
import type { AutoQueueOrder } from "../jobs/autoQueuePolicy";
import { batchRecencyComparator } from "./batchRecency";
import {
@@ -395,6 +397,14 @@ export async function runDigestBatch(
});
};
+ // LLM fan-out: an "llm" endpoint worker (a box running nothing but `ollama
+ // serve`) can take this lane's model calls. Only for an HTTP engine with a
+ // baseUrl field (ollama) — the metered CLI lane has no endpoint to swap.
+ // `llmActive` counts leases THIS run holds, so limit() stays stable at
+ // local + total-llm-capacity instead of sagging while slots are leased.
+ const fanOutEligible = app.fields.baseUrl === true;
+ let llmActive = 0;
+
const runOne = async (
candidate: Candidate,
runSignal: AbortSignal,
@@ -408,6 +418,14 @@ export async function runDigestBatch(
// a task-count average is the wrong denominator for a sweep ETA.
audioSeconds: candidate.duration,
});
+ // A lease on a verified endpoint, or null → the local/default endpoint
+ // exactly as today. ONLY baseUrl varies: model, numCtx and the rest of the
+ // config travel verbatim, so which endpoint served a call is invisible to
+ // the freshness identity.
+ const llm = fanOutEligible
+ ? await acquireLlmSlot("digest", modelRequested, task ? task.onLog : log)
+ : null;
+ if (llm) llmActive++;
try {
const outcome = await digestVideo({
paths: opts.paths,
@@ -415,7 +433,7 @@ export async function runDigestBatch(
videoId: candidate.id,
sections,
appId: app.id,
- config,
+ config: llm ? { ...config, baseUrl: llm.baseUrl } : config,
context,
timestampMode,
...(digestSettings.promptVariant
@@ -425,6 +443,9 @@ export async function runDigestBatch(
onLog: task ? task.onLog : opts.onLog,
signal: runSignal,
});
+ // The call round-tripped — the endpoint is alive. Chunk-level failures
+ // are digestVideo's per-chunk isolation and log loudly on their own.
+ if (llm) getWorkerPool().markSuccess(llm.workerId);
if (outcome.status === "fresh") {
result.fresh++;
return;
@@ -461,12 +482,19 @@ export async function runDigestBatch(
reportProgress();
} catch (err) {
if (runSignal.aborted || opts.signal?.aborted) throw err;
+ // Counted against the worker too, so a flapping endpoint degrades
+ // through the pool's existing consecutive-failure machinery.
+ if (llm) getWorkerPool().markFailure(llm.workerId);
result.attempted++;
result.failed++;
log(
`Failed ${candidate.id}: ${(err as Error)?.message ?? String(err)}`,
);
} finally {
+ if (llm) {
+ llmActive--;
+ llm.lease.release();
+ }
resolvedAudioSeconds += candidate.duration;
task?.end();
}
@@ -507,7 +535,27 @@ export async function runDigestBatch(
// is read off the live accumulator rather than from disk.
costUsd: result.costUsd,
});
+ // The REMOTE term: leases this run already holds plus free llm slots
+ // serving digest. Counting `llmActive` keeps the limit stable at
+ // local + total-llm-capacity rather than sagging as slots are leased.
+ const remoteTerm = fanOutEligible
+ ? llmActive + freeLlmSlots(["digest"])
+ : 0;
if (gate.hold) {
+ // The yield-to-transcription hold zeroes ONLY the local term: an llm
+ // endpoint contends for nothing on this box's GPU, so parking it too
+ // would idle remote capacity to protect hardware it never touches. An
+ // operator pause and the spend cap still zero BOTH terms — intent and
+ // money are global.
+ if (gate.reason === "yield" && remoteTerm > 0) {
+ if (heldReason !== "yield") {
+ heldReason = "yield";
+ log(
+ `${gate.message} Remote LLM endpoint(s) keep the lane moving meanwhile.`,
+ );
+ }
+ return remoteTerm;
+ }
// LOGGED ON THE EDGE ONLY, and the edge is tracked here rather than in
// the rule: the rule is pure and a line per poll would bury a job log
// over a multi-week run. An operator watching a lane sit at zero
@@ -525,7 +573,7 @@ export async function runDigestBatch(
log("Transcription finished; resuming the digest lane.");
}
heldReason = null;
- return concurrency;
+ return concurrency + remoteTerm;
},
signal: opts.signal ?? NEVER,
drainSignal: opts.drainSignal ?? NEVER,
diff --git a/common/controller/digestVideo.ts b/common/controller/digestVideo.ts
@@ -34,6 +34,7 @@ import {
import { getDigestApp, type DigestAppConfig } from "../lib/digestApps";
import {
digestPromptVariant,
+ isModelResolutionSuspicious,
isSectionFresh,
type DigestItem,
type DigestTimestampMode,
@@ -249,6 +250,16 @@ export async function digestVideo(
}
}
+ // A tag completion ("qwen2.5" → "qwen2.5:7b") is normal; a different model
+ // is an endpoint serving the wrong weights. Freshness compares the
+ // REQUESTED string, so this would never invalidate anything — the loud log
+ // line is the only tripwire, which is why it is here and not optional.
+ if (isModelResolutionSuspicious(modelRequested, reportedModel)) {
+ log(
+ `${opts.channelSlug}/${opts.videoId}: engine reported model "${reportedModel}" for requested "${modelRequested}" — check the serving endpoint has the right weights.`,
+ );
+ }
+
if (outputs.length === 0) {
// Nothing usable. Do NOT write a section — an empty section with current
// provenance would read as "fresh" and the video would never be retried.
diff --git a/common/controller/llmWorkers.test.ts b/common/controller/llmWorkers.test.ts
@@ -0,0 +1,122 @@
+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 { WorkerPool } from "../jobs/workerPool";
+import {
+ acquireLlmSlot,
+ clearLlmVerificationCache,
+ freeLlmSlots,
+} from "./llmWorkers";
+
+// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test controller/llmWorkers.test.ts
+//
+// A fake ollama (/api/tags) on a loopback port drives the verification path:
+// an endpoint serving the model tag is leased, one that lacks it (or is dead)
+// is DEGRADED and the next candidate tried — one probe per bad endpoint, not a
+// failure per item. Pools are private instances, never the global singleton.
+
+async function withOllamaStub(
+ models: string[],
+ fn: (baseUrl: string) => Promise<void>,
+): Promise<void> {
+ const server = http.createServer((req, res) => {
+ res.setHeader("content-type", "application/json");
+ if (req.url?.startsWith("/api/tags")) {
+ res.end(JSON.stringify({ models: models.map((name) => ({ name })) }));
+ } else {
+ res.statusCode = 404;
+ res.end("{}");
+ }
+ });
+ await new Promise<void>((resolve) => server.listen(0, "127.0.0.1", resolve));
+ const { port } = server.address() as AddressInfo;
+ try {
+ await fn(`http://127.0.0.1:${port}`);
+ } finally {
+ await new Promise<void>((resolve) => server.close(() => resolve()));
+ }
+}
+
+function llmWorker(id: string, baseUrl: string, priority = 0): Worker {
+ return {
+ id,
+ name: id,
+ kind: "llm",
+ enabled: true,
+ priority,
+ llm: { baseUrl },
+ };
+}
+
+const quiet = () => {};
+
+test("acquireLlmSlot leases a verified endpoint and hands back its baseUrl", async () => {
+ clearLlmVerificationCache();
+ await withOllamaStub(["qwen2.5:7b"], async (baseUrl) => {
+ const pool = new WorkerPool();
+ pool.reconfigure([llmWorker("mac", baseUrl)], { applyEnabled: true });
+ const slot = await acquireLlmSlot("digest", "qwen2.5:7b", quiet, pool);
+ assert.ok(slot, "expected a lease on the serving endpoint");
+ assert.equal(slot!.baseUrl, baseUrl);
+ assert.equal(freeLlmSlots(["digest"], pool), 0);
+ slot!.lease.release();
+ assert.equal(freeLlmSlots(["digest"], pool), 1);
+ });
+});
+
+test("a bare model tag matches the endpoint's :latest", async () => {
+ clearLlmVerificationCache();
+ await withOllamaStub(["qwen2.5:latest"], async (baseUrl) => {
+ const pool = new WorkerPool();
+ pool.reconfigure([llmWorker("mac", baseUrl)], { applyEnabled: true });
+ const slot = await acquireLlmSlot("digest", "qwen2.5", quiet, pool);
+ assert.ok(slot);
+ slot!.lease.release();
+ });
+});
+
+test("an endpoint missing the tag is degraded, and the next candidate serves", async () => {
+ clearLlmVerificationCache();
+ await withOllamaStub(["llama3:8b"], async (wrongUrl) => {
+ await withOllamaStub(["qwen2.5:7b"], async (rightUrl) => {
+ const pool = new WorkerPool();
+ pool.reconfigure(
+ [llmWorker("wrong", wrongUrl, 0), llmWorker("right", rightUrl, 1)],
+ { applyEnabled: true },
+ );
+ const logs: string[] = [];
+ const slot = await acquireLlmSlot(
+ "digest",
+ "qwen2.5:7b",
+ (m) => logs.push(m),
+ pool,
+ );
+ assert.ok(slot, "the second endpoint should have been leased");
+ assert.equal(slot!.workerId, "right");
+ slot!.lease.release();
+ // The wrong endpoint is OUT — degraded, one probe, no per-item retries.
+ assert.equal(
+ pool.summary().find((w) => w.id === "wrong")?.degraded,
+ true,
+ );
+ assert.ok(logs.some((m) => /does not serve/.test(m)));
+ // And a later claim goes straight to the survivor.
+ const again = await acquireLlmSlot("digest", "qwen2.5:7b", quiet, pool);
+ assert.equal(again?.workerId, "right");
+ again?.lease.release();
+ });
+ });
+});
+
+test("a dead endpoint is degraded rather than looped on", async () => {
+ clearLlmVerificationCache();
+ const pool = new WorkerPool();
+ pool.reconfigure([llmWorker("dead", "http://127.0.0.1:59599")], {
+ applyEnabled: true,
+ });
+ const slot = await acquireLlmSlot("digest", "qwen2.5:7b", quiet, pool, 500);
+ assert.equal(slot, null);
+ assert.equal(pool.summary().find((w) => w.id === "dead")?.degraded, true);
+});
diff --git a/common/controller/llmWorkers.ts b/common/controller/llmWorkers.ts
@@ -0,0 +1,107 @@
+// The LLM fan-out's claim on an endpoint worker: lease a slot, make sure the
+// endpoint actually serves the primary's model tag, hand back the baseUrl.
+//
+// WHY VERIFICATION IS NOT OPTIONAL. Freshness identity pins the CONFIGURED
+// model string (and, for digest, numCtx via the prompt variant); only baseUrl
+// may vary per call. An endpoint that silently served a different model would
+// not fail — it would write records whose provenance looks current and whose
+// content came from the wrong model. The cheap guard is ollama's own /api/tags:
+// before an endpoint's first use (and every VERIFY_TTL_MS thereafter) the
+// dispatcher confirms the tag is present, and an endpoint that lacks it — or
+// cannot be reached — is DEGRADED through the pool's existing machinery, not
+// retried per item.
+//
+// This module is deliberately thin: no queue, no state beyond the verification
+// cache. Claiming is the pool's tryAcquire (non-parking — a parked acquire
+// inside a runPool slot would deadlock the batch), and capacity for a
+// dispatcher's limit() is the pool's freeSlots.
+
+import { getWorkerPool, type Lease, type WorkerPool } from "../jobs/workerPool";
+
+export type LlmSlot = {
+ lease: Lease;
+ baseUrl: string;
+ workerId: string;
+};
+
+const VERIFY_TTL_MS = 10 * 60_000;
+
+// worker.id -> when it was last confirmed to serve which model.
+const verified = new Map<string, { at: number; model: string }>();
+
+// Exported for tests, which need a clean slate between fixtures.
+export function clearLlmVerificationCache(): void {
+ verified.clear();
+}
+
+// The remote term of a dispatcher's limit(): free llm slots serving any of the
+// given operations. Pure counting — no lease is taken.
+export function freeLlmSlots(
+ ops: readonly string[],
+ pool: WorkerPool = getWorkerPool(),
+): number {
+ return pool.freeSlots(ops, { kind: "llm" });
+}
+
+async function servesModel(
+ baseUrl: string,
+ model: string,
+ timeoutMs: number,
+): Promise<boolean> {
+ try {
+ const res = await fetch(`${baseUrl}/api/tags`, {
+ signal: AbortSignal.timeout(timeoutMs),
+ });
+ if (!res.ok) return false;
+ const body = (await res.json()) as { models?: Array<{ name?: unknown }> };
+ const names = (body.models ?? [])
+ .map((m) => m.name)
+ .filter((n): n is string => typeof n === "string");
+ // ollama names a bare tag "<model>:latest"; a config of "qwen2.5" matches
+ // it. Anything more specific must match exactly.
+ return names.some(
+ (n) => n === model || (!model.includes(":") && n === `${model}:latest`),
+ );
+ } catch {
+ return false;
+ }
+}
+
+// Claim a verified endpoint slot for one operation, or null when none is free.
+// An endpoint that fails verification (missing tag, unreachable) is degraded
+// and the next candidate tried, so one bad endpoint costs one probe rather
+// than poisoning the run.
+export async function acquireLlmSlot(
+ op: string,
+ model: string,
+ log: (msg: string) => void,
+ pool: WorkerPool = getWorkerPool(),
+ timeoutMs = 3000,
+): Promise<LlmSlot | null> {
+ for (;;) {
+ const lease = pool.tryAcquire([op], { kind: "llm" });
+ if (!lease) return null;
+ const workerId = lease.worker.id;
+ const baseUrl = (lease.worker.llm?.baseUrl ?? "").replace(/\/+$/, "");
+ if (!baseUrl) {
+ // Sanitizers should make this unreachable; degrade rather than loop on it.
+ pool.markDegraded(workerId);
+ lease.release();
+ continue;
+ }
+ const hit = verified.get(workerId);
+ const fresh =
+ hit && hit.model === model && Date.now() - hit.at < VERIFY_TTL_MS;
+ if (fresh || (await servesModel(baseUrl, model, timeoutMs))) {
+ if (!fresh) verified.set(workerId, { at: Date.now(), model });
+ return { lease, baseUrl, workerId };
+ }
+ verified.delete(workerId);
+ log(
+ `LLM endpoint ${workerId} (${baseUrl}) does not serve "${model}" (or is unreachable) — ` +
+ `degraded for this run. Pull the exact model tag there, or re-enable it from the Workers page once fixed.`,
+ );
+ pool.markDegraded(workerId);
+ lease.release();
+ }
+}
diff --git a/common/jobs/workerPool.test.ts b/common/jobs/workerPool.test.ts
@@ -363,6 +363,48 @@ test("a slot added by expansion inherits the base slot's runtime state", async (
);
});
+// --- LLM endpoint workers ---
+
+function llmWorker(id: string, slots?: number, tags?: string[]): Worker {
+ return {
+ id,
+ name: id,
+ kind: "llm",
+ enabled: true,
+ priority: 0,
+ ...(tags ? { tags } : {}),
+ llm: { baseUrl: "http://127.0.0.1:9", ...(slots ? { slots } : {}) },
+ };
+}
+
+test("an llm worker never satisfies a requirement-less claim", () => {
+ const pool = new WorkerPool();
+ pool.reconfigure([llmWorker("mac")], { applyEnabled: true });
+ // Today's transcription path: no requirement. The endpoint must be invisible
+ // to it — including to the "paused — waiting for a worker" probe.
+ assert.equal(pool.tryAcquire([]), null);
+ assert.equal(pool.hasEligibleWorker(), false);
+ // But it serves its LLM operations.
+ assert.equal(pool.hasEligibleWorker(["digest"]), true);
+ const lease = pool.tryAcquire(["digest"], { kind: "llm" });
+ assert.ok(lease);
+ 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 });
+ assert.equal(pool.freeSlots(["digest"], { kind: "llm" }), 2);
+ const l1 = pool.tryAcquire(["digest"], { kind: "llm" });
+ assert.ok(l1);
+ assert.equal(pool.freeSlots(["digest"], { kind: "llm" }), 1);
+ const l2 = pool.tryAcquire(["digest"], { kind: "llm" });
+ assert.ok(l2);
+ assert.equal(pool.tryAcquire(["digest"], { kind: "llm" }), null);
+ l1!.release();
+ l2!.release();
+});
+
test("legacy { background: true } maps to the background tier", async () => {
const { pool, lease } = await busyPool();
const order: string[] = [];
diff --git a/common/jobs/workerPool.ts b/common/jobs/workerPool.ts
@@ -173,7 +173,14 @@ export class WorkerPool {
// independently enable/drain/degrade-able entry sharing one remote
// config. Shrinking the count retires the surplus through the exact
// removed-worker path below (drain if busy, drop when idle).
- const slotCount = w.kind === "remote" ? this.remoteSlotCount(w) : 1;
+ const slotCount =
+ w.kind === "remote"
+ ? this.remoteSlotCount(w)
+ : w.kind === "llm"
+ ? // One concurrent generation per endpoint unless the operator says
+ // otherwise — no probe: an ollama box does not report a slot count.
+ Math.max(1, Math.floor(w.llm?.slots ?? 1))
+ : 1;
for (let slot = 1; slot <= slotCount; slot++) {
const cfg: Worker =
slot === 1
@@ -484,6 +491,22 @@ export class WorkerPool {
return this.grant(claimed);
}
+ // How many matching slots are FREE right now — the remote term of a
+ // dispatcher's limit() (localTerm + freeSlots(...)). Counting rather than
+ // leasing keeps limit() pure and re-readable every poll.
+ freeSlots(
+ requires: readonly string[],
+ opts?: { kind?: WorkerKind },
+ ): number {
+ this.ensureInit();
+ let n = 0;
+ for (const entry of this.entries.values()) {
+ if (opts?.kind && entry.config.kind !== opts.kind) continue;
+ if (this.eligible(entry, requires)) n++;
+ }
+ return n;
+ }
+
// --- Runtime controls (Workers page). Transient operator overrides: they
// mutate the live pool's runtime `state` only, never settings.json. A restart
// returns workers to their configured enabled state. They survive the
diff --git a/common/lib/backfillKinds.ts b/common/lib/backfillKinds.ts
@@ -67,6 +67,7 @@ import {
} from "./queueKeys";
import {
isSectionFresh,
+ type DigestAppConfig,
type DigestFreshnessTarget,
type DigestLane,
type DigestSectionKind,
@@ -223,6 +224,12 @@ export type BackfillRunOptions = {
target: unknown;
// Redo a `present` video anyway (an operator forcing a regeneration).
force?: boolean;
+ // Engine-config override for the LLM-backed kinds (attribution today). The
+ // fan-out passes the primary's resolved config with only `baseUrl` swapped
+ // to a leased endpoint — baseUrl is not part of the freshness identity, so
+ // which endpoint served a call can vary freely while everything that IS
+ // identity (model, numCtx, timeouts) travels verbatim from the primary.
+ appConfig?: DigestAppConfig;
onLog?: (msg: string) => void;
signal?: AbortSignal;
};
@@ -599,6 +606,7 @@ const attributionText: BackfillKind = {
channelSlug: opts.channelSlug,
method: "text-only",
force: opts.force,
+ appConfig: opts.appConfig,
onLog: opts.onLog,
signal: opts.signal,
}),
@@ -677,6 +685,7 @@ const attributionDiarized: BackfillKind = {
channelSlug: opts.channelSlug,
method: "diarized",
force: opts.force,
+ appConfig: opts.appConfig,
onLog: opts.onLog,
signal: opts.signal,
}),
diff --git a/common/lib/digest.ts b/common/lib/digest.ts
@@ -171,6 +171,25 @@ export type DigestAppConfig = {
timeoutMs?: number;
};
+// Whether an engine-reported model resolution looks like a DIFFERENT model
+// rather than a benign tag completion. "qwen2.5" resolving to "qwen2.5:7b" or
+// "qwen2.5:latest" is ollama filling in a tag; "qwen2.5:7b" coming back as
+// "llama3:8b" means the endpoint served the wrong weights. Freshness compares
+// the REQUESTED string on purpose (see AttributionProvenance.modelRequested),
+// so a wrong resolution would not invalidate anything — which is exactly why
+// the runners log it loudly instead.
+export function isModelResolutionSuspicious(
+ requested: string,
+ reported: string,
+): boolean {
+ if (!reported || !requested || reported === requested) return false;
+ if (reported === `${requested}:latest`) return false;
+ if (!requested.includes(":") && reported.startsWith(`${requested}:`)) {
+ return false;
+ }
+ return true;
+}
+
// The two generated sections. Per-SECTION provenance (not per-file) because the
// controller does a read-modify-write merge, so one video can legitimately hold
// local chapters and metered tags.
diff --git a/common/lib/transcribeOutcome.ts b/common/lib/transcribeOutcome.ts
@@ -13,7 +13,9 @@ export type TranscribeOutcomeRecord = {
// Which worker produced it (id + engine app), for later attribution. Optional
// so older/hand-written records still validate; appId is absent for remote
// workers (they delegate to another app instance).
- worker?: { id: string; appId?: string; kind: "local" | "remote" };
+ // `kind` mirrors WorkerKind. "llm" can never actually appear (an llm worker
+ // never matches a transcription lease) but the type tracks the worker model.
+ worker?: { id: string; appId?: string; kind: "local" | "remote" | "llm" };
// Wall-clock transcription time in milliseconds.
durationMs?: number;
};
diff --git a/common/lib/workers.test.ts b/common/lib/workers.test.ts
@@ -4,6 +4,7 @@ import {
WORKER_RESOURCE_TAGS,
sanitizeWorkers,
unknownWorkerTags,
+ validateWorkers,
workerMatches,
type Worker,
} from "./workers";
@@ -36,6 +37,61 @@ test("workerMatches: a tagged worker matches on intersection only", () => {
assert.equal(workerMatches(w, ["network"]), false);
});
+test("workerMatches: an llm worker never takes a requirement-less lease", () => {
+ // Today's transcription acquire passes no requirement; an ollama endpoint
+ // handed a transcription job cannot run it. This is the one place the
+ // untagged-is-universal rule inverts.
+ assert.equal(workerMatches({ kind: "llm" }), false);
+ assert.equal(workerMatches({ kind: "llm", tags: [] }, []), false);
+});
+
+test("workerMatches: an untagged llm worker serves all three LLM operations", () => {
+ assert.equal(workerMatches({ kind: "llm" }, ["digest"]), true);
+ assert.equal(workerMatches({ kind: "llm" }, ["attribution-text"]), true);
+ assert.equal(workerMatches({ kind: "llm" }, ["attribution-diarized"]), true);
+ // …and nothing else.
+ assert.equal(workerMatches({ kind: "llm" }, ["transcription"]), false);
+ assert.equal(workerMatches({ kind: "llm" }, ["diarization"]), false);
+});
+
+test("workerMatches: tags narrow an llm worker to specific operations", () => {
+ const w = { kind: "llm" as const, tags: ["digest"] };
+ assert.equal(workerMatches(w, ["digest"]), true);
+ assert.equal(workerMatches(w, ["attribution-text"]), false);
+});
+
+test("sanitizeWorkers: an llm worker keeps its endpoint config", () => {
+ const [w] = sanitizeWorkers([
+ {
+ name: "mac",
+ kind: "llm",
+ llm: { baseUrl: " http://mac.lan:11434 ", slots: 2.7 },
+ },
+ ]);
+ assert.equal(w.kind, "llm");
+ assert.equal(w.llm?.baseUrl, "http://mac.lan:11434");
+ assert.equal(w.llm?.slots, 2);
+ assert.equal(w.remote, undefined);
+});
+
+test("validateWorkers: an llm-only list cannot satisfy the enabled-worker rule", () => {
+ const [llm] = sanitizeWorkers([
+ { name: "mac", kind: "llm", llm: { baseUrl: "http://mac.lan:11434" } },
+ ]);
+ assert.match(
+ validateWorkers([llm]) ?? "",
+ /transcription worker .* must be enabled/i,
+ );
+ const [local] = sanitizeWorkers([{ name: "cpu", kind: "local" }]);
+ assert.equal(validateWorkers([local, llm]), null);
+});
+
+test("validateWorkers: an llm worker needs a valid http(s) base URL", () => {
+ const [bad] = sanitizeWorkers([{ name: "mac", kind: "llm", llm: {} }]);
+ const [local] = sanitizeWorkers([{ name: "cpu", kind: "local" }]);
+ assert.match(validateWorkers([local, bad]) ?? "", /needs a base URL/);
+});
+
test("sanitizeWorkers trims tags, drops empties, and de-duplicates", () => {
const [w] = sanitizeWorkers([
{
diff --git a/common/lib/workers.ts b/common/lib/workers.ts
@@ -20,7 +20,28 @@ import {
validateTranscribeArgs,
} from "./transcriptionApps";
-export type WorkerKind = "local" | "remote";
+export type WorkerKind = "local" | "remote" | "llm";
+
+// An LLM ENDPOINT worker: a machine contributing nothing but `ollama serve`.
+// The primary keeps orchestrating (chunking, prompts, parsing, the guarded
+// sidecar writes) and fans the model CALLS out to it — so the box needs no
+// repo, no editor, no token, no corpus. This is the lightest correct way to
+// distribute the digest/attribution backlog, which bottoms out in exactly one
+// ollama HTTP call per chunk.
+//
+// THE HOMOGENEITY CONSTRAINT: freshness identity pins the configured model
+// string (and, for digest, numCtx via the prompt variant). Only `baseUrl` is
+// localised per call — every other config field travels from the primary — so
+// an endpoint must serve the primary's exact model tag. One that cannot is not
+// "close enough": its output would be permanently-stale. The dispatcher
+// verifies the tag against /api/tags before first use and degrades the worker
+// when it is missing (controller/llmWorkers.ts).
+export type LlmWorkerConfig = {
+ baseUrl: string; // e.g. http://macbook.lan:11434
+ // Concurrent generations to allow this endpoint. Defaults to 1 — one model
+ // instance, one generation — unless the operator knows better.
+ slots?: number;
+};
// Where a remote worker delegates. The remote runs its OWN worker pool and picks
// among ITS local workers — so the primary stores only how to reach it, not which
@@ -73,6 +94,9 @@ export type Worker = {
// REMOTE: how to reach the delegate instance.
remote?: RemoteWorkerConfig;
+
+ // LLM: how to reach the bare model endpoint.
+ llm?: LlmWorkerConfig;
};
// The contended-resource half of the tag vocabulary — BackfillLane.contendsFor's
@@ -81,6 +105,14 @@ export type Worker = {
// A drift between the two lists is caught by a test, not by the type system.
export const WORKER_RESOURCE_TAGS = ["gpu", "cpu", "network"] as const;
+// The operations an LLM endpoint can serve — the three whose work bottoms out
+// in one model call. An UNTAGGED llm worker serves all three; tags narrow it.
+export const LLM_OPERATION_TAGS = [
+ "digest",
+ "attribution-text",
+ "attribution-diarized",
+] as const;
+
// THE MATCHING RULE, in one exported place so the pool, the dispatchers and the
// sweep console cannot disagree about who can take what.
//
@@ -94,10 +126,24 @@ export const WORKER_RESOURCE_TAGS = ["gpu", "cpu", "network"] as const;
//
// A worker with only unknown tags therefore matches requirement-less work and
// nothing else — harmlessly over-general rather than silently excluded.
+//
+// THE ONE EXCEPTION: an "llm" worker is a bare model endpoint, not an instance
+// of this app, and it can serve ONLY a requirement naming one of its LLM
+// operations. It must never take a requirement-less lease — today's
+// transcription acquire passes no requirement, and an ollama endpoint handed a
+// transcription job cannot run it. So for llm workers the untagged-is-universal
+// rule inverts: no tags means "all three LLM operations", never "anything".
export function workerMatches(
worker: Pick<Worker, "kind" | "tags">,
requires?: readonly string[],
): boolean {
+ if (worker.kind === "llm") {
+ if (!requires || requires.length === 0) return false;
+ const serves: readonly string[] = worker.tags?.length
+ ? worker.tags
+ : LLM_OPERATION_TAGS;
+ return requires.some((r) => serves.includes(r));
+ }
if (!requires || requires.length === 0) return true;
const tags = worker.tags ?? [];
if (tags.length === 0) return true;
@@ -153,6 +199,17 @@ function sanitizeRemoteConfig(value: unknown): RemoteWorkerConfig {
return cfg;
}
+function sanitizeLlmConfig(value: unknown): LlmWorkerConfig {
+ const cfg: LlmWorkerConfig = { baseUrl: "" };
+ if (!value || typeof value !== "object") return cfg;
+ const r = value as Record<string, unknown>;
+ if (typeof r.baseUrl === "string") cfg.baseUrl = r.baseUrl.trim();
+ if (typeof r.slots === "number" && Number.isFinite(r.slots) && r.slots >= 1) {
+ cfg.slots = Math.min(64, Math.floor(r.slots));
+ }
+ return cfg;
+}
+
function slugify(input: string): string {
return input
.toLowerCase()
@@ -178,7 +235,8 @@ export function sanitizeWorkers(value: unknown): Worker[] {
value.forEach((raw, index) => {
if (!raw || typeof raw !== "object") return;
const r = raw as Record<string, unknown>;
- const kind: WorkerKind = r.kind === "remote" ? "remote" : "local";
+ const kind: WorkerKind =
+ r.kind === "remote" ? "remote" : r.kind === "llm" ? "llm" : "local";
const rawId = typeof r.id === "string" ? r.id.trim() : "";
const rawName = typeof r.name === "string" ? r.name.trim() : "";
const id = uniqueId(rawId || rawName || `worker-${index + 1}`, seenIds);
@@ -212,8 +270,10 @@ export function sanitizeWorkers(value: unknown): Worker[] {
? r.appId
: DEFAULT_TRANSCRIPTION_APP_ID;
worker.config = sanitizeWorkerConfig(r.config);
- } else {
+ } else if (kind === "remote") {
worker.remote = sanitizeRemoteConfig(r.remote);
+ } else {
+ worker.llm = sanitizeLlmConfig(r.llm);
}
out.push(worker);
});
@@ -224,11 +284,28 @@ export function sanitizeWorkers(value: unknown): Worker[] {
// at least one enabled worker, each enabled local worker has a resolvable binary
// and (if it surfaces customArgs) a valid template, each remote has a valid URL.
export function validateWorkers(workers: Worker[]): string | null {
- if (!workers.some((w) => w.enabled)) {
- return "At least one worker must be enabled";
+ if (!workers.some((w) => w.enabled && w.kind !== "llm")) {
+ // An llm endpoint cannot transcribe, so a list of nothing but llm workers
+ // would park every transcription forever while reading as "enabled".
+ return workers.some((w) => w.enabled)
+ ? "At least one transcription worker (local or remote) must be enabled — an LLM endpoint only serves digest/attribution calls"
+ : "At least one worker must be enabled";
}
for (const w of workers) {
const label = w.name || w.id;
+ if (w.kind === "llm") {
+ const url = w.llm?.baseUrl?.trim();
+ if (!url) return `LLM worker "${label}" needs a base URL`;
+ try {
+ const u = new URL(url);
+ if (u.protocol !== "http:" && u.protocol !== "https:") {
+ return `LLM worker "${label}" base URL must be http(s)`;
+ }
+ } catch {
+ return `LLM worker "${label}" has an invalid base URL`;
+ }
+ continue;
+ }
if (w.kind === "local") {
if (!w.appId || !TRANSCRIPTION_APPS[w.appId]) {
return `Worker "${label}" has an unknown engine`;
diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md
@@ -1,6 +1,8 @@
# Changelog
## [Unreleased]
+- **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.
- **A remote worker is now worth what the remote can actually do, not one slot.** The worker model's own header has promised since it shipped that workers "carry a priority and a slot count" — no such field existed, so a remote box running four workers of its own took one transcription at a time from here unless you hand-copied its card four times. A remote card now has **Slots**: set a number, or leave it blank and the remote is asked directly — its health endpoint has always returned its per-worker summary, and every caller read one boolean off it and threw the rest away. The pool expands one remote into that many independently-schedulable slots (each can be enabled, drained or degraded on its own on the Workers page), re-checks a blank-slots remote about once a minute with no timer and no waiting — a stale answer costs at most one slot for one minute — and shrinking the count retires the surplus through the same drain-then-drop path a removed worker takes.
- **A worker's tags now route work instead of being stored and ignored.** The worker model has carried a `tags` field since it shipped, documented as "stored but NOT consulted by the scheduler" — so a second machine could be *described* as CPU-only or diarization-only and the scheduler would hand it anything anyway. Tags are consulted now, under one written-down rule: a tag names an **operation** (`diarization`, `digest`, …) or a **resource** (`gpu`, `cpu`, `network`); work asks for what it needs; an untagged worker still takes anything, so an install that never touched tags routes exactly as it always has. The grant is matched rather than blind — a freed CPU slot now goes to the first waiter that can actually use it instead of sitting behind a GPU waiter at the head of the queue, which was a deadlock, not a slowdown. The settings form gained a Tags field per worker that warns (and only warns) about a tag the scheduler will never match, and the Workers page shows each worker's tags beside its state. Arming the speaker-work sweep was one unlabelled button meaning *every enabled operation, every channel* — and on this archive that includes speaker-names-from-the-transcript at **~1 model call per transcript chunk**, on the order of 194,000 calls. The one setting that would have bounded it, `sweepKinds`, has been honoured by the run since the sweep was written and settable by **nothing but a hand-edit of `settings.json`**; "diarization and names-from-the-audio only" had no expression anywhere in the product. The lane now opens on a console: tick which **operations** the run covers, each stating its own backlog *and what one unit of it costs*, and read **THE PLAN** underneath — the ordered channel itinerary the sweep will actually walk, built by the same counting the run does, with the position marker on the channel it is working now. Untick the expensive lane and the total falls, the rows re-sort, and the button re-labels itself. *Newest first, across all channels* stops being a phrase in a dropdown and becomes a consequence you can see.
- **The channel rows in the plan are the channel scope.** Rather than a second picker, **Choose channels** turns the itinerary into checkboxes — off by default, so the resting state is a clean list — and the button then says *Sweep 3 channels* instead of *Sweep every channel*. Scope, plan and progress are deliberately **one view rather than three**, because the alternative has a bug in it: the scope has to be written at the instant the sweep is armed. The boot hook re-launches a sweep from settings alone, so a scope saved separately (or saved and then not armed) comes back after a restart as the corpus-wide run it was meant to replace. The control that sets the scope is therefore the control that arms.
diff --git a/editor/app/settings/components/WorkersField.tsx b/editor/app/settings/components/WorkersField.tsx
@@ -42,6 +42,18 @@ function blankRemote(): Row {
};
}
+function blankLlm(): Row {
+ return {
+ _key: keyCounter++,
+ id: "",
+ name: "",
+ kind: "llm",
+ enabled: true,
+ priority: 0,
+ llm: { baseUrl: "" },
+ };
+}
+
// Duplicate a worker into a new slot: same config, fresh id, "(copy)" name.
function copyOf(row: Row): Row {
const { _key, id, name, ...rest } = row;
@@ -54,6 +66,7 @@ function copyOf(row: Row): Row {
name: name ? `${name} (copy)` : "",
config: row.config ? { ...row.config } : undefined,
remote: row.remote ? { ...row.remote } : undefined,
+ llm: row.llm ? { ...row.llm } : undefined,
};
}
@@ -104,6 +117,15 @@ export function WorkersField({ initial, apps, name, knownTags }: Props) {
),
);
}
+ function patchLlm(key: number, change: Partial<NonNullable<Worker["llm"]>>) {
+ update(
+ rows.map((r) =>
+ r._key === key
+ ? { ...r, llm: { baseUrl: "", ...r.llm, ...change } }
+ : r,
+ ),
+ );
+ }
function remove(key: number) {
update(rows.filter((r) => r._key !== key));
}
@@ -139,6 +161,7 @@ export function WorkersField({ initial, apps, name, knownTags }: Props) {
onPatch={(c) => patch(row._key, c)}
onPatchConfig={(c) => patchConfig(row._key, c)}
onPatchRemote={(c) => patchRemote(row._key, c)}
+ onPatchLlm={(c) => patchLlm(row._key, c)}
onRemove={() => remove(row._key)}
onCopy={() => copy(i)}
onMove={(dir) => move(i, dir)}
@@ -159,6 +182,13 @@ export function WorkersField({ initial, apps, name, knownTags }: Props) {
>
+ Remote worker
</button>
+ <button
+ type="button"
+ onClick={() => update([...rows, blankLlm()])}
+ className="px-2 py-1 rounded border border-border text-xs hover:bg-muted"
+ >
+ + LLM endpoint
+ </button>
</div>
</div>
);
@@ -173,6 +203,7 @@ function WorkerCard({
onPatch,
onPatchConfig,
onPatchRemote,
+ onPatchLlm,
onRemove,
onCopy,
onMove,
@@ -185,6 +216,7 @@ function WorkerCard({
onPatch: (c: Partial<Row>) => void;
onPatchConfig: (c: Partial<NonNullable<Worker["config"]>>) => void;
onPatchRemote: (c: Partial<NonNullable<Worker["remote"]>>) => void;
+ onPatchLlm: (c: Partial<NonNullable<Worker["llm"]>>) => void;
onRemove: () => void;
onCopy: () => void;
onMove: (dir: -1 | 1) => void;
@@ -193,6 +225,7 @@ function WorkerCard({
const app = apps.find((a) => a.id === row.appId);
const cfg = row.config ?? {};
const remote = row.remote ?? { baseUrl: "" };
+ const llmCfg = row.llm ?? { baseUrl: "" };
return (
<div
aria-label={`worker ${index + 1}`}
@@ -279,9 +312,11 @@ function WorkerCard({
<span className="text-xs text-muted-foreground">
{row.kind === "local"
? "Local · one slot"
- : remote.slots !== undefined
- ? `Remote · ${remote.slots} slot${remote.slots === 1 ? "" : "s"}`
- : "Remote · slots probed from the remote"}
+ : row.kind === "llm"
+ ? `LLM endpoint · ${llmCfg.slots ?? 1} slot${(llmCfg.slots ?? 1) === 1 ? "" : "s"}`
+ : remote.slots !== undefined
+ ? `Remote · ${remote.slots} slot${remote.slots === 1 ? "" : "s"}`
+ : "Remote · slots probed from the remote"}
</span>
<TagsField
tags={row.tags ?? []}
@@ -293,7 +328,32 @@ function WorkerCard({
/>
</div>
- {row.kind === "local" ? (
+ {row.kind === "llm" ? (
+ <div className="flex flex-col gap-2 border-t border-border pt-2">
+ <CardField
+ label="Base URL"
+ value={llmCfg.baseUrl}
+ onChange={(v) => onPatchLlm({ baseUrl: v })}
+ hint="A bare ollama endpoint (e.g. http://macbook.lan:11434) — the box runs nothing but `ollama serve`. It must serve the primary's exact model tag; the scheduler checks and skips it if not."
+ id={`${uid}-llmurl`}
+ />
+ <CardField
+ label="Slots"
+ type="number"
+ value={llmCfg.slots !== undefined ? String(llmCfg.slots) : ""}
+ onChange={(v) =>
+ onPatchLlm({ slots: v.trim() === "" ? undefined : Number(v) })
+ }
+ hint="Concurrent generations to send this endpoint. Blank = 1."
+ id={`${uid}-llmslots`}
+ />
+ <p className="text-xs text-muted-foreground">
+ Takes digest and speaker-attribution model calls only — never a
+ transcription. Tags narrow which of those operations it serves;
+ blank tags mean all of them.
+ </p>
+ </div>
+ ) : row.kind === "local" ? (
<div className="flex flex-col gap-2 border-t border-border pt-2">
<CardField
label="Binary"
diff --git a/editor/app/workers/components/WorkersView.tsx b/editor/app/workers/components/WorkersView.tsx
@@ -37,7 +37,7 @@ function useNow(): number | null {
export type WorkerView = {
id: string;
name: string;
- kind: "local" | "remote";
+ kind: "local" | "remote" | "llm";
appId?: string;
priority: number;
busy: boolean;
@@ -221,7 +221,12 @@ function WorkerCard({
inDefault: boolean;
}) {
const badge = stateBadge(w);
- const engine = w.kind === "remote" ? "remote" : (w.appId ?? "local");
+ const engine =
+ w.kind === "remote"
+ ? "remote"
+ : w.kind === "llm"
+ ? "llm endpoint"
+ : (w.appId ?? "local");
return (
<li
aria-label={`worker ${w.name}`}