// Transcription workers: named, independently-configured processing slots. // // A worker is one engine the scheduler can hand a video to. The three apps in // transcriptionApps.ts (whisper-cpp/chough/parakeet) are the *types* of a LOCAL // worker — a local worker is just an (appId, AppInstanceConfig) pair promoted to // a named, repeatable slot. A REMOTE worker delegates the transcription to the // same app running on another machine on the LAN (see common/controller/ // remoteTranscribe.ts). Workers carry a priority and a slot count; the global // scheduler (common/jobs/workerPool.ts) distributes each video — per task — to // the highest-priority free worker. // // Pure definitions + sanitizers only. Dependency direction is one-way: // settings.ts -> workers.ts -> transcriptionApps.ts. Keep it that way. import { type AppInstanceConfig, DEFAULT_TRANSCRIPTION_APP_ID, TRANSCRIPTION_APPS, getTranscriptionApp, validateTranscribeArgs, } from "./transcriptionApps"; import type { FieldDocs } from "./fieldDocs"; 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). // Each field is documented in LLM_WORKER_CONFIG_FIELD_DOCS below (rendered into SETTINGS.md). export type LlmWorkerConfig = { baseUrl: string; slots?: number; }; export const LLM_WORKER_CONFIG_FIELD_DOCS: FieldDocs = { baseUrl: "e.g. http://macbook.lan:11434", slots: "Concurrent generations to allow this endpoint. Defaults to 1 — one " + "model instance, one generation — unless the operator knows better.", }; // 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 // engine to use. // Each field is documented in REMOTE_WORKER_CONFIG_FIELD_DOCS below (rendered into SETTINGS.md). export type RemoteWorkerConfig = { baseUrl: string; token?: string; sharedFs?: boolean; slots?: number; }; export const REMOTE_WORKER_CONFIG_FIELD_DOCS: FieldDocs = { baseUrl: "e.g. http://gpu-box.lan:3011", token: "Outbound bearer token sent with every /api/worker request to this " + "remote. The accepting side validates against its own WORKER_TOKEN env," + " never this.", sharedFs: "When true the remote shares the transcripts mount, so we send " + "{channelSlug, videoId} instead of uploading the audio bytes.", slots: "How many units this remote takes in parallel. The pool expands one " + "remote config into this many independently-schedulable slot entries at" + " reconfigure time (the defaultWorkersFromApps trick, applied live). " + "Absent = probed from the remote's /api/worker/health (its enabled " + "worker count) — see controller/remoteCapacity.ts; 1 until the probe " + "answers.", }; // A worker is ONE processing slot — one transcription at a time. To run N in // parallel on the same engine, define N workers (the settings editor's "Copy" // button duplicates one). This makes each slot independently togglable on the // Workers page (free up one core/GPU while the rest keep going). // // The one exception is a REMOTE worker with a slot count (remote.slots, or the // probed capacity): the pool expands it into N slot entries itself, because // asking the operator to hand-copy a remote once per slot of a machine whose // slot count the machine already reports would be busywork. // Each field is documented in WORKER_FIELD_DOCS below (rendered into SETTINGS.md). export type Worker = { id: string; name: string; kind: WorkerKind; enabled: boolean; priority: number; tags?: string[]; appId?: string; config?: AppInstanceConfig; remote?: RemoteWorkerConfig; llm?: LlmWorkerConfig; }; export const WORKER_FIELD_DOCS: FieldDocs = { id: "Stable slug; used in settings, task ids, and logs.", name: "Human label shown in the UI.", kind: "\"local\" runs an app from TRANSCRIPTION_APPS on this machine (`appId` + `config`); \"remote\" delegates to another instance on the LAN (`remote`); \"llm\" is a bare ollama endpoint serving digest/attribution calls only (`llm`).", enabled: "Whether the scheduler may hand this slot work. Each worker is one slot, so parallelism is toggled per slot on the Workers page. Anything but an explicit `false` reads as enabled.", priority: "Lower = preferred. Ties broken by array order in the scheduler.", tags: "Capability routing. A tag is an OPERATION id from the backfill catalog" + " (\"diarization\", \"attribution-text\", …) or a contended RESOURCE " + "(WORKER_RESOURCE_TAGS). The scheduler consults them through " + "workerMatches below: an untagged worker takes anything, a tagged " + "worker takes only work whose requirement intersects its tags. Unknown " + "tags are tolerated (they match nothing and warn in the settings UI), " + "never fatal.", appId: "LOCAL: an instance of a TRANSCRIPTION_APPS entry + its per-worker " + "config.", config: "LOCAL: the per-worker engine config (binary, model, device, …) — an AppInstanceConfig, see `transcriptionApps.`.", remote: "REMOTE: how to reach the delegate instance.", llm: "LLM: how to reach the bare model endpoint.", }; // The contended-resource half of the tag vocabulary — Lane.contendsFor's // three values, restated here because this module must stay client-safe and the // lane type lives in operations.ts, whose import graph reaches controllers. // 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. // // - `requires` lists acceptable QUALIFICATIONS for one piece of work — the // operation id, or the resource its lane contends for. Any one suffices. // - No requirement means "today's transcription path": every worker matches, // so an existing install with no tags routes exactly as it always has. // - An UNTAGGED worker matches every requirement. Tags narrow a worker; their // absence must never quietly bench it. // - A TAGGED worker matches when tags ∩ requires ≠ ∅. // // 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, 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; return requires.some((r) => tags.includes(r)); } // Tags that name neither a known operation nor a known resource, for the // settings UI to WARN about — never to reject. The known-operation ids live in // the backfill catalog, which this client-safe module cannot import, so the // caller passes them in. export function unknownWorkerTags( workers: readonly Worker[], knownTags: readonly string[], ): string[] { const known = new Set(knownTags); const unknown = new Set(); for (const w of workers) { for (const t of w.tags ?? []) { if (!known.has(t)) unknown.add(t); } } return [...unknown]; } // Coerce a raw per-worker engine config into a clean AppInstanceConfig, dropping // unknown/ill-typed fields. Shared with settings.ts's sanitizeTranscriptionApps. export function sanitizeWorkerConfig(value: unknown): AppInstanceConfig { const cfg: AppInstanceConfig = {}; if (!value || typeof value !== "object") return cfg; const r = value as Record; if (typeof r.bin === "string") cfg.bin = r.bin; if (typeof r.model === "string") cfg.model = r.model; if (typeof r.remoteUrl === "string") cfg.remoteUrl = r.remoteUrl; if (typeof r.chunkSize === "number" && Number.isFinite(r.chunkSize)) { cfg.chunkSize = Math.floor(r.chunkSize); } if (Array.isArray(r.customArgs)) cfg.customArgs = r.customArgs.map(String); if (typeof r.device === "string") cfg.device = r.device; return cfg; } function sanitizeRemoteConfig(value: unknown): RemoteWorkerConfig { const cfg: RemoteWorkerConfig = { baseUrl: "" }; if (!value || typeof value !== "object") return cfg; const r = value as Record; if (typeof r.baseUrl === "string") cfg.baseUrl = r.baseUrl.trim(); if (typeof r.token === "string" && r.token.trim()) cfg.token = r.token.trim(); if (r.sharedFs === true) cfg.sharedFs = true; // Clamped to a sane range; absent (the default) means "probe the remote". if (typeof r.slots === "number" && Number.isFinite(r.slots) && r.slots >= 1) { cfg.slots = Math.min(64, Math.floor(r.slots)); } return cfg; } function sanitizeLlmConfig(value: unknown): LlmWorkerConfig { const cfg: LlmWorkerConfig = { baseUrl: "" }; if (!value || typeof value !== "object") return cfg; const r = value as Record; 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() .replace(/[^a-z0-9]+/g, "-") .replace(/^-+|-+$/g, ""); } function uniqueId(base: string, seen: Set): string { let id = slugify(base) || "worker"; if (!seen.has(id)) return id; let n = 2; while (seen.has(`${id}-${n}`)) n++; return `${id}-${n}`; } // Coerce a raw settings.workers value into a clean Worker[]. Generates stable // ids, fills defaults, and guarantees unique ids. Validation of "can this worker // actually run" (resolvable bin, valid baseUrl) is separate — see validateWorkers. export function sanitizeWorkers(value: unknown): Worker[] { if (!Array.isArray(value)) return []; const out: Worker[] = []; const seenIds = new Set(); value.forEach((raw, index) => { if (!raw || typeof raw !== "object") return; const r = raw as Record; 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); seenIds.add(id); const worker: Worker = { id, name: rawName || id, kind, enabled: r.enabled !== false, priority: typeof r.priority === "number" && Number.isFinite(r.priority) ? Math.floor(r.priority) : index, }; if (Array.isArray(r.tags)) { // Trimmed, de-duplicated, empty strings dropped: "" would intersect // nothing yet still count as "tagged", quietly benching the worker. const tags = [ ...new Set( r.tags .filter((t): t is string => typeof t === "string") .map((t) => t.trim()) .filter(Boolean), ), ]; if (tags.length > 0) worker.tags = tags; } if (kind === "local") { worker.appId = typeof r.appId === "string" && TRANSCRIPTION_APPS[r.appId] ? r.appId : DEFAULT_TRANSCRIPTION_APP_ID; worker.config = sanitizeWorkerConfig(r.config); } else if (kind === "remote") { worker.remote = sanitizeRemoteConfig(r.remote); } else { worker.llm = sanitizeLlmConfig(r.llm); } out.push(worker); }); return out; } // Validate a worker list for saving. Returns an error string or null. Enforces: // 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 && 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`; } if (!w.enabled) continue; const app = getTranscriptionApp(w.appId); const cfg = w.config ?? {}; const bin = cfg.bin?.trim() || app.defaultBin(); if (!bin) return `Worker "${label}" needs a binary`; if (app.fields.customArgs && cfg.customArgs && cfg.customArgs.length > 0) { const err = validateTranscribeArgs(cfg.customArgs); if (err) return err; } } else { const url = w.remote?.baseUrl?.trim(); if (!url) return `Remote worker "${label}" needs a base URL`; try { const u = new URL(url); if (u.protocol !== "http:" && u.protocol !== "https:") { return `Remote worker "${label}" base URL must be http(s)`; } } catch { return `Remote worker "${label}" has an invalid base URL`; } } } return null; } // Synthesize the default worker list for a settings.json that predates the worker // model. Since a worker is now one slot, the old global parallelTranscriptions is // preserved by emitting that many copies of the active app (priority 0), each a // separate enabled slot. Every other configured app becomes one disabled worker, // ready to enable. Keeps existing installs' parallelism while surfacing prior // configuration. export function defaultWorkersFromApps( activeAppId: string, appConfigs: Record, parallelSlots: number, ): Worker[] { const activeApp = getTranscriptionApp(activeAppId); const slots = Math.max(1, Math.floor(parallelSlots) || 1); const workers: Worker[] = []; const seen = new Set(); for (let i = 0; i < slots; i++) { const id = uniqueId(slots > 1 ? `${activeApp.id}-${i + 1}` : activeApp.id, seen); seen.add(id); workers.push({ id, name: slots > 1 ? `${activeApp.label} #${i + 1}` : activeApp.label, kind: "local", enabled: true, priority: 0, appId: activeApp.id, config: appConfigs[activeApp.id] ?? {}, }); } let priority = 1; for (const [appId, config] of Object.entries(appConfigs)) { if (appId === activeApp.id || !TRANSCRIPTION_APPS[appId]) continue; const id = uniqueId(appId, seen); seen.add(id); workers.push({ id, name: TRANSCRIPTION_APPS[appId].label, kind: "local", enabled: false, priority: priority++, appId, config, }); } return workers; }