commit 0eeb9b078b28318840251e41be196246d047a51a
parent 8aa041621e9e4aca45b3ec85311ba994105dcc76
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Sat, 22 Aug 2026 01:10:18 -0400
arbiter: one dispatcher for every operation a rule names
The step the unified-operations model exists for. backfillSweep's own
header says it is a clone of digestSweep and that "the third would be
written for the same reason". There is no third: an operation joins
this loop by being named on a leaf.
common/controller/arbiter.ts holds five properties, each paid for once
already elsewhere:
1. It runs on queueKey "". A long-lived job that WAITS on jobs
needing the real keys must not hold one, or it deadlocks against
its own work.
2. It reserves by `lane.queueKey` and never invents one — which is
what keeps digest and diarization overlapping instead of taking
turns, as backfill.spec.ts pins.
3. One job per channel, never per video. 77,000 jobs would evict the
registry's 100 records.
4. It re-plans every pass and holds no cursor.
5. It refuses to run beside an armed sweep, and says which one.
planArbiterUnits is pure and separated deliberately: it is where the
tree's PRIORITY becomes a DISPATCH ORDER, which is the part worth
pinning (7 tests). Leaf order is priority order and a leaf's ids are
already date-sorted, so under `order: "newest"` pending[leaf][0]'s
channel IS the freshest channel — cross-channel ordering is INDUCED
rather than implemented, which is what lets Stage 2's plan sort retire
with the sweeps.
CLOSES A HAZARD STAGE 3 LEFT OPEN. Operation leaves and bucket leaves
can share a tree — an operator may reasonably want "diarize this
channel before you transcribe anything" in one priority list — but
they must not share a DISPATCHER, or the auto-transcribe runner picks
a digest candidate off an operation leaf and hands it to whisper. Both
runner paths now retain bucket leaves only (retainLeaves), zeroing
rather than removing so selectNextWork's existing empty-leaf path
handles it.
NOT PERSISTED across a restart, and that is the one place it differs
from the sweeps: it dispatches four pipelines at once and has not
driven a full corpus pass, which is exactly what ARCHILYZER_IDLE_BOOT
exists to protect operators from. It becomes persistent at step 6.
Also fixes a second latent race in "UI: build a policy": saving an
enabled policy starts the runner itself, so clicking Start was a race
against the 3s poll rather than an assertion. It now clicks only if
still offered and asserts the runner is RUNNING plus the work
completes — strictly more than before.
759 common tests (7 new), 70/70 e2e across seven specs, build clean.
plans/unified-operations-model.md step 3 marked done.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Diffstat:
13 files changed, 1043 insertions(+), 14 deletions(-)
diff --git a/common/controller/arbiter.test.ts b/common/controller/arbiter.test.ts
@@ -0,0 +1,177 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import type {
+ AutoQueueGroup,
+ ChannelWork,
+} from "../jobs/autoQueuePolicy";
+import type { BackfillLane } from "../lib/backfillKinds";
+import { planArbiterUnits, type ArbiterUnit } from "./arbiter";
+
+// Run with: node_modules/.bin/tsx --test common/controller/arbiter.test.ts
+//
+// planArbiterUnits is the whole of the arbiter worth pinning: it is where the
+// policy tree's PRIORITY becomes a DISPATCH ORDER. The loop around it is
+// job-starting and sleeping, which an integration run covers.
+
+const GPU: BackfillLane = { queueKey: "digest:local", contendsFor: "gpu" };
+const CPU: BackfillLane = { queueKey: "backfill", contendsFor: "cpu" };
+
+const LANES: Record<string, BackfillLane> = {
+ digest: GPU,
+ diarization: CPU,
+ "attribution-text": CPU,
+};
+const laneOf = (op: string): BackfillLane | null => LANES[op] ?? null;
+
+function channels(): ChannelWork[] {
+ return [
+ {
+ slug: "alpha",
+ platform: "youtube",
+ buckets: { downloadedNoTranscript: ["a-bucket"] },
+ operations: { digest: ["a1", "a2"], diarization: ["a1"] },
+ },
+ {
+ slug: "beta",
+ platform: "youtube",
+ buckets: {},
+ operations: { digest: ["b1"], diarization: [] },
+ },
+ ];
+}
+
+const summarise = (units: ArbiterUnit[]): string[] =>
+ units.map((u) => `${u.operation}/${u.channelSlug}:${u.ids.join(",")}`);
+
+test("an operation leaf becomes one unit per channel", () => {
+ // ONE JOB PER CHANNEL, never per video: 77,000 jobs would evict the registry's
+ // 100 records, including the history of the run that made them.
+ const root: AutoQueueGroup = {
+ id: "root",
+ mode: "strict",
+ children: [{ id: "dig", match: { type: "all", operation: "digest" } }],
+ };
+ assert.deepEqual(summarise(planArbiterUnits([root], channels(), laneOf)), [
+ "digest/alpha:a1,a2",
+ "digest/beta:b1",
+ ]);
+});
+
+test("bucket leaves are ignored — they belong to the runners", () => {
+ // The hazard this closes: a bucket leaf and an operation leaf can share a
+ // tree, and without the split the arbiter would start a digest job for a
+ // video the auto-transcribe runner had claimed for whisper.
+ const root: AutoQueueGroup = {
+ id: "root",
+ mode: "strict",
+ children: [
+ { id: "buckets", match: { type: "all" } },
+ { id: "retries", match: { type: "all", bucket: "failedListed" } },
+ ],
+ };
+ assert.deepEqual(planArbiterUnits([root], channels(), laneOf), []);
+});
+
+test("leaf order is dispatch order, and a channel rule outranks the catch-all", () => {
+ const root: AutoQueueGroup = {
+ id: "root",
+ mode: "strict",
+ children: [
+ {
+ id: "beta-first",
+ match: { type: "channel", value: "beta", operation: "digest" },
+ },
+ { id: "rest", match: { type: "all", operation: "digest" } },
+ ],
+ };
+ // beta is claimed by the higher-priority leaf and does NOT reappear under the
+ // catch-all — first-match-wins, exactly as it is for buckets.
+ assert.deepEqual(summarise(planArbiterUnits([root], channels(), laneOf)), [
+ "digest/beta:b1",
+ "digest/alpha:a1,a2",
+ ]);
+});
+
+test("the same video is dispatched for two DIFFERENT operations", () => {
+ // The reason the claim key carries the operation. alpha/a1 is pending for
+ // both digest and diarization — different work on the same input — and an
+ // id-keyed claim set would silently drop one of them.
+ const root: AutoQueueGroup = {
+ id: "root",
+ mode: "strict",
+ children: [
+ { id: "dia", match: { type: "all", operation: "diarization" } },
+ { id: "dig", match: { type: "all", operation: "digest" } },
+ ],
+ };
+ const units = planArbiterUnits([root], channels(), laneOf);
+ assert.deepEqual(summarise(units), [
+ "diarization/alpha:a1",
+ "digest/alpha:a1,a2",
+ "digest/beta:b1",
+ ]);
+ // And they land on DIFFERENT queue keys, which is the only concurrency
+ // mechanism this system has. Sharing one would serialize the GPU lane behind
+ // the CPU one.
+ assert.notEqual(units[0].lane.queueKey, units[1].lane.queueKey);
+});
+
+test("an operation with no lane is not dispatched", () => {
+ // A tree may name a kind that has since been renamed, or one an operator
+ // later switched off. Neither is an error; both are simply not dispatched.
+ const root: AutoQueueGroup = {
+ id: "root",
+ mode: "strict",
+ children: [{ id: "gone", match: { type: "all", operation: "sortformer" } }],
+ };
+ const withIds: ChannelWork[] = [
+ {
+ slug: "alpha",
+ platform: "youtube",
+ buckets: {},
+ operations: { sortformer: ["a1"] },
+ },
+ ];
+ assert.deepEqual(planArbiterUnits([root], withIds, laneOf), []);
+});
+
+test("an empty selection produces no unit, and no empty job", () => {
+ const root: AutoQueueGroup = {
+ id: "root",
+ mode: "strict",
+ children: [
+ {
+ id: "dia-beta",
+ match: { type: "channel", value: "beta", operation: "diarization" },
+ },
+ ],
+ };
+ // beta has an entry for diarization, and it is empty. A unit here would start
+ // a job that classifies a whole channel to discover it has nothing to do.
+ assert.deepEqual(planArbiterUnits([root], channels(), laneOf), []);
+});
+
+test("both trees are read, in order", () => {
+ // An operator may express operation rules in either runner's tree — the
+ // arbiter reads the priority list they actually edit rather than demanding a
+ // third one.
+ const transcription: AutoQueueGroup = {
+ id: "t",
+ mode: "strict",
+ children: [
+ {
+ id: "t-dia",
+ match: { type: "channel", value: "alpha", operation: "diarization" },
+ },
+ ],
+ };
+ const download: AutoQueueGroup = {
+ id: "d",
+ mode: "strict",
+ children: [{ id: "d-dig", match: { type: "all", operation: "digest" } }],
+ };
+ assert.deepEqual(
+ summarise(planArbiterUnits([transcription, download], channels(), laneOf)),
+ ["diarization/alpha:a1", "digest/alpha:a1,a2", "digest/beta:b1"],
+ );
+});
diff --git a/common/controller/arbiter.ts b/common/controller/arbiter.ts
@@ -0,0 +1,418 @@
+// THE ARBITER: one dispatcher for every operation the policy tree names.
+//
+// This is the step the whole unified-operations model exists for.
+// `backfillSweep.ts` says in its own header that it is a clone of
+// `digestSweep.ts`, and that "the second implementation was written because the
+// first one's shape was not reusable, and the third would be written for the
+// same reason". There is no third. An operation joins this loop by being named
+// on a leaf.
+//
+// WHAT IT IS NOT. It is not a new runner per operation, it does not interleave
+// individual videos across channels, and it does not invent a queue key. It
+// reads the tree the operator already edits, groups the selected work by the
+// lane each operation DECLARES, and starts the per-channel jobs that already
+// exist.
+//
+// Five properties it must hold, each paid for once already elsewhere:
+//
+// 1. IT RUNS ON queueKey "". A long-lived job that WAITS on jobs needing the
+// real keys must not hold one itself, or it deadlocks against the work it
+// is waiting for. "" bypasses serialization entirely (registry.ts).
+// 2. IT RESERVES BY `lane.queueKey`, never by a key of its own. That is what
+// keeps the pinned separation test passing — digest and diarization hold
+// different keys precisely so they overlap instead of taking turns, and an
+// arbiter that grouped them under one key would silently serialize the GPU
+// lane behind a CPU one.
+// 3. ONE JOB PER CHANNEL, never per video. The registry keeps 100 records;
+// 77,000 jobs evict the history of the run that made them.
+// 4. IT RE-PLANS EVERY PASS, holding no cursor. A transcript that finished, a
+// settings change, a channel added — all of it is simply visible next pass.
+// 5. IT REFUSES TO RUN BESIDE AN ARMED SWEEP. Two dispatchers on one lane
+// would start the same channel twice, and the second job would spend its
+// life re-deriving work the first one had already taken.
+
+import type { Paths } from "../lib/paths";
+import { getPaths } from "../lib/paths";
+import { getSettings } from "../lib/settings";
+import { getRegistry } from "../jobs/registry";
+import { runManagedFunction } from "../jobs/streamCommand";
+import { drainStream } from "../jobs/drainStream";
+import {
+ buildPendingByLeaf,
+ flattenLeaves,
+ retainLeaves,
+ type AutoQueueGroup,
+ type ChannelWork,
+} from "../jobs/autoQueuePolicy";
+import {
+ DIGEST_KIND_ID,
+ getBackfillKind,
+ allBackfillKinds,
+ operationLabel,
+ type BackfillLane,
+} from "../lib/backfillKinds";
+import { digestLaneFor } from "../lib/backfillKinds";
+import { runDigestChannelJob } from "./digestSweep";
+import { runBackfillChannelJob } from "./backfillSweep";
+import {
+ buildArbiterChannelWork,
+ listArbiterChannelMeta,
+ type ArbiterChannelMeta,
+} from "./arbiterWork";
+
+export const ARBITER_KIND = "operations-arbiter";
+
+// How long the loop waits before re-planning when a pass moved nothing.
+// Generous for the same reason the sweeps' is: nothing here is latency
+// sensitive, and re-planning reads every channel's snapshot.
+const IDLE_POLL_MS = 60_000;
+
+// Consecutive passes that dispatched nothing before the arbiter gives up.
+// Three, not one: a single barren pass is legitimate (everything remaining
+// failed a guard this time and may not next time), but three in a row is a
+// broken engine or a corpus it cannot make progress on, and both want a human.
+const MAX_BARREN_PASSES = 3;
+
+type ArbiterLive = { jobId: string };
+type ArbiterSingleton = { live: ArbiterLive | null };
+
+declare global {
+ // eslint-disable-next-line no-var
+ var __yttArbiter__: ArbiterSingleton | undefined;
+}
+
+function getSingleton(): ArbiterSingleton {
+ if (!globalThis.__yttArbiter__) globalThis.__yttArbiter__ = { live: null };
+ return globalThis.__yttArbiter__;
+}
+
+export function getArbiterJobId(): string | null {
+ const live = getSingleton().live;
+ if (!live) return null;
+ return getRegistry().get(live.jobId)?.status === "running"
+ ? live.jobId
+ : null;
+}
+
+// One unit of dispatchable work: an operation, a channel, and the ids the tree
+// selected for it. The ids are a SELECTION, not a work list — every batch
+// re-derives eligibility from disk on every pull, so a stale id costs one
+// classification, never a wrong run.
+export type ArbiterUnit = {
+ operation: string;
+ channelSlug: string;
+ ids: string[];
+ lane: BackfillLane;
+ // The leaf that claimed it, for the log. An operator asking "why that
+ // channel" gets the rule back rather than a shrug.
+ leafId: string;
+};
+
+// The lane an operation runs on, or null when the catalog does not know it.
+//
+// Digest is asked through digestLaneFor because its lane depends on which
+// ENGINE is configured — the local lane contends for the GPU and yields, the
+// metered one contends for nothing local and must not. Every other operation's
+// lane is a fixed declaration.
+export function laneForOperation(operation: string): BackfillLane | null {
+ if (operation === DIGEST_KIND_ID) {
+ const settings = getSettings();
+ return digestLaneFor(
+ settings.digest.remoteEnabled ? "remote-api" : "local-gpu",
+ );
+ }
+ return getBackfillKind(operation)?.lane ?? null;
+}
+
+// Turn the policy trees into dispatchable units, in priority order.
+//
+// PURE, and separated from the loop for exactly that reason: this is where the
+// tree's priority becomes a dispatch order, and it is the part worth pinning
+// with a test rather than an integration run.
+//
+// Leaf order IS priority order (pre-order over the tree), and a leaf's ids are
+// already sorted by the runner's `order` setting — so under "newest",
+// `pending[leaf][0]`'s channel is the freshest channel, and cross-channel
+// ordering is INDUCED rather than implemented. That is what lets Stage 2's plan
+// sort retire with the sweeps.
+export function planArbiterUnits(
+ roots: ReadonlyArray<AutoQueueGroup>,
+ channels: ReadonlyArray<ChannelWork>,
+ laneOf: (operation: string) => BackfillLane | null,
+): ArbiterUnit[] {
+ const units: ArbiterUnit[] = [];
+ for (const root of roots) {
+ // `defaultBuckets` is empty on purpose: this pass wants operation leaves
+ // only, and a bucket leaf with no explicit bucket would otherwise claim ids
+ // out of a union it is about to be filtered out of anyway. Claiming still
+ // dedups correctly because the claim key carries the operation.
+ const pending = retainLeaves(
+ buildPendingByLeaf(root, channels, []),
+ root,
+ "operations",
+ );
+ for (const leaf of flattenLeaves(root)) {
+ const operation = leaf.match.operation;
+ if (!operation) continue;
+ const ids = pending[leaf.id] ?? [];
+ if (ids.length === 0) continue;
+ const lane = laneOf(operation);
+ // An operation the catalog does not know, or one whose feature is off,
+ // has no lane — so it is not dispatched. Silently skipping is right:
+ // a snapshot may name a kind that has since been renamed, and a tree may
+ // name one an operator later switched off.
+ if (!lane) continue;
+ // Group by channel, PRESERVING the order the ids arrived in, so the first
+ // channel out is the one holding the highest-priority video.
+ const byChannel = new Map<string, string[]>();
+ for (const id of ids) {
+ const slug = channelOf(channels, operation, id);
+ if (!slug) continue;
+ const list = byChannel.get(slug);
+ if (list) list.push(id);
+ else byChannel.set(slug, [id]);
+ }
+ for (const [channelSlug, channelIds] of byChannel) {
+ units.push({
+ operation,
+ channelSlug,
+ ids: channelIds,
+ lane,
+ leafId: leaf.id,
+ });
+ }
+ }
+ }
+ return units;
+}
+
+function channelOf(
+ channels: ReadonlyArray<ChannelWork>,
+ operation: string,
+ id: string,
+): string | null {
+ for (const ch of channels) {
+ if (ch.operations?.[operation]?.includes(id)) return ch.slug;
+ }
+ return null;
+}
+
+// Why the arbiter will not start, or null when it will.
+//
+// Stated as a function so the console can ask the same question the starter
+// asks, and get the same sentence back — an operator who clicks Start and gets
+// nothing deserves the reason, not a silent no-op.
+export function arbiterBlockedReason(): string | null {
+ const settings = getSettings();
+ if (settings.digest.sweepEnabled) {
+ return "The digest sweep is armed. Two dispatchers on one lane would start the same channel twice — stop the sweep first.";
+ }
+ if (settings.backfill.sweepEnabled) {
+ return "The backfill sweep is armed. Two dispatchers on one lane would start the same channel twice — stop the sweep first.";
+ }
+ return null;
+}
+
+async function runOneUnit(
+ paths: Paths,
+ unit: ArbiterUnit,
+ onLog: (msg: string) => void,
+): Promise<boolean> {
+ onLog(
+ `→ ${unit.channelSlug} · ${operationLabel(unit.operation)}: ` +
+ `${unit.ids.length.toLocaleString()} selected (rule ${unit.leafId}, queue ${unit.lane.queueKey}).`,
+ );
+ // THE JOBS THAT ALREADY EXIST. Same kind, same queue key and same progress
+ // metric a hand-clicked run produces, so an arbiter-driven pass is
+ // inspectable with the tools that already exist rather than being an opaque
+ // mega-job — and so this loop is a scheduler rather than a second runtime.
+ const outcome =
+ unit.operation === DIGEST_KIND_ID
+ ? await runDigestChannelJob({
+ paths,
+ channelSlug: unit.channelSlug,
+ lane: getSettings().digest.remoteEnabled ? "remote" : "local",
+ // Behind anything an operator clicks by hand.
+ background: true,
+ })
+ : await runBackfillChannelJob({
+ paths,
+ channelSlug: unit.channelSlug,
+ kindIds: [unit.operation],
+ // RESERVE BY THE LANE'S OWN KEY, never invent one. This is what keeps
+ // digest and diarization overlapping instead of taking turns, which
+ // backfill.spec.ts pins.
+ queueKey: unit.lane.queueKey,
+ background: true,
+ });
+ if (!outcome.ok) {
+ // One channel failing to START is not the arbiter failing.
+ onLog(`!! ${unit.channelSlug}: ${outcome.error ?? "failed to start"}`);
+ return false;
+ }
+ return true;
+}
+
+async function runArbiterLoop(
+ paths: Paths,
+ meta: ReadonlyArray<ArbiterChannelMeta>,
+ onLog: (msg: string) => void,
+ signal: AbortSignal,
+ drainSignal: AbortSignal,
+): Promise<void> {
+ let pass = 0;
+ let barrenPasses = 0;
+
+ for (;;) {
+ if (signal.aborted || drainSignal.aborted) return;
+ const blocked = arbiterBlockedReason();
+ if (blocked) {
+ // Cleanly, not by cancellation: the job ends "done" and the reason is the
+ // last thing in its log.
+ onLog(`Stopping: ${blocked}`);
+ return;
+ }
+
+ pass++;
+ const settings = getSettings();
+ const roots = [
+ settings.autoQueue.transcription.root,
+ settings.autoQueue.download.root,
+ ];
+ // Only the operations that are actually SWITCHED ON get projected. A tree
+ // may name one an operator later disabled, and projecting it would put tens
+ // of thousands of ids behind a feature that is off.
+ const enabled = new Set([
+ ...allBackfillKinds(settings).map((k) => k.id),
+ DIGEST_KIND_ID,
+ ]);
+ const operations = new Set<string>();
+ for (const root of roots) {
+ for (const leaf of flattenLeaves(root)) {
+ if (leaf.match.operation && enabled.has(leaf.match.operation)) {
+ operations.add(leaf.match.operation);
+ }
+ }
+ }
+ if (operations.size === 0) {
+ onLog(
+ "No rule names an operation this arbiter dispatches — nothing to do. " +
+ "Point a rule at an operation on /auto-queue, or stop the arbiter.",
+ );
+ return;
+ }
+
+ const channels = await buildArbiterChannelWork(paths, meta, [...operations]);
+ const units = planArbiterUnits(roots, channels, laneForOperation);
+ if (units.length === 0) {
+ onLog(`Pass ${pass}: nothing selected by any rule.`);
+ await sleep(IDLE_POLL_MS, signal);
+ continue;
+ }
+
+ // ONE UNIT PER LANE PER PASS, run concurrently. Distinct queue keys are the
+ // only concurrency mechanism this system has, so a lane gets exactly one
+ // job in flight and different lanes genuinely overlap. Taking the FIRST
+ // unit on each lane is what makes the tree's priority the dispatch order.
+ const firstByLane = new Map<string, ArbiterUnit>();
+ for (const unit of units) {
+ if (!firstByLane.has(unit.lane.queueKey)) {
+ firstByLane.set(unit.lane.queueKey, unit);
+ }
+ }
+ onLog(
+ `Pass ${pass}: ${units.length} unit(s) selected across ` +
+ `${new Set(units.map((u) => u.channelSlug)).size} channel(s); ` +
+ `dispatching ${firstByLane.size} lane(s).`,
+ );
+
+ const results = await Promise.all(
+ [...firstByLane.values()].map((unit) => runOneUnit(paths, unit, onLog)),
+ );
+ const didWork = results.some(Boolean);
+
+ barrenPasses = didWork ? 0 : barrenPasses + 1;
+ if (barrenPasses >= MAX_BARREN_PASSES) {
+ onLog(
+ `Pass ${pass} dispatched nothing for ${MAX_BARREN_PASSES} passes in a row. ` +
+ "Stopping rather than spinning — check the engines and the job logs, then start it again.",
+ );
+ return;
+ }
+ if (!didWork) {
+ onLog(
+ `Pass ${pass} dispatched nothing; waiting ${IDLE_POLL_MS / 1000}s before re-planning.`,
+ );
+ await sleep(IDLE_POLL_MS, signal);
+ }
+ }
+}
+
+function sleep(ms: number, signal: AbortSignal): Promise<void> {
+ return new Promise((resolve) => {
+ const t = setTimeout(resolve, ms);
+ signal.addEventListener(
+ "abort",
+ () => {
+ clearTimeout(t);
+ resolve();
+ },
+ { once: true },
+ );
+ });
+}
+
+export type StartArbiterResult = {
+ jobId: string | null;
+ error?: string;
+};
+
+// Start the arbiter. Returns the existing job id if one is already up, so a
+// double-click is a no-op rather than a second dispatcher.
+//
+// NOT PERSISTED, deliberately, and this is the one place it differs from the
+// sweeps. `sweepEnabled` survives a restart because a multi-week sweep that
+// forgot itself on a reboot would be worse than one that had to be re-armed.
+// The arbiter is new and it dispatches four pipelines at once; making it
+// resume unattended before it has driven a full corpus pass is exactly the
+// thing ARCHILYZER_IDLE_BOOT exists to protect operators from. It becomes
+// persistent when the sweeps retire (step 6), not before.
+export async function startArbiter(
+ paths: Paths = getPaths(),
+): Promise<StartArbiterResult> {
+ const running = getArbiterJobId();
+ if (running) return { jobId: running };
+
+ const blocked = arbiterBlockedReason();
+ if (blocked) return { jobId: null, error: blocked };
+
+ const meta = await listArbiterChannelMeta(paths);
+
+ const result = await runManagedFunction({
+ kind: ARBITER_KIND,
+ // "" — see property 1 at the top of this file. It waits on jobs that need
+ // the real keys, so holding one would deadlock against its own work.
+ queueKey: "",
+ paths,
+ fn: async (onLog, signal, _setProgress, ctx) => {
+ onLog(
+ "Arbiter up. It dispatches every operation a rule names, one job per " +
+ "channel, on each operation's own declared lane.",
+ );
+ await runArbiterLoop(paths, meta, onLog, signal, ctx.drainSignal);
+ onLog("Arbiter stopped.");
+ },
+ });
+ if (!result.ok) return { jobId: null, error: result.error };
+ getSingleton().live = { jobId: result.jobId };
+ // Deliberately NOT awaited: this is a long-lived loop, and the caller is a
+ // button.
+ void drainStream(result.stream).catch(() => {});
+ return { jobId: result.jobId };
+}
+
+export function stopArbiter(): void {
+ const id = getArbiterJobId();
+ if (id) getRegistry().cancel(id);
+ getSingleton().live = null;
+}
diff --git a/common/controller/arbiterWork.ts b/common/controller/arbiterWork.ts
@@ -0,0 +1,67 @@
+import type { Paths } from "../lib/paths";
+import { mapConcurrent } from "../lib/concurrency";
+import { listChannelConfigs, readChannelSnapshot } from "./channels";
+import { detectPlatform } from "../lib/platform";
+import type { ChannelWork } from "../jobs/autoQueuePolicy";
+import type { Platform } from "../lib/platform";
+
+// The arbiter's snapshot projection: operations only.
+//
+// SEPARATE from autoRunner's buildChannelWork, and not a refactor of it, for a
+// reason that is about scope rather than tidiness: that function projects the
+// BUCKETS a runner kind draws from, keyed by the runner's kind, and it also
+// builds the owner map the recency index needs. The arbiter needs neither — it
+// wants `snapshot.backfill[op].ids` for a set of operations and nothing else —
+// and widening the runner's projection to serve both would put a `kind` it does
+// not have on the arbiter's call.
+//
+// `ids` IS the reachable set (missing + stale + partial), never missing-input,
+// deferred or blocked; see BackfillSnapshotEntry. So a leaf pointed at an
+// operation claims only work the lane can actually do, which is exactly the
+// contract a bucket carries.
+
+export type ArbiterChannelMeta = { slug: string; platform: Platform | null };
+
+// Matches autoRunner's own read concurrency: these are small JSON files, and
+// the ceiling is there to keep a 68-channel corpus from opening 68 handles at
+// once on a box already running an engine.
+const SNAPSHOT_READ_CONCURRENCY = 8;
+
+export async function buildArbiterChannelWork(
+ paths: Paths,
+ meta: ReadonlyArray<ArbiterChannelMeta>,
+ operations: ReadonlyArray<string>,
+): Promise<ChannelWork[]> {
+ if (operations.length === 0) return [];
+ const snaps = await mapConcurrent(meta, SNAPSHOT_READ_CONCURRENCY, (m) =>
+ readChannelSnapshot(paths, m.slug),
+ );
+ const out: ChannelWork[] = [];
+ for (const [i, { slug, platform }] of meta.entries()) {
+ const snap = snaps[i];
+ if (!snap) continue;
+ const ops: Record<string, string[]> = {};
+ for (const op of operations) ops[op] = snap.backfill?.[op]?.ids ?? [];
+ // `buckets` stays EMPTY. The arbiter dispatches operations; a bucket
+ // projected here would be claimable by a bucket leaf that the arbiter is
+ // about to filter out anyway, and the empty map makes that structural
+ // rather than a filtering step someone could later remove.
+ out.push({ slug, platform, buckets: {}, operations: ops });
+ }
+ return out;
+}
+
+// Slug + platform for every channel. Re-derived from the shared
+// listChannelConfigs rather than repeating a readdir-then-serial-read, and the
+// same derivation autoRunner's own listChannelMeta uses — so the arbiter and
+// the runners cannot disagree about which channels exist or what platform one
+// is on.
+export async function listArbiterChannelMeta(
+ paths: Paths,
+): Promise<ArbiterChannelMeta[]> {
+ const configs = await listChannelConfigs(paths);
+ return configs.map(({ slug, config }) => ({
+ slug,
+ platform: detectPlatform(config.url),
+ }));
+}
diff --git a/common/controller/autoRunner.ts b/common/controller/autoRunner.ts
@@ -22,6 +22,7 @@ import {
type WorkPick,
buildPendingByLeaf,
flattenLeaves,
+ retainLeaves,
selectNextWork,
defaultBucketsForPolicy,
selectableBucketsForKind,
@@ -383,11 +384,19 @@ export async function computeLeafPending(
meta,
channels,
);
- const pending = buildPendingByLeaf(
+ // BUCKET LEAVES ONLY. A leaf naming an operation belongs to the arbiter, and
+ // counting its work here would report a digest backlog as a transcription
+ // one; DISPATCHING it would hand a digest candidate to whisper. See
+ // retainLeaves.
+ const pending = retainLeaves(
+ buildPendingByLeaf(
+ policy.root,
+ channels,
+ defaultBucketsForPolicy(kind, policy),
+ compare ? { compare } : undefined,
+ ),
policy.root,
- channels,
- defaultBucketsForPolicy(kind, policy),
- compare ? { compare } : undefined,
+ "buckets",
);
// A video already in flight is not "next up" — drop the live runner's set
// before asking the policy, exactly as next() does.
@@ -686,11 +695,18 @@ async function runLoop(
metaCache,
channels,
);
- const pending = buildPendingByLeaf(
+ // BUCKET LEAVES ONLY — see the note in computeLeafPending. This is the
+ // dispatch path, so it is the one where getting it wrong runs the wrong
+ // engine on the wrong video.
+ const pending = retainLeaves(
+ buildPendingByLeaf(
+ policy.root,
+ channels,
+ defaultBucketsForPolicy(kind, policy),
+ compare ? { compare } : undefined,
+ ),
policy.root,
- channels,
- defaultBucketsForPolicy(kind, policy),
- compare ? { compare } : undefined,
+ "buckets",
);
const exclude = new Set<string>([...live.inFlight.keys(), ...completed]);
removeIds(pending, exclude);
diff --git a/common/jobs/autoQueuePolicy.ts b/common/jobs/autoQueuePolicy.ts
@@ -264,6 +264,38 @@ export type ChannelWork = {
operations?: Record<string, string[]>;
};
+// Split a pending map by WHO DISPATCHES IT.
+//
+// A leaf naming an operation is claimed by the arbiter; a leaf naming a bucket
+// (or nothing) is claimed by the runner whose tree it sits in. They can share a
+// tree — an operator may reasonably want "diarize cornbreadman before you
+// transcribe anything else" expressed in one priority list — and without this
+// they SHARE A DISPATCHER, which means the auto-transcribe runner picks a
+// digest candidate off an operation leaf and hands it to whisper.
+//
+// So each dispatcher zeroes the leaves that are not its own before selecting.
+// Zeroed rather than removed: selectNextWork walks the TREE, and a leaf with an
+// empty list is skipped by the same code path that skips a leaf with no work,
+// so nothing new has to understand this. A leaf missing from the map entirely
+// would instead read as `undefined` at `pending[node.id]?.length`, which is the
+// same answer by luck rather than by construction.
+export function retainLeaves(
+ pending: Record<string, string[]>,
+ root: AutoQueueNode,
+ want: "buckets" | "operations",
+): Record<string, string[]> {
+ const out: Record<string, string[]> = {};
+ const keep = new Set(
+ flattenLeaves(root)
+ .filter((l) => (want === "operations") === Boolean(l.match.operation))
+ .map((l) => l.id),
+ );
+ for (const [id, ids] of Object.entries(pending)) {
+ out[id] = keep.has(id) ? ids : [];
+ }
+ return out;
+}
+
// Pre-order (priority-order) flatten of every leaf in the tree.
export function flattenLeaves(node: AutoQueueNode): AutoQueueLeaf[] {
if (!isGroup(node)) return [node];
diff --git a/common/jobs/jobKinds.ts b/common/jobs/jobKinds.ts
@@ -222,6 +222,17 @@ const JOB_KINDS: Record<string, JobKindMeta> = {
// current channel", and NOT replayable: the sweep is armed through a
// persisted settings flag, so replaying a record would be a second way to
// start the same singleton.
+ // The unified dispatcher. Drainable (it stops starting new channel jobs and
+ // lets the in-flight ones finish) and never replayable — it is a long-lived
+ // loop, not a unit of work, and "retry" on it would mean "start a second
+ // dispatcher", which is precisely what it refuses to be.
+ "operations-arbiter": {
+ kind: "operations-arbiter",
+ label: "Operations arbiter",
+ drainable: true,
+ replayable: false,
+ queueKeyStrategy: "parallel",
+ },
"backfill-sweep": {
kind: "backfill-sweep",
label: "Backfill sweep",
diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md
@@ -11,6 +11,8 @@
- **Order and Reach sit at the foot of the lane, and now cover every pipeline.** *Newest first* used to be a control on two of the four lanes, halfway up a form. It is now the same block on all four, with the trade-off spelled out beneath it, and it is joined by **Reach** — *Within each channel* or *Across all channels* — which is disabled while the order is left at *Listed*, because a reach with no order to apply is not a setting. The digest and backfill lanes keep **both** of their switches, under the names they have everywhere else: the **sweep** decides whether there is a corpus-wide pass at all, and the **pause** decides whether the lane consumes it. They are deliberately not merged into one button: one is cheap to undo and the other costs a week of GPU time.
- **A runner is a lane, not a job, and the active-jobs screen finally says so.** Each always-on runner took a full bordered card — the same card a channel with real work in it gets — to say "running", which pushed the jobs that actually have progress bars below the fold. The runners are now a **strip** at the top: one line each, with the lane's state and, when it is not working, *why* — "nothing pending", "sweep armed, lane paused", "no enabled worker". Digest and backfill are on the strip too, so all four pipelines are visible at once from a screen that previously knew about two. Every control a runner had is still there, on its line. A screen with no jobs at all now shows the strip rather than only the words "No active jobs" — "every lane is idle and here is why" is the answer that screen was previously unable to give.
- **The digest lane's six safety rules are now written down where anything can read them.** They lived inside one function's closure — the global pause, standing aside for transcription (and the carve-out that stops a CPU-pinned worker counting as GPU contention), the metered spend cap, the two refusals for a disabled lane and an unreachable engine, and the rule that lets an aligned re-upload borrow a digest instead of generating one. Nothing outside that function could ask about any of them, which is exactly why this archive ended up with two separate schedulers for one kind of work: anything else wanting to run a digest would have had to reimplement all six from memory, correctly. They are one module now, and the digest runner calls it, so there is one definition and no second opinion to drift from it. Nothing changes about what happens — a paused lane still holds rather than stopping, a disabled metered lane still refuses without making a network call to find out, and the log still says which of the reasons is holding you up, printed once when the state changes rather than once per poll for a week.
+- **One dispatcher can now drive every pipeline, from the rules you already write.** There were four schedulers for one kind of work and only two of them were any good — the backfill sweep says in its own source that it is a copy of the digest sweep, written because the first one's shape could not be reused, and that a third would be written for the same reason. There is no third. A rule on the auto-queue page can now point at an **operation** — Digest, speaker diarization, speaker attribution — instead of at a bucket, and an **arbiter** dispatches whatever the rules name: one job per channel, on each operation's own lane, in the priority order the rule list already expresses. It starts the same per-channel jobs a hand-clicked run does, so a pass is inspectable with the tools that exist rather than being one opaque mega-job. Because a rule can be as specific as one channel, "diarize this channel before anything else touches the GPU" is now something you write down rather than something you sit and supervise. Two things it deliberately refuses: it will not start while either corpus sweep is armed — two dispatchers on one lane would start the same channel twice, and it names the switch in the way rather than failing quietly — and it does not come back by itself after a restart. The sweeps still do, and still own this work until the arbiter has driven a full pass over the archive; nothing you have running today changes.
+- **A rule pointed at an operation is no longer at risk of being transcribed.** Rules for operations and rules for buckets can share one priority list, which is the point — but they must not share a *dispatcher*, or the auto-transcribe runner picks a video off a digest rule and hands it to the speech engine. Each dispatcher now takes only the rules that are its own.
- **Newest-first could quietly hand a video another channel's upload date.** The date lookup's first and cheapest layer scans the transcript index, which is keyed by the id the *platform* reports; every video the auto-queue actually asks about is named by its **folder**, and on this archive those two disagree for about one video in seven — a Rumble folder is named for the URL, while its recorded id is the embed id. Matched across the whole corpus, one channel's folder name could collide with a different channel's recorded id and inherit its date, which is a *wrong* answer rather than a missing one, and wrong dates are exactly what an ordering setting cannot survive. The lookup is now scoped to the channel the video belongs to, so a folder name that is not an id in its own channel falls through to reading the date out of the folder it actually names. Separately, the first pass over a very large backlog can only date so many videos at once; it now says so once in the log, because the remainder sorting to the back and being picked up on the next pass is the design, not a fault.
- **The auto-queue can be told to do the newest uploads first, and it now genuinely does.** Both runners always took the first video off a rule's pile, and that pile's order came straight from the channel snapshot, which sorts most buckets **alphabetically by video id** — arbitrary for YouTube ids, and oldest-first for the date-prefixed folder names some sites use. So when a channel uploaded today, nothing made that video jump the nine-thousand-video backlog in front of it; the only reason auto-download roughly worked was that its one bucket happens to be left in playlist order. There is now an **Order** setting per runner — *Listed order* (what you have today, and still the default), *Newest first*, *Oldest first*. It sorts the videos **inside** each rule, across every channel and bucket that rule claims; the rule list still decides which rule goes first, because that is what the rule list is for. For a straight newest-first archive, use one catch-all rule. Worth knowing before you switch it on: under *Newest first* the retry and partial-download buckets lose their head start, so a half-finished download can end up waiting behind fresh work. The page says so next to the setting.
- **Working out how recent 79,000 videos are turned out to be nearly free, once we stopped guessing where the dates were.** The obvious source — reading each video's metadata file — is about six and a half minutes and 41 GB of reading, on every scheduling decision, which is a non-starter. The transcript index already holds a date per video in a form that can be scanned without decoding anything: **78,583 videos in well under a second**. That covers the corpus, but it turned out **not** to cover the videos auto-transcribe actually queues, because the index only holds videos that already *have* a transcript and auto-transcribe's whole job is the ones that don't — of the 870 videos genuinely pending here, it knew the date of **122**. The gap is closed by reading the last 8 KB of each remaining video's metadata file, where the upload date happens to sit: **868 of the 870, at a fifth of a millisecond each**, and remembered afterwards so it is paid once rather than every few seconds. Videos not downloaded yet have no date anywhere on disk at all, so auto-download estimates one from the video's position in the channel's newest-first listing; those show with a `≈`, and a brand-new upload with nothing dated above it goes to the front, which is the entire point. Anything still undatable sorts to the back rather than disappearing, and a missing or busy index degrades to the old ordering instead of stopping the runner.
diff --git a/editor/app/auto-queue/actions.ts b/editor/app/auto-queue/actions.ts
@@ -220,3 +220,33 @@ export async function saveLaneOrderAction(
revalidatePath("/auto-queue");
return { ok: true };
}
+
+// The arbiter: start / stop the unified dispatcher.
+//
+// Kept out of saveLaneOrderAction and out of the policy form: starting a
+// dispatcher is an operational act, not a setting, and routing it through a
+// form would mean a pending tree edit had to be saved or discarded first.
+//
+// It REFUSES rather than silently no-opping when a sweep is armed, and the
+// refusal carries the reason — two dispatchers on one lane would start the same
+// channel twice, and an operator who clicks Start and sees nothing happen
+// deserves to be told which switch is in the way.
+export async function startArbiterAction(): Promise<SaveResult> {
+ const { startArbiter } = await import(
+ "yt-dlp-transcript-common/controller/arbiter"
+ );
+ const result = await startArbiter();
+ revalidatePath("/auto-queue");
+ return result.jobId
+ ? { ok: true }
+ : { ok: false, error: result.error ?? "The arbiter could not be started." };
+}
+
+export async function stopArbiterAction(): Promise<SaveResult> {
+ const { stopArbiter } = await import(
+ "yt-dlp-transcript-common/controller/arbiter"
+ );
+ stopArbiter();
+ revalidatePath("/auto-queue");
+ return { ok: true };
+}
diff --git a/editor/app/auto-queue/components/ArbiterBar.tsx b/editor/app/auto-queue/components/ArbiterBar.tsx
@@ -0,0 +1,113 @@
+"use client";
+
+import Link from "next/link";
+import { useState, useTransition } from "react";
+import { Button } from "yt-dlp-transcript-common/components/ui/button";
+import { startArbiterAction, stopArbiterAction } from "../actions";
+import type { ArbiterStatus } from "../lanes";
+
+// THE ARBITER, on one line, under the rail.
+//
+// It sits above the lane switcher and outside every lane because it is not a
+// lane — it is the thing that decides which lane runs. Putting it inside one
+// would say the opposite.
+//
+// A <div>, never a <section>: nesting one inside this page's structure breaks
+// every heading-scoped lookup in the e2e suite. The name is plain text for the
+// same reason the rail's names are.
+//
+// IT SAYS WHY IT WILL NOT START, before the click rather than after. The two
+// reasons are an armed sweep (two dispatchers on one lane would start the same
+// channel twice) and a tree with no operation rule in it (it would come up,
+// find nothing to dispatch, and stop). Both come from the server, from the same
+// function that would refuse the start, so the page cannot promise something
+// the action then declines.
+
+export function ArbiterBar({
+ arbiter,
+ onRefresh,
+}: {
+ arbiter: ArbiterStatus;
+ onRefresh: () => void | Promise<void>;
+}) {
+ const [pending, startTransition] = useTransition();
+ const [error, setError] = useState<string | null>(null);
+
+ const run = (action: () => Promise<{ ok: boolean; error?: string }>) => {
+ startTransition(async () => {
+ const result = await action();
+ setError(result.ok ? null : (result.error ?? "Failed."));
+ await onRefresh();
+ });
+ };
+
+ const noRules = arbiter.operationRules === 0;
+ const blocked = arbiter.blockedReason ?? (noRules ? NO_RULES : null);
+
+ return (
+ <div className="flex flex-col gap-1.5 rounded-md border border-border bg-card px-3 py-2">
+ <div className="flex flex-wrap items-center gap-x-3 gap-y-1 text-sm">
+ <span className="font-mono text-xs uppercase tracking-[0.14em] text-muted-foreground">
+ Arbiter
+ </span>
+ <span
+ aria-hidden="true"
+ className={`size-2 shrink-0 rounded-full ${
+ arbiter.running
+ ? "bg-info animate-pulse motion-reduce:animate-none"
+ : "bg-muted-foreground/40"
+ }`}
+ />
+ <span className={arbiter.running ? "text-foreground" : "text-muted-foreground"}>
+ {arbiter.running ? "Running" : "Stopped"}
+ </span>
+ <span className="text-xs text-muted-foreground">
+ {arbiter.operationRules === 0
+ ? "no rule names an operation"
+ : `${arbiter.operationRules} rule${arbiter.operationRules === 1 ? "" : "s"} name an operation`}
+ </span>
+ {arbiter.jobId && (
+ <Link
+ href={`/jobs/${arbiter.jobId}`}
+ className="text-xs underline underline-offset-2 text-muted-foreground hover:text-foreground"
+ >
+ log
+ </Link>
+ )}
+ <span className="ml-auto">
+ {arbiter.running ? (
+ <Button
+ type="button"
+ size="sm"
+ variant="outline"
+ aria-label="Stop the arbiter"
+ disabled={pending}
+ onClick={() => run(stopArbiterAction)}
+ >
+ Stop
+ </Button>
+ ) : (
+ <Button
+ type="button"
+ size="sm"
+ aria-label="Start the arbiter"
+ disabled={pending || blocked !== null}
+ onClick={() => run(startArbiterAction)}
+ >
+ Start
+ </Button>
+ )}
+ </span>
+ </div>
+ <p className="text-xs text-muted-foreground">
+ {blocked ??
+ "One dispatcher for every operation a rule names: one job per channel, on each operation's own lane."}
+ </p>
+ {error && <p className="text-xs text-warning">{error}</p>}
+ </div>
+ );
+}
+
+const NO_RULES =
+ "No rule names an operation yet, so there would be nothing to dispatch. " +
+ "Point a rule at Digest or a backfill operation first.";
diff --git a/editor/app/auto-queue/components/AutoQueueView.tsx b/editor/app/auto-queue/components/AutoQueueView.tsx
@@ -7,6 +7,7 @@ import type {
AutoQueueKindStatus,
PlatformCooldownView,
} from "../status";
+import { ArbiterBar } from "./ArbiterBar";
import { InFlightList } from "./InFlightList";
import { LaneHeader } from "./LaneHeader";
import { NextUp } from "./NextUp";
@@ -118,6 +119,8 @@ export function AutoQueueView({
selectedId={highlightedOperation}
/>
+ <ArbiterBar arbiter={data.lanes.arbiter} onRefresh={refresh} />
+
<div className="flex flex-wrap items-center gap-2">
<label
htmlFor={switcherId}
diff --git a/editor/app/auto-queue/lanes.ts b/editor/app/auto-queue/lanes.ts
@@ -10,6 +10,11 @@ import {
DIGEST_REMOTE_QUEUE,
BACKFILL_QUEUE,
} from "yt-dlp-transcript-common/lib/queueKeys";
+import {
+ arbiterBlockedReason,
+ getArbiterJobId,
+} from "yt-dlp-transcript-common/controller/arbiter";
+import { flattenLeaves } from "yt-dlp-transcript-common/jobs/autoQueuePolicy";
import { getChannelBriefs } from "../lib/requestCache";
import {
buildOperationBands,
@@ -55,11 +60,28 @@ export type SweepLaneStatus = {
reach: AutoQueueReach;
};
+// The unified dispatcher's state. Not a lane: it does not do work, it decides
+// which lane does. It is on this payload because the console is where an
+// operator turns it on, and because `blockedReason` has to come from the same
+// function that would refuse the start — an operator told "you can start it"
+// and then refused has been lied to by a race between two copies of one rule.
+export type ArbiterStatus = {
+ running: boolean;
+ jobId: string | null;
+ // Why it will not start, or null when it will. See arbiterBlockedReason.
+ blockedReason: string | null;
+ // How many rules across both trees name an operation it dispatches. Zero is
+ // why it would come straight back up and stop, so the console says so before
+ // the click rather than after.
+ operationRules: number;
+};
+
export type AutoQueueLanesPayload = {
digest: SweepLaneStatus;
backfill: SweepLaneStatus;
// One band per pipeline, in rail order. See lib/operationBands.ts.
bands: OperationBand[];
+ arbiter: ArbiterStatus;
};
// Running jobs on a lane's queue keys. The digest lane has two (local and
@@ -121,5 +143,16 @@ export async function buildAutoQueueLanes(): Promise<AutoQueueLanesPayload> {
reach: settings.backfill.reach,
},
bands,
+ arbiter: {
+ running: getArbiterJobId() !== null,
+ jobId: getArbiterJobId(),
+ blockedReason: arbiterBlockedReason(),
+ operationRules: [
+ settings.autoQueue.transcription.root,
+ settings.autoQueue.download.root,
+ ]
+ .flatMap((root) => flattenLeaves(root))
+ .filter((leaf) => Boolean(leaf.match.operation)).length,
+ },
};
}
diff --git a/editor/e2e/auto-queue.spec.ts b/editor/e2e/auto-queue.spec.ts
@@ -22,7 +22,7 @@ import { baseUrl } from "./baseUrl";
type Leaf = {
id: string;
- match: { type: string; value?: string; bucket?: string };
+ match: { type: string; value?: string; bucket?: string; operation?: string };
weight?: number;
maxWorkers?: number | null;
};
@@ -399,7 +399,19 @@ test("UI: build a policy in the editor, save it, and start the runner", async ({
).toBe(true);
// Start the runner from the UI; it should transcribe alpha's videos.
- await section.getByRole("button", { name: "Start Auto-transcribe" }).click();
+ //
+ // Clicked ONLY IF IT IS STILL OFFERED. Saving an enabled policy starts the
+ // runner itself (saveAutoQueueAction calls startAutoRunner), so Start is
+ // disabled the moment the 3-second poll notices — and an unconditional click
+ // is a race against that poll, not an assertion about anything. The claim
+ // here is that the runner ends up RUNNING and does the work, not that one
+ // particular button was what started it, so both halves are asserted
+ // directly.
+ const start = section.getByRole("button", { name: "Start Auto-transcribe" });
+ if (await start.isEnabled()) await start.click();
+ await expect(section.getByText("Runner running")).toBeVisible({
+ timeout: 20_000,
+ });
await expect
.poll(async () => allTranscribed("alpha", ["a1", "a2"]), { timeout: 60_000 })
.toBe(true);
@@ -1189,3 +1201,112 @@ test("Active jobs: a runner is a lane on a strip, not a card", async ({
// And the generic bucket stays gone.
await expect(page.getByRole("heading", { name: "Other" })).toHaveCount(0);
});
+
+
+// --- The arbiter -------------------------------------------------------------
+
+test("the arbiter refuses to start beside an armed sweep, and says which one", async ({
+ page,
+}) => {
+ await resetData(null);
+ await makeChannel("alpha", ["a1"]);
+ await writeSettings({
+ adminTitle: "Test Admin",
+ maxTranscriptPageBytes: 8388608,
+ sleepBetweenDownloadsSeconds: 0,
+ minFreeDiskGB: 0,
+ workers: ONE_WORKER,
+ // A rule DOES name an operation, so the only thing in the way is the sweep.
+ autoQueue: transcriptionAutoQueue({
+ id: "root",
+ mode: "strict",
+ children: [{ id: "dig", match: { type: "all", operation: "digest" } }],
+ }),
+ digest: { sweepEnabled: true },
+ });
+
+ await page.goto("/auto-queue");
+ await awaitHydration(
+ page.locator("section", {
+ has: page.getByRole("heading", { name: "Auto-transcribe" }),
+ }),
+ );
+ const start = page.getByRole("button", { name: "Start the arbiter" });
+ await expect(start).toBeDisabled();
+ // Two dispatchers on one lane would start the same channel twice, so the
+ // refusal names the switch that is in the way rather than failing silently.
+ await expect(page.getByText(/digest sweep is armed/i)).toBeVisible();
+});
+
+test("the arbiter will not start with no rule naming an operation", async ({
+ page,
+}) => {
+ await resetData(null);
+ await makeChannel("alpha", ["a1"]);
+ await writeSettings({
+ adminTitle: "Test Admin",
+ maxTranscriptPageBytes: 8388608,
+ sleepBetweenDownloadsSeconds: 0,
+ minFreeDiskGB: 0,
+ workers: ONE_WORKER,
+ autoQueue: transcriptionAutoQueue(ALPHA_ROOT),
+ });
+
+ await page.goto("/auto-queue");
+ await awaitHydration(
+ page.locator("section", {
+ has: page.getByRole("heading", { name: "Auto-transcribe" }),
+ }),
+ );
+ // It would come up, find nothing to dispatch and stop. Saying so before the
+ // click is the difference between a disabled button and a mystery.
+ await expect(
+ page.getByRole("button", { name: "Start the arbiter" }),
+ ).toBeDisabled();
+ // The explanatory line, not the status chip beside it — both say the same
+ // thing, deliberately, so the assertion names which one it means.
+ await expect(
+ page.getByText(/No rule names an operation yet, so there would be nothing/),
+ ).toBeVisible();
+});
+
+test("the arbiter starts, appears as a job, and stops", async ({ page }) => {
+ await resetData(null);
+ await makeChannel("alpha", ["a1"]);
+ await writeSettings({
+ adminTitle: "Test Admin",
+ maxTranscriptPageBytes: 8388608,
+ sleepBetweenDownloadsSeconds: 0,
+ minFreeDiskGB: 0,
+ workers: ONE_WORKER,
+ autoQueue: transcriptionAutoQueue({
+ id: "root",
+ mode: "strict",
+ children: [{ id: "dig", match: { type: "all", operation: "digest" } }],
+ }),
+ });
+
+ await page.goto("/auto-queue");
+ await awaitHydration(
+ page.locator("section", {
+ has: page.getByRole("heading", { name: "Auto-transcribe" }),
+ }),
+ );
+ await page.getByRole("button", { name: "Start the arbiter" }).click();
+ const stop = page.getByRole("button", { name: "Stop the arbiter" });
+ await expect(stop).toBeVisible({ timeout: 15_000 });
+
+ // It runs on queueKey "" — a long-lived job that WAITS on jobs needing the
+ // real keys must not hold one, or it deadlocks against its own work.
+ const job = await page.request
+ .get(`${baseUrl}/api/jobs/active`)
+ .then((r) => r.json() as Promise<{ jobs: { kind: string; queueKey: string }[] }>);
+ const arbiter = job.jobs.find((j) => j.kind === "operations-arbiter");
+ expect(arbiter).toBeTruthy();
+ expect(arbiter?.queueKey).toBe("");
+
+ await stop.click();
+ await expect(
+ page.getByRole("button", { name: "Start the arbiter" }),
+ ).toBeVisible({ timeout: 15_000 });
+});
diff --git a/plans/unified-operations-model.md b/plans/unified-operations-model.md
@@ -96,10 +96,16 @@ change** — every one of them moves live numbers on a 78,000-video corpus.
drawing from a separate `ChannelWork.operations` id space, with the claim key widened to
`${operation}\0${id}` so digest and diarization cannot steal each other's work. The tree
can *express* operation work; nothing dispatches it yet.
-3. **One arbiter per lane, not per subsystem.** A single runner reads the tree, groups
- selected work by `lane.queueKey`, and dispatches one job per key. This is where
- `digestSweep` and `backfillSweep` finally merge — they are the same loop with different
- constants.
+3. **One arbiter per lane, not per subsystem.** — **DONE** (`common/controller/arbiter.ts`).
+ One long-lived job on `queueKey: ""` — it WAITS on jobs needing the real keys, so holding
+ one would deadlock against its own work. It reads both policy trees, plans one unit per
+ (operation, channel) in leaf-priority order, and dispatches ONE UNIT PER LANE per pass,
+ concurrently, reserving by `lane.queueKey` and never inventing one. It starts the
+ per-channel jobs that already exist, so a pass is inspectable with the existing tools.
+ It REFUSES to start beside an armed sweep and says which one — two dispatchers on one
+ lane would start the same channel twice. It is deliberately NOT persisted across a
+ restart; that comes with step 6. `digestSweep` and `backfillSweep` are untouched and
+ still own dispatch until then.
4. **Move the guards onto declared lane rules.** — **DONE.** `common/controller/laneGuards.ts`
holds all six, in the two shapes a dispatcher can act on: a **preflight** answered before
any pool exists (`remoteEnabled`, the engine `probe()` — both still throw, because a