commit 39714d6514f7243a687670f4bce499a3f3569ccd
parent 26689d10b832efe9345baa20d62a27facfeb089e
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Fri, 11 Sep 2026 12:06:12 -0400
Merge channel-priority/s1 — the dispatch slice
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Diffstat:
3 files changed, 623 insertions(+), 16 deletions(-)
diff --git a/common/controller/autoRunner.test.ts b/common/controller/autoRunner.test.ts
@@ -0,0 +1,99 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import { focusHoldLine, makeFocusHoldReporter } from "./autoRunner";
+import type { FocusSummary } from "../lib/channelPriority";
+
+// THE RUNNER'S HALF OF S1 (plans/channel-priority.md), and it is here rather
+// than in `jobs/channelPriorityCompile.test.ts` for one reason:
+// `../architecture.test.ts` forbids `jobs/ -> controller/`, including from a
+// test file, and its ALLOWED list may only shrink. The engine-level proof — a
+// focus holding and releasing a real `buildPendingByLeaf` + `selectNextWork`
+// pair — lives there; what lives here is the thing that can only be said about
+// the runner: it writes ONE line per transition.
+//
+// WHY A LINE AND NOT AN IDLE REASON. A lane whose focus group holds the rest is
+// not idle, it is dispatching focus work, so `AutoRunnerIdleReason` gains
+// nothing and `no-pending` stays the true answer for an empty tree. The log
+// line is what makes "why is only jeralyzer moving?" answerable from the job
+// log, and the banner (S4) is what makes it answerable from the page.
+
+function summary(over: Partial<FocusSummary> = {}): FocusSummary {
+ return {
+ kind: "channels",
+ siteId: null,
+ slugs: ["slow-a"],
+ channelCount: 1,
+ active: true,
+ focusPending: 2,
+ otherPending: 5,
+ holding: true,
+ heldChannels: 1,
+ ...over,
+ };
+}
+
+test("the focus-hold line fires once per state change, not once per tick", () => {
+ const lines: string[] = [];
+ const report = makeFocusHoldReporter("download", (line) => lines.push(line));
+
+ // No focus at all: silence, however many ticks run.
+ report(null);
+ report(null);
+ assert.deepEqual(lines, []);
+
+ // Entering the hold: one line.
+ report(summary());
+ report(summary());
+ report(summary());
+ assert.equal(lines.length, 1);
+ assert.match(lines[0], /Auto-download: focus \(1 channel\) is holding/);
+
+ // THE COUNTS MOVE ON EVERY GRANT AND MUST NOT RE-FIRE IT. The state is the
+ // three-valued thing; the counts are in the message only, which is why the
+ // reporter keys on the state and not on the line it printed.
+ report(summary({ focusPending: 1 }));
+ report(summary({ focusPending: 1, heldChannels: 2, otherPending: 9 }));
+ assert.equal(lines.length, 1);
+
+ // The focus runs out: the release line, once.
+ report(summary({ focusPending: 0, holding: false }));
+ report(summary({ focusPending: 0, holding: false }));
+ assert.equal(lines.length, 2);
+ assert.match(lines[1], /has no work left in this lane/);
+
+ // New focus work retakes the lane: the hold line again.
+ report(summary());
+ assert.equal(lines.length, 3);
+ assert.match(lines[2], /is holding/);
+
+ // The focus ends (or resolves to nothing): back to silence, no line.
+ report(null);
+ report(null);
+ assert.equal(lines.length, 3);
+});
+
+test("an inactive focus is the same state as no focus", () => {
+ const lines: string[] = [];
+ const report = makeFocusHoldReporter("digest", (line) => lines.push(line));
+ // What a `{kind:"site"}` focus naming an unknown site resolves to: a summary
+ // exists, but nothing is focused, so the lane is not holding for anyone.
+ report(summary({ active: false, channelCount: 0, slugs: [], holding: false }));
+ assert.deepEqual(lines, []);
+ assert.equal(
+ focusHoldLine("digest", summary({ active: false })),
+ null,
+ );
+});
+
+test("the line names the lane and pluralises the focus set", () => {
+ assert.match(
+ focusHoldLine("transcription", summary({ channelCount: 30 })) ?? "",
+ /Auto-transcription: focus \(30 channels\) is holding this lane — 2 focus unit\(s\) pending, 1 channel\(s\) held\./,
+ );
+ assert.match(
+ focusHoldLine("backfill", summary({ holding: false, otherPending: 7 })) ??
+ "",
+ /Auto-backfill: focus \(1 channel\) has no work left in this lane — 7 unit\(s\) released/,
+ );
+ assert.equal(focusHoldLine("download", null), null);
+});
diff --git a/common/controller/autoRunner.ts b/common/controller/autoRunner.ts
@@ -78,6 +78,16 @@ import { type DownloadOutcomeStatus } from "../lib/downloadOutcome";
import { downloadQueueKey } from "../lib/queueKeys";
import { isGateHeld } from "../lib/pauseGates";
import {
+ type ChannelPriority,
+ type FocusSummary,
+ type SiteChannelIndex,
+ compileLaneRoot,
+ focusSummary,
+ isChannelPaused,
+ resolveFocusSlugs,
+} from "../lib/channelPriority";
+import { listSites } from "../lib/site";
+import {
listChannelConfigs,
readChannelConfig,
readChannelSnapshotShared,
@@ -285,12 +295,199 @@ const SNAPSHOT_READ_CONCURRENCY = 64;
// Re-derived from the shared listChannelConfigs rather than repeating the
// readdir-then-serial-read here.
-async function listChannelMeta(paths: Paths): Promise<ChannelMeta[]> {
+//
+// PAUSED IS A FILTER ON THE CHANNEL LIST, not a shape in the tree, and this is
+// the one predicate that makes it so (plans/channel-priority.md, decision 3). A
+// tree cannot express exclusion — an `{type:"all"}` catch-all matches
+// everything, and first-match-wins would let a catch-all placed above the Low
+// group swallow Low's work — so a paused channel is removed from the LIST every
+// leaf draws from instead. This function is the single source of that list for
+// BOTH the runner loop and `computeLeafPending`, so a paused channel is absent
+// from the lane's draw and from the status panel's pending counts in one edit.
+//
+// PER LANE, through `isChannelPaused(model, slug, lane)` — the EFFECTIVE tier
+// for the lane being listed, so the per-operation override map decides: a
+// channel with `{tier:"normal", overrides:{sync:"paused"}}` (what all 15 live
+// `excludeFromSync` channels migrate to) is still drawn by the download lane,
+// and one with `{tier:"paused"}` or `overrides:{download:"paused"}` is not.
+// Re-evaluated on CHANNEL_LIST_TTL_MS, which is the clock for this decision.
+async function listChannelMeta(
+ paths: Paths,
+ kind: AutoQueueKind,
+ priority: ChannelPriority,
+): Promise<ChannelMeta[]> {
const configs = await listChannelConfigs(paths);
- return configs.map(({ slug, config }) => ({
- slug,
- platform: detectPlatform(config.url),
- }));
+ return configs
+ .filter(({ slug }) => !isChannelPaused(priority, slug, kind))
+ .map(({ slug, config }) => ({
+ slug,
+ platform: detectPlatform(config.url),
+ }));
+}
+
+// --- The compiled priority trees -------------------------------------------
+//
+// THE MODEL COMPILES, IT IS NOT CONSULTED. `settings.channelPriority` holds one
+// tier per channel plus one focus selector; `compileLaneRoot` turns that into
+// the lane's `AutoQueueGroup` (focus > normal > low > catch-all, all strict),
+// and dispatch runs the ordinary engine over it. Nothing in
+// `buildPendingByLeaf`, `selectNextWork`, `operationBatch`, `laneLimit` or
+// `pauseGates` learns a second priority mechanism, and "a zero limit is a hold,
+// never a stop" (controller/operationBatch.ts:22-25) is preserved trivially
+// because nothing here ever returns a limit.
+//
+// A FOCUS HOLDS THE REST BECAUSE STRICT DESCENT ALREADY DOES, and no second
+// mechanism is added: `pick()` (jobs/autoQueuePolicy.ts) filters a strict
+// group's children to those WITH WORK and descends into the first of them, and
+// it is re-asked on every grant. So focus work present => nothing below it is
+// picked; focus work exhausted => the next group runs; new focus work arriving
+// => the very next pick retakes the lane.
+//
+// WHERE THE HAND-EDITED TREES GO. The compiled root REPLACES
+// `settings.autoQueue[lane].root` at `laneDispatchRoot` below — the stored tree
+// is not consulted at all while a model exists. That is what makes the compiler
+// the ONE WRITER of channel priority: a `PolicyTreeEditor` save can still put a
+// channel leaf in the stored tree, but it cannot change what this lane
+// dispatches, so the two cannot fight — the compiler simply wins. (S3's
+// `saveChannelPriorityAction` then also PERSISTS the compiled roots, so the
+// stored tree and this one agree on disk; S4 makes the editor read-only for
+// compiled groups, which is the UI catching up with this fact.)
+//
+// AN ABSENT MODEL CHANGES NOTHING, BYTE FOR BYTE. `isDefaultChannelPriority`
+// below is the gate: no focus and no channel entries means the compiler never
+// runs and the stored trees stand exactly as they are today.
+
+const PRIORITY_CONTEXT_TTL_MS = 60_000;
+
+// The resolved half of the model — the half that costs I/O. Rebuilt when the
+// stored document changes or the TTL lapses, the same two triggers the
+// operation lanes' run context uses (settings key + 60 s), and for the same
+// reason: `resolveFocusSlugs` reads `transcripts/sites/*/site.json` through
+// `listSites`, and a `{kind:"site"}` focus tracks that file's membership rather
+// than freezing a list. The TTL is what makes a channel added to the focused
+// site join the focus without a settings write.
+type PriorityContext = {
+ key: string;
+ at: number;
+ model: ChannelPriority;
+ focusSlugs: string[];
+};
+
+let priorityContext: PriorityContext | null = null;
+
+// No focus and no per-channel entry: the document says nothing, so the compiler
+// must not run. Local to the runner rather than in lib/channelPriority.ts
+// because it is a dispatch-side question ("is there anything to compile"), not
+// part of the model's contract.
+function isDefaultChannelPriority(model: ChannelPriority): boolean {
+ return (
+ model.focus.kind === "none" && Object.keys(model.channels).length === 0
+ );
+}
+
+function priorityContextFor(paths: Paths): PriorityContext {
+ const model = getSettings().channelPriority;
+ // The paths go in the key so two worktrees' runners in one process cannot
+ // share a focus resolved against the other's sites directory.
+ const key = JSON.stringify([model, paths.sitesDir]);
+ const now = Date.now();
+ if (
+ priorityContext &&
+ priorityContext.key === key &&
+ now - priorityContext.at < PRIORITY_CONTEXT_TTL_MS
+ ) {
+ return priorityContext;
+ }
+ // ONLY A SITE FOCUS READS THE SITES DIRECTORY. `{kind:"channels"}` and
+ // `{kind:"none"}` resolve from the document alone, so the common case pays
+ // nothing. `listSites` is the existing reader (lib/site.ts) — there is no
+ // second one here.
+ const siteChannels: Record<string, string[]> = {};
+ if (model.focus.kind === "site") {
+ for (const site of listSites(paths)) {
+ siteChannels[site.siteId] = site.channels.map((c) => c.slug);
+ }
+ }
+ const index: SiteChannelIndex = siteChannels;
+ priorityContext = {
+ key,
+ at: now,
+ model,
+ focusSlugs: resolveFocusSlugs(model, index),
+ };
+ return priorityContext;
+}
+
+// THE ROOT THIS LANE ACTUALLY DISPATCHES FROM. One function, called by the
+// runner loop AND by computeLeafPending, so the status panel can never name a
+// leaf the runner does not have.
+//
+// `slugs` is the lane's own (already paused-filtered) channel list;
+// `compileLaneRoot` re-applies the same per-lane predicate, so handing it the
+// filtered list and handing it every slug produce the identical tree for THIS
+// lane. Compiling is pure and O(channels) — ~69 string pushes against the
+// ~6.5 MB of snapshot JSON the same tick folds — so it happens per tick and
+// only the focus resolution above is cached.
+function laneDispatchRoot(
+ kind: AutoQueueKind,
+ policy: AutoQueuePolicy,
+ ctx: PriorityContext,
+ slugs: readonly string[],
+): AutoQueueGroup {
+ if (isDefaultChannelPriority(ctx.model)) return policy.root;
+ return compileLaneRoot(kind, ctx.model, slugs, ctx.focusSlugs);
+}
+
+// The once-per-state-change line a lane writes while a focus is holding it.
+//
+// NOT AN IDLE REASON, and deliberately not: a lane whose focus group holds the
+// rest is not idle, it is dispatching focus work — `AutoRunnerIdleReason` stays
+// exactly as it is, and `no-pending` remains the true answer when the whole
+// tree is empty. This is the runner's LOG saying which of three states it is
+// in, so "why is only jeralyzer moving?" is answerable from the job log alone.
+//
+// The state is the three-valued thing, NOT the counts: the counts are in the
+// message but never in the key, or every completed focus unit would re-fire the
+// line. Same shape as the snooze line and the disk-gate line above.
+export function focusHoldState(summary: FocusSummary | null): string {
+ if (!summary || !summary.active) return "none";
+ return summary.holding ? "hold" : "free";
+}
+
+export function focusHoldLine(
+ kind: AutoQueueKind,
+ summary: FocusSummary | null,
+): string | null {
+ if (!summary || !summary.active) return null;
+ const channels = `${summary.channelCount} channel${summary.channelCount === 1 ? "" : "s"}`;
+ if (summary.holding) {
+ return (
+ `Auto-${kind}: focus (${channels}) is holding this lane — ` +
+ `${summary.focusPending} focus unit(s) pending, ` +
+ `${summary.heldChannels} channel(s) held.`
+ );
+ }
+ return (
+ `Auto-${kind}: focus (${channels}) has no work left in this lane — ` +
+ `${summary.otherPending} unit(s) released to the rest of the corpus.`
+ );
+}
+
+// The gate itself, as a closure so the runner keeps one line per transition and
+// the test can drive the transitions without a corpus. Logs on entering "hold"
+// and on entering "free"; says nothing while there is no active focus.
+export function makeFocusHoldReporter(
+ kind: AutoQueueKind,
+ onLog: (line: string) => void,
+): (summary: FocusSummary | null) => void {
+ let state = "none";
+ return (summary) => {
+ const next = focusHoldState(summary);
+ if (next === state) return;
+ state = next;
+ const line = focusHoldLine(kind, summary);
+ if (line) onLog(line);
+ };
}
// Read each channel's snapshot and project the buckets this runner kind cares
@@ -576,7 +773,15 @@ export async function computeLeafPending(
paths: Paths = getPaths(),
): Promise<LeafPending> {
const policy = getSettings().autoQueue[kind];
- const meta = await listChannelMeta(paths);
+ // THE SAME TWO PRIORITY DECISIONS THE RUNNER MAKES, in the same order: the
+ // paused filter on the channel list, then the compiled root. This function's
+ // whole contract is that its numbers are the runner's numbers, so both sides
+ // of channel priority have to be here too — otherwise the panel would count
+ // pending work for a paused channel, or attribute it to a stored leaf the
+ // runner is not dispatching from.
+ const ctx = priorityContextFor(paths);
+ const meta = await listChannelMeta(paths, kind, ctx.model);
+ const root = laneDispatchRoot(kind, policy, ctx, meta.map((m) => m.slug));
const laneOperations = laneOperationIds(kind);
const { channels, owner } = await buildChannelWork(
paths,
@@ -596,7 +801,7 @@ export async function computeLeafPending(
// leaf naming anything else finds no list and comes back empty — the zero
// retainLeaves used to apply afterwards, reached by construction.
const pending = buildPendingByLeaf(
- policy.root,
+ root,
channels,
defaultDrawsForPolicy(kind, policy, bucketLaneOperationId(kind)),
{ ...(compare ? { compare } : {}), defaultOperations: laneOperations },
@@ -633,14 +838,14 @@ export async function computeLeafPending(
currentWeights: { ...state[kind].runtime.currentWeights },
};
const pick = selectNextWork(
- policy.root,
+ root,
pending,
runtime,
live ? { ...live.active } : {},
);
let nextUp: NextUpView | null = null;
if (pick) {
- const order = flattenLeaves(policy.root);
+ const order = flattenLeaves(root);
const at = order.findIndex((l) => l.id === pick.leafId);
nextUp = {
videoId: pick.videoId,
@@ -806,12 +1011,16 @@ async function runLoop(
let metaCache: ChannelMeta[] = [];
let metaAt = 0;
+ // The root the last pick was made from. next() sets it every tick;
+ // runOperationPick reads it to resolve the leaf its own pick named.
+ let dispatchRoot: AutoQueueGroup = getSettings().autoQueue[kind].root;
// Whether the last next() saw the disk gate closed. next() runs on every
// scheduling tick, so without this the log fills with one identical line per
// tick for as long as the disk is full — which is precisely the situation in
// which the log needs to stay readable. Logged on each transition instead.
let diskIdle = false;
+ const reportFocusHold = makeFocusHoldReporter(kind, onLog);
// A graceful stop (the job record is gone after an e2e reset, or the policy was
// disabled) is modeled as a soft drain: stop picking, let in-flight finish.
@@ -1075,11 +1284,25 @@ async function runLoop(
}
// Refresh the (rarely-changing) channel list/platforms on a TTL.
+ //
+ // THE PAUSED FILTER RIDES THIS CLOCK. `listChannelMeta` drops every channel
+ // whose effective tier for THIS lane is `paused`, so the 30 s TTL is also
+ // how long a pause takes to reach dispatch. `metaAt === 0` rather than
+ // `metaCache.length === 0` is the cache-miss test now: a lane on which
+ // every channel is paused has a legitimately empty list, and the old
+ // sentinel would re-read 68 configs on every three-second tick for it.
const now = Date.now();
- if (now - metaAt > CHANNEL_LIST_TTL_MS || metaCache.length === 0) {
- metaCache = await listChannelMeta(paths);
+ const ctx = priorityContextFor(paths);
+ if (metaAt === 0 || now - metaAt > CHANNEL_LIST_TTL_MS) {
+ metaCache = await listChannelMeta(paths, kind, ctx.model);
metaAt = now;
}
+ // THE COMPILED ROOT ENTERS HERE, and this is the only place it does for the
+ // dispatch path: everything below — the projection, the completed filter,
+ // the pick and the leaf lookup in runOperationPick — reads `dispatchRoot`,
+ // never `policy.root`. See laneDispatchRoot.
+ const root = laneDispatchRoot(kind, policy, ctx, metaCache.map((m) => m.slug));
+ dispatchRoot = root;
const slugToPlatform = new Map(
metaCache.map((m) => [m.slug, m.platform ?? "unknown"]),
);
@@ -1105,7 +1328,7 @@ async function runLoop(
// the dispatch path, so it is the one where getting it wrong runs the wrong
// engine on the wrong video.
const pending = buildPendingByLeaf(
- policy.root,
+ root,
channels,
defaultDrawsForPolicy(kind, policy, bucketLaneOperationId(kind)),
{ ...(compare ? { compare } : {}), defaultOperations: laneOperations },
@@ -1115,8 +1338,16 @@ async function runLoop(
// also how the dependency order is kept — diarization finishes before
// attribution-diarized is offered the same video.
removeIds(pending, new Set(live.inFlight.keys()));
- dropCompleted(pending, policy.root, laneOperations);
+ dropCompleted(pending, root, laneOperations);
const pendingBeforeGates = countPending(pending);
+ // Say — ONCE per transition — whether a focus is holding this lane. Read off
+ // the `prio-*` leaf ids in the map just built, so it costs one pass over
+ // keys and no new read, and it is skipped entirely while no focus resolves.
+ reportFocusHold(
+ ctx.focusSlugs.length > 0
+ ? focusSummary(ctx.model, ctx.focusSlugs, pending)
+ : null,
+ );
// THE RUNNER JOB'S OWN BAR, on the metrics the per-channel jobs already
// use — so a lane's runner row reads like the manual verb's row rather than
// like an opaque long-lived loop. `target` moves as the corpus does (this
@@ -1167,7 +1398,7 @@ async function runLoop(
}
}
- const pick = selectNextWork(policy.root, pending, runtime, live.active);
+ const pick = selectNextWork(root, pending, runtime, live.active);
if (!pick) {
// Attribute the idleness. Nothing pending at all is a different situation
// from work that exists but is unreachable, and both differ from work the
@@ -1237,8 +1468,10 @@ async function runLoop(
const laneOperations = laneOperationIds(kind);
const opened = laneRun.run;
if (!opened) return { outcome: "skipped" };
- const policy = getSettings().autoQueue[kind];
- const leaf = flattenLeaves(policy.root).find(
+ // THE ROOT THE PICK CAME FROM, not the stored one: a compiled leaf id
+ // (`prio-<tier>-<slug>`) does not exist in the stored tree, and looking it
+ // up there would silently fall back to the lane's whole operation union.
+ const leaf = flattenLeaves(dispatchRoot).find(
(l) => l.id === picked.pick.leafId,
);
const wanted = new Set(
diff --git a/common/jobs/channelPriorityCompile.test.ts b/common/jobs/channelPriorityCompile.test.ts
@@ -0,0 +1,275 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import {
+ channelsForOperation,
+ compileLaneRoot,
+ focusSummary,
+ isChannelPaused,
+ prioLeafId,
+ sanitizeChannelPriority,
+} from "../lib/channelPriority";
+import type { ChannelPriority } from "../lib/channelPriority";
+import {
+ type AutoQueueGroup,
+ type ChannelWork,
+ buildPendingByLeaf,
+ emptyAutoQueueRuntime,
+ selectNextWork,
+} from "./autoQueuePolicy";
+
+// THE DESIGN, ASSERTED THROUGH THE REAL ENGINE.
+//
+// `plans/channel-priority.md`'s decision is COMPILE, not consult: the model
+// becomes an ordinary `AutoQueueGroup` and dispatch runs the engine it always
+// ran. So the claim "a focus holds the rest until its work is done" is not a
+// claim about `channelPriority.ts` at all — it is a claim about what
+// `buildPendingByLeaf` + `selectNextWork` do to a COMPILED tree, and it can
+// only be tested by driving both.
+//
+// IN jobs/, NOT lib/, for the reason `laneMigration.test.ts` is: this file must
+// import `jobs/autoQueuePolicy`, and `../architecture.test.ts` forbids
+// `lib/ -> jobs/` from a test file like any other. It must also stay out of the
+// controller layer's reach — `jobs/ -> controller/` is forbidden too — which is
+// why the runner-side halves (the `listChannelMeta` filter as the runner calls
+// it, and the once-per-transition log line) are asserted in
+// `controller/autoRunner.test.ts` instead.
+//
+// The sanitizer round-trip half lives in `channelPrioritySanitize.test.ts` (S0).
+
+// One bucket, spelled once. Which bucket a lane draws is `defaultDrawsForPolicy`'s
+// business and autoQueuePolicy.test.ts's; nothing here depends on the name.
+const BUCKET = "downloadedNoTranscript";
+
+function work(slug: string, ids: string[]): ChannelWork {
+ return { slug, platform: "youtube", buckets: { [BUCKET]: [...ids] } };
+}
+
+function model(value: unknown): ChannelPriority {
+ return sanitizeChannelPriority(value);
+}
+
+// EXACTLY WHAT THE RUNNER DOES, in the runner's order: `listChannelMeta` filters
+// the channel list by the lane's effective tier, then `laneDispatchRoot`
+// compiles the root from what survived. Both halves are here so a test that
+// says "paused is never drawn" is testing the pair the runner actually runs,
+// catch-all included.
+function dispatch(
+ lane: "transcription" | "download" | "digest" | "backfill",
+ priority: ChannelPriority,
+ channels: ChannelWork[],
+ focusSlugs: readonly string[] = [],
+): { root: AutoQueueGroup; channels: ChannelWork[] } {
+ const slugs = channelsForOperation(
+ priority,
+ channels.map((c) => c.slug),
+ lane,
+ );
+ const live = channels.filter((c) => slugs.includes(c.slug));
+ return {
+ root: compileLaneRoot(lane, priority, slugs, focusSlugs),
+ channels: live,
+ };
+}
+
+// Drain a compiled tree the way the runner does: re-ask the policy after every
+// grant, consuming the chosen id. Returns the owning channel of each pick, in
+// served order — the focus question is "whose video came next", not "which
+// leaf".
+function drainOwners(
+ root: AutoQueueGroup,
+ pending: Record<string, string[]>,
+ max = 100,
+): string[] {
+ const runtime = emptyAutoQueueRuntime();
+ const owners: string[] = [];
+ for (let i = 0; i < max; i++) {
+ const pick = selectNextWork(root, pending, runtime);
+ if (!pick) break;
+ owners.push(pick.leafId);
+ const ids = pending[pick.leafId];
+ assert.equal(ids[0], pick.videoId, "pick should be head of leaf queue");
+ ids.shift();
+ }
+ return owners;
+}
+
+test("a focus holds a non-focus channel with work, then releases it when the focus is exhausted", () => {
+ const priority = model({ focus: { kind: "channels", slugs: ["slow-a"] } });
+ const channels = [work("slow-a", ["a1", "a2"]), work("slow-b", ["b1", "b2"])];
+ const { root, channels: live } = dispatch(
+ "transcription",
+ priority,
+ channels,
+ ["slow-a"],
+ );
+ const pending = buildPendingByLeaf(root, live, [BUCKET]);
+
+ // Both channels have work RIGHT NOW — the hold is not "b has nothing".
+ assert.deepEqual(pending[prioLeafId("focus", "slow-a")], ["a1", "a2"]);
+ assert.deepEqual(pending[prioLeafId("normal", "slow-b")], ["b1", "b2"]);
+
+ // Strict descent serves every one of A's before it ever reaches B's group.
+ assert.deepEqual(drainOwners(root, pending), [
+ prioLeafId("focus", "slow-a"),
+ prioLeafId("focus", "slow-a"),
+ prioLeafId("normal", "slow-b"),
+ prioLeafId("normal", "slow-b"),
+ ]);
+});
+
+test("new focus work retakes the lane on the very next pick", () => {
+ const priority = model({ focus: { kind: "channels", slugs: ["slow-a"] } });
+ const channels = [work("slow-a", []), work("slow-b", ["b1", "b2"])];
+ const { root, channels: live } = dispatch(
+ "transcription",
+ priority,
+ channels,
+ ["slow-a"],
+ );
+ const pending = buildPendingByLeaf(root, live, [BUCKET]);
+ const runtime = emptyAutoQueueRuntime();
+
+ // Mid-"slow-b batch": the focus group is empty, so strict descent skips it.
+ let pick = selectNextWork(root, pending, runtime);
+ assert.equal(pick?.leafId, prioLeafId("normal", "slow-b"));
+ pending[pick!.leafId].shift();
+
+ // A snapshot regen gives the focus channel a video. The runner rebuilds
+ // `pending` every tick and re-asks, so the very next pick is the focus's —
+ // no second mechanism, no preemption of the unit already running.
+ pending[prioLeafId("focus", "slow-a")].push("a1");
+ pick = selectNextWork(root, pending, runtime);
+ assert.equal(pick?.leafId, prioLeafId("focus", "slow-a"));
+});
+
+test("a focus that resolves to nothing compiles no focus group and holds no one", () => {
+ // An unknown siteId is the live case: `resolveFocusSlugs` answers [], so a
+ // typo must leave the tree exactly as it would be with no focus at all.
+ const priority = model({ focus: { kind: "site", siteId: "no-such-site" } });
+ const channels = [work("slow-a", ["a1"]), work("slow-b", ["b1"])];
+ const { root, channels: live } = dispatch(
+ "transcription",
+ priority,
+ channels,
+ [],
+ );
+ assert.equal(
+ root.children.some((c) => c.id === "prio-focus"),
+ false,
+ );
+ const pending = buildPendingByLeaf(root, live, [BUCKET]);
+ assert.deepEqual(drainOwners(root, pending).sort(), [
+ prioLeafId("normal", "slow-a"),
+ prioLeafId("normal", "slow-b"),
+ ]);
+});
+
+test("a paused channel is drawn by no leaf, the catch-all included", () => {
+ const priority = model({ channels: { gone: { tier: "paused" } } });
+ const channels = [work("kept", ["k1"]), work("gone", ["g1", "g2"])];
+ const { root, channels: live } = dispatch("download", priority, channels);
+
+ // The FILTER is what does it: a tree cannot express exclusion, so the paused
+ // channel never reaches the projection at all.
+ assert.deepEqual(
+ live.map((c) => c.slug),
+ ["kept"],
+ );
+ const pending = buildPendingByLeaf(root, live, [BUCKET]);
+ const everyId = Object.values(pending).flat();
+ assert.deepEqual(everyId, ["k1"]);
+ assert.deepEqual(pending["prio-all"], []);
+ assert.equal(
+ Object.keys(pending).includes(prioLeafId("paused", "gone")),
+ false,
+ );
+});
+
+test("sync-only pause leaves the download lane drawing the channel", () => {
+ // What all 15 live `excludeFromSync` channels migrate to. The base tier and
+ // the rank stand; only the `sync` operation loses the channel, so no lane's
+ // membership moves — which is what makes that migration lossless.
+ const priority = model({
+ channels: {
+ omnivods: { tier: "normal", rank: 7, overrides: { sync: "paused" } },
+ },
+ });
+ assert.equal(isChannelPaused(priority, "omnivods", "sync"), true);
+ assert.equal(isChannelPaused(priority, "omnivods", "download"), false);
+
+ const channels = [work("omnivods", ["o1"])];
+ const { root, channels: live } = dispatch("download", priority, channels);
+ const pending = buildPendingByLeaf(root, live, [BUCKET]);
+ assert.deepEqual(pending[prioLeafId("normal", "omnivods")], ["o1"]);
+ assert.deepEqual(channelsForOperation(priority, ["omnivods"], "sync"), []);
+});
+
+test("a per-lane override pauses one lane and leaves the others drawing", () => {
+ const priority = model({
+ channels: { noisy: { tier: "normal", overrides: { download: "paused" } } },
+ });
+ const channels = [work("noisy", ["n1"])];
+
+ const dl = dispatch("download", priority, channels);
+ assert.deepEqual(dl.channels, []);
+ assert.deepEqual(
+ Object.values(buildPendingByLeaf(dl.root, dl.channels, [BUCKET])).flat(),
+ [],
+ );
+
+ const tr = dispatch("transcription", priority, channels);
+ const pending = buildPendingByLeaf(tr.root, tr.channels, [BUCKET]);
+ assert.deepEqual(pending[prioLeafId("normal", "noisy")], ["n1"]);
+});
+
+test("a focused channel paused for that lane is not drawn by it", () => {
+ // Focus wins over the stored tier; paused wins over focus, per operation.
+ const priority = model({
+ focus: { kind: "channels", slugs: ["star"] },
+ channels: { star: { tier: "normal", overrides: { digest: "paused" } } },
+ });
+ const channels = [work("star", ["s1"]), work("other", ["o1"])];
+
+ const digest = dispatch("digest", priority, channels, ["star"]);
+ assert.deepEqual(
+ digest.channels.map((c) => c.slug),
+ ["other"],
+ );
+ assert.equal(
+ digest.root.children.some((c) => c.id === "prio-focus"),
+ false,
+ );
+
+ const backfill = dispatch("backfill", priority, channels, ["star"]);
+ const pending = buildPendingByLeaf(backfill.root, backfill.channels, [
+ BUCKET,
+ ]);
+ assert.deepEqual(pending[prioLeafId("focus", "star")], ["s1"]);
+});
+
+test("focusSummary reads the hold off the same pending map the pick used", () => {
+ const priority = model({ focus: { kind: "channels", slugs: ["slow-a"] } });
+ const channels = [work("slow-a", ["a1"]), work("slow-b", ["b1", "b2"])];
+ const { root, channels: live } = dispatch(
+ "digest",
+ priority,
+ channels,
+ ["slow-a"],
+ );
+ const pending = buildPendingByLeaf(root, live, [BUCKET]);
+
+ const holding = focusSummary(priority, ["slow-a"], pending);
+ assert.equal(holding.active, true);
+ assert.equal(holding.holding, true);
+ assert.equal(holding.focusPending, 1);
+ assert.equal(holding.otherPending, 2);
+ assert.equal(holding.heldChannels, 1);
+
+ // Drain the focus's only unit: the lane is released, and the summary says so
+ // off the same map — no second read, no new clock.
+ pending[prioLeafId("focus", "slow-a")].shift();
+ const released = focusSummary(priority, ["slow-a"], pending);
+ assert.equal(released.holding, false);
+ assert.equal(released.heldChannels, 0);
+ assert.equal(released.otherPending, 2);
+});