Archilyzer · Source

archilyzer

Archilyzer
git clone https://archilyzer.pages.dev/source/archilyzer.git
Log | Files | Refs | README | LICENSE

commit de6327699740739953508658754a6f91392c5794
parent 30eb2e1577ccbb29b4b1313ab4f0ff9d49141419
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Mon, 24 Aug 2026 09:44:23 -0400

workers: tags route work — matched grants, a requirement on acquire, tryAcquire

Worker.tags stops being "stored but NOT consulted": the matching rule lives in
lib/workers.ts (untagged matches everything; tagged matches on intersection;
the requirement lists acceptable qualifications), the pool enforces it in
eligible/pickFree/hasEligibleWorker, and pump()'s blind head-of-queue grant
becomes a scan for the first waiter the freed slot can serve — without which a
GPU waiter at the head parks a CPU waiter forever. acquire() takes `requires`
and rejects a requirement no configured worker matches rather than parking on
it; tryAcquire() is the non-parking claim the fan-out dispatchers need (with a
kind filter so they can never grab a local slot). Tags surface in the settings
form (with an unknown-tag warning that never blocks a save) and on /workers.
Also: the degrade log message named MAX_WORKER_ATTEMPTS where the real
threshold is WORKER_DEGRADE_THRESHOLD.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

Diffstat:
Mcommon/controller/transcribeOne.ts | 7+++++--
Mcommon/jobs/workerPool.test.ts | 153+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Mcommon/jobs/workerPool.ts | 133++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------
Acommon/lib/workers.test.ts | 72++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/lib/workers.ts | 66+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---
Meditor/CHANGELOG.md | 2+-
Meditor/app/settings/components/SettingsForm.tsx | 12++++++++++--
Meditor/app/settings/components/WorkersField.tsx | 71++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Meditor/app/settings/page.tsx | 8++++++++
Meditor/app/workers/components/WorkersView.tsx | 10++++++++++
10 files changed, 505 insertions(+), 29 deletions(-)

diff --git a/common/controller/transcribeOne.ts b/common/controller/transcribeOne.ts @@ -5,7 +5,7 @@ import type { Paths } from "../lib/paths"; import type { Worker } from "../lib/workers"; import { getTranscriptionApp } from "../lib/transcriptionApps"; import { detectTranscriptFormat } from "../lib/whisper"; -import { getWorkerPool } from "../jobs/workerPool"; +import { WORKER_DEGRADE_THRESHOLD, getWorkerPool } from "../jobs/workerPool"; import type { TaskTracker } from "../jobs/taskHooks"; import { normalizeTranscript } from "./normalizeTranscript"; import { diarizeOneVideo } from "./diarizeOne"; @@ -396,8 +396,11 @@ export async function transcribeWithWorker( ); } } else if (pool.markFailure(worker.id)) { + // The degrade rule is the pool's WORKER_DEGRADE_THRESHOLD, not this + // loop's per-video attempt cap — the two happen to be equal today, + // and the message used to name the wrong one. log( - `Worker ${worker.id} auto-disabled after ${MAX_WORKER_ATTEMPTS} consecutive failures; re-enable it from the Workers page.`, + `Worker ${worker.id} auto-disabled after ${WORKER_DEGRADE_THRESHOLD} consecutive failures; re-enable it from the Workers page.`, ); } log( diff --git a/common/jobs/workerPool.test.ts b/common/jobs/workerPool.test.ts @@ -12,8 +12,15 @@ import { WorkerPool } from "./workerPool"; // seeds it with reconfigure(workers, {applyEnabled:true}), which avoids reading // settings.json so the scheduler is fully isolated. -function worker(id: string, priority = 0): Worker { - return { id, name: id, kind: "local", enabled: true, priority }; +function worker(id: string, priority = 0, tags?: string[]): Worker { + return { + id, + name: id, + kind: "local", + enabled: true, + priority, + ...(tags ? { tags } : {}), + }; } // A single-slot pool with its only worker already leased out (running), so every @@ -141,6 +148,148 @@ test("tier ordering: urgent jumps ahead of foreground, which jumps ahead of back assert.deepEqual(order, ["U", "F", "B"]); }); +// --- Capability matching (Worker.tags + acquire requires) --- + +test("an untagged worker satisfies any requirement", async () => { + const pool = new WorkerPool(); + pool.reconfigure([worker("w1")], { applyEnabled: true }); + const lease = await pool.acquire(undefined, { requires: ["diarization"] }); + assert.equal(lease.worker.id, "w1"); + lease.release(); +}); + +test("a tagged worker matches when the requirement intersects its tags", async () => { + const pool = new WorkerPool(); + pool.reconfigure([worker("cpu-box", 0, ["cpu", "diarization"])], { + applyEnabled: true, + }); + // Any one qualification suffices — the requirement lists acceptable ones. + const lease = await pool.acquire(undefined, { + requires: ["diarization", "gpu"], + }); + assert.equal(lease.worker.id, "cpu-box"); + lease.release(); +}); + +test("a requirement no configured worker matches rejects instead of parking forever", async () => { + const pool = new WorkerPool(); + pool.reconfigure([worker("cpu-box", 0, ["cpu"])], { applyEnabled: true }); + await assert.rejects( + pool.acquire(undefined, { requires: ["gpu"] }), + /no configured worker matches/, + ); +}); + +test("a freed slot skips a head waiter it cannot serve (no head-of-line deadlock)", async () => { + const pool = new WorkerPool(); + pool.reconfigure( + [worker("gpu-w", 0, ["gpu"]), worker("cpu-w", 1, ["cpu"])], + { applyEnabled: true }, + ); + const gpuLease = await pool.acquire(undefined, { requires: ["gpu"] }); + const cpuLease = await pool.acquire(undefined, { requires: ["cpu"] }); + const order: string[] = []; + // The GPU waiter parks FIRST — a blind head grant would hand it the freed + // CPU slot (which it cannot use) or park the CPU waiter behind it forever. + const pGpu = pool + .acquire(undefined, { requires: ["gpu"] }) + .then((l) => (order.push("gpu"), l)); + const pCpu = pool + .acquire(undefined, { requires: ["cpu"] }) + .then((l) => (order.push("cpu"), l)); + + cpuLease.release(); // frees the CPU slot → the CPU waiter, not the head + await pCpu; + assert.deepEqual(order, ["cpu"]); + + gpuLease.release(); + await pGpu; + assert.deepEqual(order, ["cpu", "gpu"]); +}); + +test("tier order is preserved within the waiters a slot can serve", async () => { + const pool = new WorkerPool(); + pool.reconfigure([worker("cpu-w", 0, ["cpu"])], { applyEnabled: true }); + const lease = await pool.acquire(undefined, { requires: ["cpu"] }); + const order: string[] = []; + const b = pool + .acquire(undefined, { tier: "background", requires: ["cpu"] }) + .then((l) => (order.push("B"), l)); + const f = pool + .acquire(undefined, { tier: "foreground", requires: ["cpu"] }) + .then((l) => (order.push("F"), l)); + + let lease_ = lease; + for (const p of [f, b]) { + lease_.release(); + lease_ = await p; + } + assert.deepEqual(order, ["F", "B"]); +}); + +test("tryAcquire claims a free matching slot and never parks", async () => { + const pool = new WorkerPool(); + pool.reconfigure([worker("cpu-w", 0, ["cpu"])], { applyEnabled: true }); + const lease = pool.tryAcquire(["cpu"]); + assert.ok(lease, "expected a lease on the free matching worker"); + assert.equal(lease!.worker.id, "cpu-w"); + // Busy now — a second claim yields null rather than parking. + assert.equal(pool.tryAcquire(["cpu"]), null); + // And a non-matching requirement never claims it. + lease!.release(); + assert.equal(pool.tryAcquire(["gpu"]), null); +}); + +test("tryAcquire's kind filter keeps a fan-out off local workers", async () => { + const pool = new WorkerPool(); + pool.reconfigure([worker("w1")], { applyEnabled: true }); + // The untagged local worker matches the requirement, but the claim is + // restricted to remote capacity — the fan-out must never grab a local slot + // out from under the transcription scheduler. + assert.equal(pool.tryAcquire(["digest"], { kind: "remote" }), null); + const lease = pool.tryAcquire(["digest"]); + assert.ok(lease, "unrestricted claim still matches the untagged local"); + lease!.release(); +}); + +test("a slot freed while an unmatchable waiter parks stays claimable by others", async () => { + const pool = new WorkerPool(); + pool.reconfigure( + [worker("gpu-w", 0, ["gpu"]), worker("cpu-w", 1, ["cpu"])], + { applyEnabled: true }, + ); + const gpuLease = await pool.acquire(undefined, { requires: ["gpu"] }); + const cpuLease = await pool.acquire(undefined, { requires: ["cpu"] }); + const pGpu = pool.acquire(undefined, { requires: ["gpu"] }); + // Freeing the CPU slot cannot serve the parked GPU waiter… + cpuLease.release(); + // …so it stays free, and a claim that does match it succeeds without + // starving the waiter (which the freed slot could never serve). + const claimed = pool.tryAcquire(["cpu"]); + assert.ok(claimed, "the free CPU slot should be claimable"); + claimed!.release(); + gpuLease.release(); + (await pGpu).release(); +}); + +test("hasEligibleWorker honours a requirement", () => { + const pool = new WorkerPool(); + pool.reconfigure([worker("cpu-w", 0, ["cpu"])], { applyEnabled: true }); + assert.equal(pool.hasEligibleWorker(), true); + assert.equal(pool.hasEligibleWorker(["cpu"]), true); + assert.equal(pool.hasEligibleWorker(["gpu"]), false); +}); + +test("summary() carries tags and keeps its shape for untagged workers", () => { + const pool = new WorkerPool(); + pool.reconfigure([worker("plain"), worker("tagged", 1, ["cpu"])], { + applyEnabled: true, + }); + const summary = pool.summary(); + assert.equal(summary.find((w) => w.id === "plain")?.tags, undefined); + assert.deepEqual(summary.find((w) => w.id === "tagged")?.tags, ["cpu"]); +}); + 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 @@ -24,7 +24,7 @@ // jumps ahead of queued auto work — the running auto unit isn't interrupted, it // drains, and the freed slot goes to the manual waiter. -import type { Worker } from "../lib/workers"; +import { workerMatches, type Worker, type WorkerKind } from "../lib/workers"; import { getSettings } from "../lib/settings"; import { getPaths } from "../lib/paths"; import { readWorkerDefaults } from "./workerDefaults"; @@ -55,6 +55,10 @@ export type WorkerSummary = { state: WorkerRuntimeState; degraded: boolean; enabled: boolean; // persisted intent (config.enabled) + // Capability tags, verbatim from the config. What the sweep console reads to + // say where an operation can run — through the same workerMatches rule the + // scheduler grants by, so the display and the dispatch cannot disagree. + tags?: string[]; // The engine's compute device, verbatim from the per-worker config ("cpu", // "cuda:0", …). UNDEFINED MEANS UNKNOWN, not CPU: it is the engine binary's own // default, which on parakeet-cli may well be the GPU. @@ -88,6 +92,10 @@ type Waiter = { // queue behind foreground (manual) ones for free slots; an urgent waiter would // jump ahead of both. Same compareTier ordering the registry's scheduler uses. tier: SchedulerTier; + // Capability requirement (see workerMatches). A freed slot is granted to the + // FIRST waiter it matches, not blindly to the head — pump() scans past a + // waiter whose requirement this slot cannot satisfy. + requires?: readonly string[]; }; export class WorkerPool { @@ -247,15 +255,21 @@ export class WorkerPool { return this.pausedSnapshot !== null; } - private eligible(entry: PoolEntry): boolean { - return entry.state === "enabled" && !entry.degraded && !entry.busy; + private eligible(entry: PoolEntry, requires?: readonly string[]): boolean { + return ( + entry.state === "enabled" && + !entry.degraded && + !entry.busy && + workerMatches(entry.config, requires) + ); } - // Highest-priority free entry: (priority asc, insertion order asc). - private pickFree(): PoolEntry | null { + // Highest-priority free entry matching the requirement: + // (priority asc, insertion order asc). + private pickFree(requires?: readonly string[]): PoolEntry | null { let best: PoolEntry | null = null; for (const entry of this.entries.values()) { - if (!this.eligible(entry)) continue; + if (!this.eligible(entry, requires)) continue; if (!best || entry.config.priority < best.config.priority) best = entry; } return best; @@ -278,14 +292,27 @@ export class WorkerPool { return { worker: entry.config, release }; } - // Try to satisfy parked waiters in FIFO order while free slots remain. + // Try to satisfy parked waiters while free slots remain. The queue is + // tier-ordered (see enqueueWaiter), and the grant is MATCHED, not blind: each + // pass grants the earliest waiter for which a matching slot is free. Scanning + // past a waiter whose requirement no free slot satisfies is what stops a GPU + // waiter at the head from parking a CPU waiter forever when a CPU slot frees — + // a deadlock, not a slowdown. Tier order is preserved WITHIN the set of + // waiters a given slot can serve, which is the only order that means anything. private pump(): void { - while (this.waiters.length > 0) { - const entry = this.pickFree(); - if (!entry) break; - const waiter = this.waiters.shift()!; - waiter.onAbort?.(); - waiter.resolve(this.grant(entry)); + for (;;) { + let granted = false; + for (let i = 0; i < this.waiters.length; i++) { + const waiter = this.waiters[i]; + const entry = this.pickFree(waiter.requires); + if (!entry) continue; + this.waiters.splice(i, 1); + waiter.onAbort?.(); + waiter.resolve(this.grant(entry)); + granted = true; + break; + } + if (!granted) return; } } @@ -313,22 +340,50 @@ export class WorkerPool { // cancel. // opts.tier sets the priority tier; the legacy opts.background is still // accepted and maps to the "background" tier (default "foreground"). + // opts.requires is a capability requirement (see workerMatches): the lease is + // granted only on a worker that matches it. A requirement NO configured worker + // matches rejects immediately — parking on it would never resolve, and the + // caller should treat the work as unrunnable rather than wedge. acquire( signal?: AbortSignal, - opts?: { tier?: SchedulerTier; background?: boolean }, + opts?: { + tier?: SchedulerTier; + background?: boolean; + requires?: readonly string[]; + }, ): Promise<Lease> { this.ensureInit(); if (signal?.aborted) { return Promise.reject(abortError()); } - const entry = this.pickFree(); + const requires = opts?.requires; + if (requires && requires.length > 0) { + let satisfiable = false; + for (const e of this.entries.values()) { + // Config-level, ignoring busy/disabled/degraded on purpose: a matching + // worker that is merely busy or switched off is a reason to park, not + // to give up. + if (workerMatches(e.config, requires)) { + satisfiable = true; + break; + } + } + if (!satisfiable) { + return Promise.reject( + new Error( + `no configured worker matches requirement [${requires.join(", ")}]`, + ), + ); + } + } + const entry = this.pickFree(requires); if (entry && this.waiters.length === 0) { return Promise.resolve(this.grant(entry)); } const tier: SchedulerTier = opts?.tier ?? (opts?.background ? "background" : "foreground"); return new Promise<Lease>((resolve, reject) => { - const waiter: Waiter = { resolve, reject, tier }; + const waiter: Waiter = { resolve, reject, tier, requires }; if (signal) { const onAbort = () => { const idx = this.waiters.indexOf(waiter); @@ -345,6 +400,36 @@ export class WorkerPool { }); } + // Non-parking claim: a matching free slot right now, or null. For dispatchers + // that already have a local way to run the work (the digest/backfill fan-out): + // a parked acquire inside one of their runPool slots would deadlock the batch, + // so they ask, take a lease when one is free, and otherwise run locally. + // + // opts.kind restricts the claim to one worker kind — the fan-out paths lease + // REMOTE capacity only ("llm" endpoints, "remote" unit executors) and must + // never grab a local slot out from under the transcription scheduler. + // + // 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 }, + ): 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 (!this.eligible(entry, requires)) continue; + if (!best || entry.config.priority < best.config.priority) best = entry; + } + if (!best) return null; + const claimed = best; + if (this.waiters.some((w) => workerMatches(claimed.config, w.requires))) { + return null; + } + return this.grant(claimed); + } + // --- 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 @@ -459,10 +544,19 @@ export class WorkerPool { // fullness — a full but enabled pool is working at capacity, not paused). Its // negation is what surfaces the "paused — waiting for a worker" state: a running // batch whose every worker is disabled/draining/degraded, so nothing can start. - hasEligibleWorker(): boolean { + // `requires` narrows the question to workers matching that capability + // requirement — "could a digest fan-out ever land anywhere?" is a different + // question from "can transcription run?". + hasEligibleWorker(requires?: readonly string[]): boolean { this.ensureInit(); for (const entry of this.entries.values()) { - if (entry.state === "enabled" && !entry.degraded) return true; + if ( + entry.state === "enabled" && + !entry.degraded && + workerMatches(entry.config, requires) + ) { + return true; + } } return false; } @@ -489,6 +583,9 @@ export class WorkerPool { state: e.state, degraded: e.degraded, enabled: e.config.enabled, + ...(e.config.tags && e.config.tags.length > 0 + ? { tags: [...e.config.tags] } + : {}), device: e.config.config?.device, })) .sort((a, b) => a.priority - b.priority || a.id.localeCompare(b.id)); diff --git a/common/lib/workers.test.ts b/common/lib/workers.test.ts @@ -0,0 +1,72 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { + WORKER_RESOURCE_TAGS, + sanitizeWorkers, + unknownWorkerTags, + workerMatches, + type Worker, +} from "./workers"; + +// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test lib/workers.test.ts +// +// Pins the capability-matching rule (workerMatches) and the tag sanitizer. The +// rule is enforced by the worker pool but DECLARED here, and the sweep console +// reads it too — so the pinned semantics are: no requirement matches everyone, +// an untagged worker matches everything, a tagged worker matches on a +// non-empty intersection. + +test("workerMatches: no requirement matches every worker", () => { + assert.equal(workerMatches({ kind: "local" }), true); + assert.equal(workerMatches({ kind: "local", tags: ["cpu"] }), true); + assert.equal(workerMatches({ kind: "remote", tags: ["gpu"] }, []), true); +}); + +test("workerMatches: an untagged worker matches any requirement", () => { + assert.equal(workerMatches({ kind: "local" }, ["diarization"]), true); + assert.equal(workerMatches({ kind: "remote", tags: [] }, ["gpu"]), true); +}); + +test("workerMatches: a tagged worker matches on intersection only", () => { + const w = { kind: "remote" as const, tags: ["cpu", "diarization"] }; + assert.equal(workerMatches(w, ["diarization"]), true); + // Any one acceptable qualification suffices. + assert.equal(workerMatches(w, ["gpu", "cpu"]), true); + assert.equal(workerMatches(w, ["gpu"]), false); + assert.equal(workerMatches(w, ["network"]), false); +}); + +test("sanitizeWorkers trims tags, drops empties, and de-duplicates", () => { + const [w] = sanitizeWorkers([ + { + name: "box", + kind: "local", + tags: [" cpu ", "", "cpu", "diarization", 7, " "], + }, + ]); + assert.deepEqual(w.tags, ["cpu", "diarization"]); +}); + +test("sanitizeWorkers omits tags entirely when nothing survives", () => { + const [w] = sanitizeWorkers([ + { name: "box", kind: "local", tags: ["", " "] }, + ]); + assert.equal(w.tags, undefined); +}); + +test("unknownWorkerTags names only tags outside the vocabulary", () => { + const workers: Worker[] = [ + { + id: "a", + name: "a", + kind: "local", + enabled: true, + priority: 0, + tags: ["cpu", "diarizatoin"], + }, + { id: "b", name: "b", kind: "local", enabled: true, priority: 1 }, + ]; + const known = ["diarization", ...WORKER_RESOURCE_TAGS]; + assert.deepEqual(unknownWorkerTags(workers, known), ["diarizatoin"]); + assert.deepEqual(unknownWorkerTags(workers, [...known, "diarizatoin"]), []); +}); diff --git a/common/lib/workers.ts b/common/lib/workers.ts @@ -48,8 +48,12 @@ export type Worker = { enabled: boolean; // Lower = preferred. Ties broken by array order in the scheduler. priority: number; - // Reserved for future capability routing (e.g. "only long audio to the GPU"). - // Stored but NOT consulted by the scheduler in v1. + // 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. tags?: string[]; // LOCAL: an instance of a TRANSCRIPTION_APPS entry + its per-worker config. @@ -60,6 +64,53 @@ export type Worker = { remote?: RemoteWorkerConfig; }; +// The contended-resource half of the tag vocabulary — BackfillLane.contendsFor's +// three values, restated here because this module must stay client-safe and the +// lane type lives in backfillKinds.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 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. +export function workerMatches( + worker: Pick<Worker, "kind" | "tags">, + requires?: readonly string[], +): boolean { + 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<string>(knownTags); + const unknown = new Set<string>(); + 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 { @@ -128,7 +179,16 @@ export function sanitizeWorkers(value: unknown): Worker[] { : index, }; if (Array.isArray(r.tags)) { - const tags = r.tags.filter((t): t is string => typeof t === "string"); + // 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") { diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -1,7 +1,7 @@ # Changelog ## [Unreleased] -- **A corpus sweep is a commitment, and you can finally read it before making it.** 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. +- **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. - **The digest sweep got the same console, and a scope it could not previously be given at all.** Its channel scope was reachable in the controller and unreachable from the application — the action that armed it took no arguments. It takes one now. It offers no operation list, because it *is* one operation and a permanently-ticked checkbox would imply a choice that does not exist. - **A scope naming an operation that no longer exists no longer arms a sweep that runs forever doing nothing.** Settings sanitation keeps an unknown operation id, and the resolver then matches nothing with it — so the sweep starts, reports itself armed, and holds at zero. Unknown ids are now dropped at the moment of arming, with the console saying which; a scope that names *only* unknown operations is refused outright rather than started. diff --git a/editor/app/settings/components/SettingsForm.tsx b/editor/app/settings/components/SettingsForm.tsx @@ -42,9 +42,12 @@ type Props = { // Built on the server: digestApps.ts reaches process.env and imports execa, so // it must never end up in the client bundle. Same reason as `apps`. digestApps: DigestAppDescriptor[]; + // Known worker-tag vocabulary (operation ids + resources), built on the + // server for the same bundle reason — the catalog lives in backfillKinds.ts. + workerTags?: string[]; }; -export function SettingsForm({ initial, apps, digestApps }: Props) { +export function SettingsForm({ initial, apps, digestApps, workerTags }: Props) { const [state, formAction] = useActionState<SaveResult | undefined, FormData>( saveSettingsAction, undefined, @@ -92,7 +95,12 @@ export function SettingsForm({ initial, apps, digestApps }: Props) { </a>{" "} page — handy for freeing one core/GPU while the rest keep going. </p> - <WorkersField initial={initial.workers} apps={apps} name="workersJson" /> + <WorkersField + initial={initial.workers} + apps={apps} + name="workersJson" + knownTags={workerTags} + /> </fieldset> <Field label="Cookies from browser" diff --git a/editor/app/settings/components/WorkersField.tsx b/editor/app/settings/components/WorkersField.tsx @@ -71,9 +71,14 @@ type Props = { initial: Worker[]; apps: TranscriptionAppDescriptor[]; name: string; + // The known tag vocabulary (operation ids + resource names), built on the + // server: the catalog lives in backfillKinds.ts, whose import graph must not + // reach the client bundle. Used only to WARN about a tag the scheduler will + // never match — an unknown tag saves fine. + knownTags?: string[]; }; -export function WorkersField({ initial, apps, name }: Props) { +export function WorkersField({ initial, apps, name, knownTags }: Props) { const [rows, setRows] = useState<Row[]>(() => withKeys(initial)); const defaultAppId = apps[0]?.id ?? "whisper-cpp"; @@ -130,6 +135,7 @@ export function WorkersField({ initial, apps, name }: Props) { index={i} total={rows.length} apps={apps} + knownTags={knownTags} onPatch={(c) => patch(row._key, c)} onPatchConfig={(c) => patchConfig(row._key, c)} onPatchRemote={(c) => patchRemote(row._key, c)} @@ -163,6 +169,7 @@ function WorkerCard({ index, total, apps, + knownTags, onPatch, onPatchConfig, onPatchRemote, @@ -174,6 +181,7 @@ function WorkerCard({ index: number; total: number; apps: TranscriptionAppDescriptor[]; + knownTags?: string[]; onPatch: (c: Partial<Row>) => void; onPatchConfig: (c: Partial<NonNullable<Worker["config"]>>) => void; onPatchRemote: (c: Partial<NonNullable<Worker["remote"]>>) => void; @@ -271,6 +279,14 @@ function WorkerCard({ <span className="text-xs text-muted-foreground"> {row.kind === "local" ? "Local · one slot" : "Remote · one slot"} </span> + <TagsField + tags={row.tags ?? []} + knownTags={knownTags} + index={index} + onChange={(tags) => + onPatch({ tags: tags.length > 0 ? tags : undefined }) + } + /> </div> {row.kind === "local" ? ( @@ -380,6 +396,59 @@ function WorkerCard({ ); } +// Comma-separated capability tags. Local text state so a trailing comma or +// space survives typing; the parsed list is pushed up on every change. Unknown +// tags WARN and still save — the scheduler tolerates them, it just never +// matches them. +function TagsField({ + tags, + knownTags, + index, + onChange, +}: { + tags: string[]; + knownTags?: string[]; + index: number; + onChange: (tags: string[]) => void; +}) { + const [text, setText] = useState(tags.join(", ")); + const unknown = knownTags ? tags.filter((t) => !knownTags.includes(t)) : []; + return ( + <label className="flex flex-col gap-1 text-xs flex-1 min-w-48"> + <span className="font-medium">Tags</span> + <input + type="text" + aria-label={`worker ${index + 1} tags`} + value={text} + placeholder="e.g. cpu, diarization" + onChange={(e) => { + setText(e.target.value); + onChange([ + ...new Set( + e.target.value + .split(",") + .map((s) => s.trim()) + .filter(Boolean), + ), + ]); + }} + className="rounded border border-border bg-card px-2 py-1 text-sm" + /> + <span className="text-muted-foreground"> + Capability routing: operation ids (e.g. <code>diarization</code>) or + resources (<code>gpu</code>, <code>cpu</code>, <code>network</code>), + comma-separated. Blank = this worker takes anything. + </span> + {unknown.length > 0 && ( + <span className="text-warning"> + Unknown tag{unknown.length > 1 ? "s" : ""}: {unknown.join(", ")} — + saved, but the scheduler will never match {unknown.length > 1 ? "them" : "it"}. + </span> + )} + </label> + ); +} + function CardField({ label, value, diff --git a/editor/app/settings/page.tsx b/editor/app/settings/page.tsx @@ -3,6 +3,8 @@ import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { getSettings } from "yt-dlp-transcript-common/lib/settings"; import { listTranscriptionApps } from "yt-dlp-transcript-common/lib/transcriptionApps"; import { listDigestApps } from "yt-dlp-transcript-common/lib/digestApps"; +import { operationCatalog } from "yt-dlp-transcript-common/lib/backfillKinds"; +import { WORKER_RESOURCE_TAGS } from "yt-dlp-transcript-common/lib/workers"; import { readXSessionStatus } from "yt-dlp-transcript-common/social/xSessionBroker"; import { SettingsForm } from "./components/SettingsForm"; import { XSessionSection } from "./components/XSessionSection"; @@ -47,6 +49,12 @@ export default async function SettingsPage() { initial={settings} apps={listTranscriptionApps()} digestApps={listDigestApps()} + workerTags={[ + ...new Set([ + ...operationCatalog().map((o) => o.id), + ...WORKER_RESOURCE_TAGS, + ]), + ]} /> </section> diff --git a/editor/app/workers/components/WorkersView.tsx b/editor/app/workers/components/WorkersView.tsx @@ -44,6 +44,8 @@ export type WorkerView = { state: WorkerRuntimeState; degraded: boolean; enabled: boolean; + // Capability tags (see lib/workers.ts workerMatches). Absent = takes anything. + tags?: string[]; // True when the in-flight transcription can be stopped into a partial result. canStopPartial: boolean; tasks: WorkerTask[]; @@ -234,6 +236,14 @@ function WorkerCard({ </span> <span className="text-xs text-muted-foreground">{engine}</span> <span className="text-xs text-muted-foreground">priority {w.priority}</span> + {w.tags && w.tags.length > 0 && ( + <span + className="text-xs text-muted-foreground" + title="Capability tags — this worker only takes work requiring one of these" + > + tags: {w.tags.join(", ")} + </span> + )} {inDefault && ( <span className="text-[11px] px-1.5 py-0.5 rounded-full border bg-info-soft text-info border-info/30"