Archilyzer · Source

archilyzer

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

commit f424e881bb48c44393cd72270e8b72d0dbe8c0eb
parent 67c7f85de076c0fc4c5070e4deb57ac8380a6626
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Mon, 14 Sep 2026 17:15:01 -0400

views: the auto-queue console reads its state once, at the edge

The four-lane console's payload was three files that each did their own
reading, and the cost of that was visible on the wire: `buildKind` ran once per
lane and opened with getPaths()/getSettings(), so `readAutoQueueState(paths)`
parsed the SAME document FOUR TIMES on a three-second poll (status.ts:122).
Nothing needed it four times; the function simply had no way to be handed it.

Three moves, one shape:

  * `common/views/autoQueueLanes.ts` takes { settings, briefs, workerSummary }
    and is sync. It no longer resolves the channel briefs or the worker pool —
    the status shell has both already.
  * `common/views/channelPriority.ts` is the pure half of the priority view:
    the `PriorityView` type and the two folds over it. `readPriorityView` stays
    in the editor, because resolving one is a channel listing and, for a site
    focus, a sites read.
  * `common/views/autoQueueStatus.ts` takes the settings, the pool, the clock,
    the priority view, the per-lane runner status, the auto-queue document, the
    per-lane pending counts and the rail — and is sync. `getAutoRunnerStatus`
    CONSTRUCTS a runner, which is the sharpest reason a status poll's fold must
    not do its own reading.

`editor/app/operations/status.ts` is the shell that does all of it: one
priority resolution, ONE `readAutoQueueState`, four `computeLeafPending`, four
`getAutoRunnerStatus`, the briefs off the per-request cache, and the lanes
builder folded in beside them. Same payload type names, same wire fields, same
route. `operations/lanes.ts` is three type re-exports until sub-slice D
repoints its two client importers.

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

Diffstat:
Acommon/views/autoQueueLanes.test.ts | 102+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/views/autoQueueLanes.ts | 139+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/views/autoQueueStatus.test.ts | 179+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/views/autoQueueStatus.ts | 181+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/views/channelPriority.ts | 54++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/app/operations/channelPriorityView.ts | 53++++++++++-------------------------------------------
Meditor/app/operations/lanes.ts | 140++++++-------------------------------------------------------------------------
Meditor/app/operations/status.ts | 215++++++++++++++++++++-----------------------------------------------------------
8 files changed, 729 insertions(+), 334 deletions(-)

diff --git a/common/views/autoQueueLanes.test.ts b/common/views/autoQueueLanes.test.ts @@ -0,0 +1,102 @@ +import test from "node:test"; +import assert from "node:assert/strict"; +import { defaultSiteSettings } from "../lib/settings"; +import type { WorkerSummary } from "../jobs/workerPool"; +import { buildAutoQueueLanes } from "./autoQueueLanes"; + +// WHERE A LANE'S WORK CAN RUN, and the three rules that decide it. The rows are +// read straight off `workerMatches`, the same predicate the dispatchers grant +// by, so what this test really pins is the FOLD around it: one row per +// configured worker, not one per slot and not one per (operation, slot) pair. + +const worker = (over: Partial<WorkerSummary> & { id: string }): WorkerSummary => ({ + name: over.id, + kind: "remote", + priority: 0, + busy: false, + state: "enabled", + degraded: false, + enabled: true, + ...over, +}); + +// The default corpus has exactly one lane operation — `digest`, on the digest +// lane — which is all these rules need. +const lanes = (workerSummary: WorkerSummary[]) => + buildAutoQueueLanes({ + settings: defaultSiteSettings(), + briefs: [], + workerSummary, + }); + +test("the digest lane names its operation, and the empty lanes stay empty", () => { + const payload = lanes([]); + assert.deepEqual( + payload.operations.digest.map((o) => o.id), + ["digest"], + ); + assert.deepEqual(payload.operations.transcription, []); + assert.deepEqual(payload.runsOn.digest, []); + // The rail is the CATALOG view: one band per catalogued operation whatever + // queue it runs on, on the two media bands the pipeline always has. + assert.deepEqual( + payload.bands.map((b) => b.id), + ["download", "transcription", "digest"], + ); +}); + +test("a worker's slots fold back to one row", () => { + // Slot expansion writes `gpu#1` / `gpu#2` with names to match. They are ONE + // machine with two slots: two rows would advertise two workers, and counting + // the slots per operation would advertise capacity that does not exist. + const payload = lanes([ + worker({ id: "gpu#1", name: "GPU box #1", tags: ["digest"] }), + worker({ id: "gpu#2", name: "GPU box #2", tags: ["digest"] }), + ]); + assert.deepEqual(payload.runsOn.digest, [ + { id: "gpu", name: "GPU box", kind: "remote", slots: 2, available: 2 }, + ]); +}); + +test("only slots that could take work right now are available", () => { + const payload = lanes([ + worker({ id: "gpu#1", name: "GPU box #1", tags: ["digest"] }), + worker({ id: "gpu#2", name: "GPU box #2", tags: ["digest"], degraded: true }), + worker({ id: "gpu#3", name: "GPU box #3", tags: ["digest"], state: "disabled" }), + ]); + assert.deepEqual(payload.runsOn.digest, [ + { id: "gpu", name: "GPU box", kind: "remote", slots: 3, available: 1 }, + ]); +}); + +test("an untagged remote worker is not listed", () => { + // Tagging IS the unit-dispatch opt-in for a remote worker: untagged means + // "transcriptions only", and `workerMatches` says true for an untagged remote + // because that predicate is also asked about transcription leases. The lane + // rows have to apply the opt-in themselves. + assert.deepEqual(lanes([worker({ id: "plain", name: "plain" })]).runsOn.digest, []); + assert.deepEqual( + lanes([worker({ id: "other", name: "other", tags: ["diarization"] })]).runsOn + .digest, + [], + ); +}); + +test("an llm worker is listed by tag", () => { + const payload = lanes([ + worker({ id: "ollama", name: "ollama", kind: "llm", tags: ["digest"] }), + worker({ id: "elsewhere", name: "elsewhere", kind: "llm", tags: ["attribution-text"] }), + ]); + assert.deepEqual(payload.runsOn.digest, [ + { id: "ollama", name: "ollama", kind: "llm", slots: 1, available: 1 }, + ]); +}); + +test("a local worker is never listed", () => { + // The lane always runs here; saying so would be noise on every row. + assert.deepEqual( + lanes([worker({ id: "local", name: "local", kind: "local", tags: ["digest"] })]) + .runsOn.digest, + [], + ); +}); diff --git a/common/views/autoQueueLanes.ts b/common/views/autoQueueLanes.ts @@ -0,0 +1,139 @@ +import type { SiteSettings } from "../lib/settings"; +import { LANES, type AutoQueueKind } from "../lib/autoQueueTypes"; +import { + allOperations, + operationLabel, + operationsForLane, +} from "../lib/operations"; +import type { WorkerSummary } from "../jobs/workerPool"; +import type { ChannelSnapshot } from "../controller/channelSnapshot"; +import { workerMatches } from "../lib/workers"; +import { + buildOperationBands, + type OperationBand, +} from "./pipeline/buildBands"; + +// WHAT THE CONSOLE KNOWS ABOUT A LANE THAT IS NOT ITS RUNNER'S STATE. +// +// This file used to be the two SWEEP-fed lanes: a state, a reason, a scope, a +// plan and an arm button apiece, for the two pipelines the auto-queue runner +// could not dispatch. It can dispatch them now (slice 1.2) and the sweeps are +// gone (1.3), so a lane's state, its work list and its controls all come from +// `/api/auto-queue/status` like the other two lanes'. +// +// TWO THINGS ARE LEFT, and both are properties of the REGISTRY rather than of a +// runner, which is why they are still built here: +// +// * the BANDS — one per catalogued operation, its five populations off the +// channel snapshots. The rail draws them, and they are per OPERATION where +// everything else on this page is per LANE. +// * WHICH OPERATIONS EACH LANE DISPATCHES, with their labels and the delegate +// workers that can take them. Settings-dependent (`digest.remoteEnabled` +// moves the digest operation between two queue keys that are both the +// digest lane), so it is a server read, and the console is a client +// component. + +export type LaneOperationView = { + id: string; + label: string; +}; + +export type LaneWorker = { + // The configured worker's id (slot-expansion suffixes folded back together). + id: string; + name: string; + kind: "remote" | "llm"; + slots: number; + // Slots currently enabled and not degraded — 0 across every delegate is the + // "configured but nothing can take it right now" warning state. + available: number; +}; + +export type AutoQueueLanesPayload = { + // One band per pipeline, in rail order. See components/pipelines/buildBands.ts. + bands: OperationBand[]; + // The operations each lane dispatches, in dependency order. Empty for the two + // bucket lanes, whose work list is a snapshot bucket rather than an operation. + operations: Record<AutoQueueKind, LaneOperationView[]>; + // WHERE A LANE'S WORK CAN RUN besides this machine, folded to one row per + // configured worker with its slot count. READ-ONLY and never part of any + // policy — delegation is a worker-pool concern decided per item at dispatch, + // and the console only reports it. Derived through the SAME workerMatches rule + // the dispatchers grant by, so this row and the actual routing cannot + // disagree. + runsOn: Record<AutoQueueKind, LaneWorker[]>; +}; + +export type AutoQueueLanesInputs = { + settings: SiteSettings; + // Only the snapshots are read — the bands are the five populations each + // channel's report already carries. Typed as the half that is used so a + // caller need not have a full brief to ask. + briefs: ReadonlyArray<{ snapshot: ChannelSnapshot | null }>; + // `pool.summary()`, taken once by the shell rather than per lane. + workerSummary: readonly WorkerSummary[]; +}; + +export function buildAutoQueueLanes( + inputs: AutoQueueLanesInputs, +): AutoQueueLanesPayload { + const { settings, briefs, workerSummary } = inputs; + // allOperations, not one lane's members: the rail is the CATALOG view — every + // operation that is switched on, whatever queue it runs on. + const kinds = allOperations(settings); + + const bands = buildOperationBands({ + snapshots: briefs.map((c) => c.snapshot ?? null), + operationIds: kinds.map((k) => k.id), + }); + + const contendsForOf = new Map( + kinds.map((k) => [k.id, (k.laneFor?.(settings) ?? k.lane).contendsFor]), + ); + // An llm worker serves an operation its tags (or the llm default set) name; a + // remote worker takes units only when TAGGED (the unit-dispatch opt-in) and + // its tags intersect [operation, contended resource]. Local workers are not + // listed — the lane always runs here, and saying so would be noise. + // ONE PASS PER WORKER SLOT, not per (operation, slot) pair. A worker serving + // two of the lane's operations is one row with one slot count — the slots are + // the machine's, not the operation's, and counting them twice would advertise + // capacity that does not exist. + const runsOnFor = (ids: readonly string[]): LaneWorker[] => { + const byBase = new Map<string, LaneWorker>(); + for (const w of workerSummary) { + const takesAny = ids.some((id) => + w.kind === "llm" + ? workerMatches({ kind: "llm", tags: w.tags }, [id]) + : w.kind === "remote" && + (w.tags?.length ?? 0) > 0 && + workerMatches({ kind: "remote", tags: w.tags }, [ + id, + contendsForOf.get(id) ?? "cpu", + ]), + ); + if (!takesAny) continue; + const base = w.id.split("#")[0]; + const entry = byBase.get(base) ?? { + id: base, + name: w.name.replace(/ #\d+$/, ""), + kind: w.kind as "remote" | "llm", + slots: 0, + available: 0, + }; + entry.slots++; + if (w.state === "enabled" && !w.degraded) entry.available++; + byBase.set(base, entry); + } + return [...byBase.values()]; + }; + + const operations = {} as Record<AutoQueueKind, LaneOperationView[]>; + const runsOn = {} as Record<AutoQueueKind, LaneWorker[]>; + for (const lane of LANES) { + const ids = operationsForLane(lane, settings).map((op) => op.id); + operations[lane] = ids.map((id) => ({ id, label: operationLabel(id) })); + runsOn[lane] = runsOnFor(ids); + } + + return { bands, operations, runsOn }; +} diff --git a/common/views/autoQueueStatus.test.ts b/common/views/autoQueueStatus.test.ts @@ -0,0 +1,179 @@ +import test from "node:test"; +import assert from "node:assert/strict"; +import { defaultSiteSettings, type SiteSettings } from "../lib/settings"; +import { defaultChannelPriority, compileLaneRoot } from "../lib/channelPriority"; +import { LANES, type AutoQueueKind } from "../lib/autoQueueTypes"; +import { emptyAutoQueueState } from "../jobs/autoQueueState"; +import type { AutoRunnerStatus, LeafPending } from "../controller/autoRunner"; +import { buildAutoQueueLanes } from "./autoQueueLanes"; +import type { PriorityView } from "./channelPriority"; +import { + buildAutoQueueStatusPayload, + type AutoQueueStatusInputs, +} from "./autoQueueStatus"; + +// THE FOUR LANES' CONSOLE, as a fold over state that is handed to it. Every +// reader this used to do itself — the auto-queue document, the pending counts, +// the runner status, the clock — is an argument now, which is what makes the +// three rules below testable at all. + +const runnerFor = (kind: AutoQueueKind): AutoRunnerStatus => ({ + kind, + running: false, + jobId: null, + startedAt: null, + inFlight: [], + activeByNode: {}, + idleReason: "stopped", +}); + +const pendingFor = (counts: Record<string, number> = {}): LeafPending => ({ + counts, + head: {}, + owner: {}, + recency: {}, + nextUp: null, +}); + +const byLane = <T,>(make: (kind: AutoQueueKind) => T): Record<AutoQueueKind, T> => + Object.fromEntries(LANES.map((l) => [l, make(l)])) as Record<AutoQueueKind, T>; + +const NOW = 1_700_000_000_000; + +function inputs(over: Partial<AutoQueueStatusInputs> = {}): AutoQueueStatusInputs { + const settings = defaultSiteSettings(); + const priority: PriorityView = { + model: defaultChannelPriority(), + compiled: false, + slugs: [], + focusSlugs: [], + name: null, + }; + return { + settings, + pool: { + summary: () => [], + isPaused: () => false, + canStopPartial: () => false, + }, + now: () => NOW, + priority, + runner: byLane(runnerFor), + state: emptyAutoQueueState(), + pending: byLane(() => pendingFor()), + lanes: buildAutoQueueLanes({ settings, briefs: [], workerSummary: [] }), + ...over, + }; +} + +test("every lane in LANES gets an entry, and the rail rides along", () => { + const i = inputs(); + const payload = buildAutoQueueStatusPayload(i); + for (const lane of LANES) assert.equal(payload[lane].kind, lane); + assert.equal(payload.lanes, i.lanes); +}); + +test("a cooldown is filtered by the injected clock, newest first", () => { + const state = emptyAutoQueueState(); + state.download.platformBackoff = { + // Lapsed a second ago: the operator is not waiting on it and it must not + // show as a reason the lane is quiet. + youtube: { until: NOW - 1_000, fails: 9 }, + rumble: { until: NOW + 60_000, fails: 2 }, + twitch: { until: NOW + 600_000, fails: 5 }, + // The boundary is strictly greater-than: a cooldown lapsing exactly now has + // lapsed. + kick: { until: NOW, fails: 1 }, + }; + const payload = buildAutoQueueStatusPayload(inputs({ state })); + assert.deepEqual(payload.download.cooldowns, [ + { platform: "twitch", untilMs: NOW + 600_000, fails: 5 }, + { platform: "rumble", untilMs: NOW + 60_000, fails: 2 }, + ]); + // The clock is a function, so the same state read later says something else. + assert.deepEqual( + buildAutoQueueStatusPayload(inputs({ state, now: () => NOW + 120_000 })) + .download.cooldowns, + [{ platform: "twitch", untilMs: NOW + 600_000, fails: 5 }], + ); +}); + +test("transcription's hold is the pool's, every other lane's is the gate", () => { + // The asymmetry is deliberate and is the reason `pool` is on these inputs at + // all: transcription's hold is LIVE on the worker pool, while every other + // lane's stored flag IS its gate. + const paused = buildAutoQueueStatusPayload( + inputs({ + pool: { + summary: () => [], + isPaused: () => true, + canStopPartial: () => false, + }, + }), + ); + assert.equal(paused.transcription.held, true); + assert.equal(paused.download.held, false); + assert.equal(paused.digest.held, false); + + const settings: SiteSettings = defaultSiteSettings(); + settings.autoQueue = { + ...settings.autoQueue, + download: { ...settings.autoQueue.download, held: true }, + transcription: { ...settings.autoQueue.transcription, held: true }, + }; + const gated = buildAutoQueueStatusPayload(inputs({ settings })); + assert.equal(gated.download.held, true); + // Settings say held; the live pool says otherwise, and the pool wins. + assert.equal(gated.transcription.held, false); +}); + +test("policy.root is the tree the lane DISPATCHES from, stored or compiled", () => { + const stored = buildAutoQueueStatusPayload(inputs()); + const settings = defaultSiteSettings(); + assert.equal(stored.digest.policyCompiled, false); + assert.deepEqual(stored.digest.policy.root, settings.autoQueue.digest.root); + // Everything else on the policy comes through the spread untouched. + assert.equal(stored.digest.policy.enabled, settings.autoQueue.digest.enabled); + assert.equal(stored.digest.policy.order, settings.autoQueue.digest.order); + + const model = defaultChannelPriority(); + model.channels = { alpha: { tier: "low" } }; + const priority: PriorityView = { + model, + compiled: true, + slugs: ["alpha", "beta"], + focusSlugs: ["beta"], + name: "beta", + }; + const compiled = buildAutoQueueStatusPayload(inputs({ priority })); + assert.equal(compiled.digest.policyCompiled, true); + assert.deepEqual( + compiled.digest.policy.root, + compileLaneRoot("digest", model, priority.slugs, priority.focusSlugs), + ); + assert.notDeepEqual(compiled.digest.policy.root, settings.autoQueue.digest.root); + assert.equal(compiled.digest.focusName, "beta"); +}); + +test("the pending fold is per lane and reaches the focus banner", () => { + const pending = byLane(() => pendingFor()); + pending.digest = pendingFor({ "prio-focus-beta": 3, "prio-normal": 7 }); + const model = defaultChannelPriority(); + const priority: PriorityView = { + model: { ...model, focus: { kind: "channels", slugs: ["beta"] } }, + compiled: true, + slugs: ["alpha", "beta"], + focusSlugs: ["beta"], + name: "beta", + }; + const payload = buildAutoQueueStatusPayload(inputs({ pending, priority })); + assert.deepEqual(payload.digest.pendingByLeaf, { + "prio-focus-beta": 3, + "prio-normal": 7, + }); + assert.equal(payload.digest.focus.active, true); + assert.equal(payload.digest.focus.focusPending, 3); + assert.equal(payload.digest.focus.otherPending, 7); + // A lane with nothing pending is still an entry, with the same shape. + assert.deepEqual(payload.download.pendingByLeaf, {}); +}); diff --git a/common/views/autoQueueStatus.ts b/common/views/autoQueueStatus.ts @@ -0,0 +1,181 @@ +import type { LiveInputs } from "./inputs"; +import { + type AutoRunnerStatus, + type LeafPending, + type NextUpView, + type RecencyKeyView, + laneDispatchRoot, +} from "../controller/autoRunner"; +import type { + AutoQueueKind, + AutoQueuePick, + AutoQueueState, +} from "../jobs/autoQueueState"; +import { LANES } from "../lib/autoQueueTypes"; +import type { AutoQueuePolicy } from "../jobs/autoQueuePolicy"; +import { isGateHeld } from "../lib/pauseGates"; +import type { FocusSummary } from "../lib/channelPriority"; +import type { AutoQueueLanesPayload } from "./autoQueueLanes"; +import { type PriorityView, laneFocusSummary } from "./channelPriority"; + +// Read-only payload for the Auto-Queue panel: per-kind runner status (running?, +// what's in flight, in-flight counts per tree node), the effective policy, the +// recent pick log, and snapshot-derived pending counts per leaf — so the UI can +// show "cornbreadman: 12 pending" and which leaf is currently being serviced. +// Backs a passive UI poll, like /api/scheduler/status. + +// A platform currently in a rate-limit (429) / network cooldown. `untilMs` is +// the epoch ms the cooldown lapses; the client renders remaining time itself so +// it counts down live. Only the download kind ever has these. +export type PlatformCooldownView = { + platform: string; + untilMs: number; + fails: number; +}; + +export type AutoQueueKindStatus = { + kind: AutoQueueKind; + // THE POLICY AS THE LANE DISPATCHES IT, not as it is stored. Every field is + // the stored one except `root`, which is the COMPILED tree whenever + // `settings.channelPriority` says anything (see ./channelPriority.ts). The + // claim ladder is drawn from this, so a rung is a rule the runner actually + // has — and with a priority model set the stored tree is not one. + policy: AutoQueuePolicy; + runner: AutoRunnerStatus; + pendingByLeaf: Record<string, number>; + // The first few pending ids of each rule, in the order the runner would hand + // them out — the drill-down behind each rung's count, and the only way to + // check an ordering change by eye. + pendingHeadByLeaf: Record<string, string[]>; + // videoId -> owning channel slug, for the drill-down's links. Covers exactly + // the ids in pendingHeadByLeaf (plus nextUp's). + ownerByVideo: Record<string, string>; + // videoId -> recency sort key, for the same ids, when the policy orders by + // recency. Absent entries just render without a date. + recencyByVideo: Record<string, RecencyKeyView>; + // What the policy would dispatch next, and which rules it passed over. + nextUp: NextUpView | null; + picks: AutoQueuePick[]; + // Platforms paused by a rate-limit/network backoff (download kind only). + cooldowns: PlatformCooldownView[]; + // This runner's lane gate — what the page's pause button draws, and what lets + // the operations rail show a runner HOLDING for the first time. + // + // THE TWO KINDS ANSWER IT FROM DIFFERENT PLACES, on purpose. Transcription's + // hold is LIVE, on the worker pool; settings.transcriptionsPaused is only what + // the boot hook re-applies after a restart, and the e2e harness rewrites + // settings wholesale between tests while the pool keeps its pausedSnapshot. + // Download has no live counterpart: its flag IS the gate, read at dispatch. + held: boolean; + // THE TREE ABOVE IS GENERATED, so the editor for it is read-only and the + // channel priorities on /channels are where it is edited. False for a corpus + // with no priorities set, which is every corpus until one is. + policyCompiled: boolean; + // THE FOCUS BANNER'S NUMBERS, FOR THIS LANE. Always present and inert when + // `active` is false, so the banner is one component with one early return + // rather than a conditional on the payload. The counts are this lane's own — + // `focusPending`/`otherPending` differ per lane by construction. + focus: FocusSummary; + // The focus's display name, resolved on the server (the model stores a + // siteId, and the banner does no I/O). Null when nothing is focused. + focusName: string | null; +}; + +export type AutoQueueStatusPayload = Record< + AutoQueueKind, + AutoQueueKindStatus +> & { + // The sweep half of the digest and backfill pipelines, plus the comparison + // rail's bands. On the SAME payload as the runners rather than a second + // endpoint, because the rail's whole purpose is that four lanes are read + // together — two polls would let the rail and the focused lane disagree about + // the same moment. + // + // Every lane now has a `status[lane]` entry too. The two are not duplicates: + // this is what the SWEEPS are doing (armed, scoped, planned), and the entry + // is what the lane's TREE says. Slice 1.3 deletes this half. + lanes: AutoQueueLanesPayload; +}; + +export type AutoQueueStatusInputs = Pick< + LiveInputs, + "settings" | "pool" | "now" +> & { + // Resolved ONCE for all four lanes by the shell: a focus set costs a channel + // listing and, for a site focus, a sites read, and the four lanes compile from + // the same one. + priority: PriorityView; + // Per-lane runner status. Injected because asking for it CONSTRUCTS the + // runner, which a status poll must never do. + runner: Record<AutoQueueKind, AutoRunnerStatus>; + // The auto-queue state document, read once for all four lanes (it used to be + // read once per lane, four parses of the same file per poll). + state: AutoQueueState; + // Per-lane snapshot-derived pending work. + pending: Record<AutoQueueKind, LeafPending>; + lanes: AutoQueueLanesPayload; +}; + +function buildKind( + kind: AutoQueueKind, + inputs: AutoQueueStatusInputs, +): AutoQueueKindStatus { + const { settings, pool, priority } = inputs; + const stored = settings.autoQueue[kind]; + // The stored policy, with the DISPATCHED root in place of the stored one. + // Spread rather than rebuilt, so `held`, `snoozeUntil`, `enabled`, `order` + // and `maxWorkers` come through untouched — the same rule every writer of a + // policy in this repo follows. + const policy: AutoQueuePolicy = { + ...stored, + root: laneDispatchRoot(kind, stored, priority, priority.slugs), + }; + const runner = inputs.runner[kind]; + const pending = inputs.pending[kind]; + const now = inputs.now(); + // Surface platforms still inside their cooldown window so the operator can see + // why an otherwise-pending platform isn't being serviced (and a manual sync on + // it would be refused). A manual sync/download 429 writes into this same map. + const cooldowns: PlatformCooldownView[] = Object.entries( + inputs.state[kind].platformBackoff, + ) + .filter(([, e]) => e.until > now) + .map(([platform, e]) => ({ platform, untilMs: e.until, fails: e.fails })) + .sort((a, b) => b.untilMs - a.untilMs); + return { + kind, + policy, + runner, + pendingByLeaf: pending.counts, + pendingHeadByLeaf: pending.head, + ownerByVideo: pending.owner, + recencyByVideo: pending.recency, + nextUp: pending.nextUp, + picks: inputs.state[kind].picks, + cooldowns, + // THE ASYMMETRY IS TRANSCRIPTION'S ALONE, and it is not a special case for + // "the first lane": its hold is LIVE on the worker pool, while every other + // lane's flag IS its gate. See lib/pauseGates.ts. + held: kind === "transcription" ? pool.isPaused() : isGateHeld(settings, kind), + policyCompiled: priority.compiled, + // Keyed by COMPILED leaf id (`prio-focus-<slug>`), which is why this is + // computed here and not on the client: it is only meaningful against the + // counts of the tree the lane dispatches from. + focus: laneFocusSummary(priority, pending.counts), + focusName: priority.name, + }; +} + +// ONE ENTRY PER LANE, off LANES rather than off two hand-written names — so a +// lane added to the model appears on this payload with no edit here, which is +// the whole point of the widened type. +export function buildAutoQueueStatusPayload( + inputs: AutoQueueStatusInputs, +): AutoQueueStatusPayload { + return { + ...(Object.fromEntries( + LANES.map((lane) => [lane, buildKind(lane, inputs)]), + ) as Record<AutoQueueKind, AutoQueueKindStatus>), + lanes: inputs.lanes, + }; +} diff --git a/common/views/channelPriority.ts b/common/views/channelPriority.ts @@ -0,0 +1,54 @@ +import { + type ChannelPriority, + type FocusSummary, + type PendingByLeaf, + focusSummary, +} from "../lib/channelPriority"; + +// THE STATUS PAYLOAD'S HALF OF CHANNEL PRIORITY — the shape of "which tree does +// this lane dispatch from", and the focus banner's numbers off it. +// +// The RESOLUTION of that shape is a read (a channel listing, and for a site +// focus a sites read), so it stays in the editor as `readPriorityView`. What +// lives here is everything downstream of it: the view type the four lanes are +// compiled against, and the two pure folds the payload applies to it. + +export type PriorityView = { + model: ChannelPriority; + // The model says something, so the four trees are compiled rather than + // stored. The lane consoles read this to go read-only on the tree. + compiled: boolean; + // Every channel slug, the population `compileLaneRoot` filters per lane. + // Empty when nothing is compiled — it is never read in that case. + slugs: string[]; + focusSlugs: string[]; + // A display name for the focus: the site's title for a site focus, the slugs + // for a channel focus. Resolved by the reader because the model stores a + // siteId and the banner does no I/O. + name: string | null; +}; + +// `focusSummary` reads nothing but each leaf's COUNT, and `computeLeafPending` +// throws the arrays away before this layer sees them (it returns counts plus a +// truncated head). Rather than widen that return type — which would put the +// whole pending set of every leaf on a three-second poll's heap — the counts are +// re-presented as arrays of the right length. `new Array(n)` allocates no +// elements; only `.length` is ever read. +export function pendingByLeafFromCounts( + counts: Record<string, number>, +): PendingByLeaf { + const out: Record<string, readonly string[]> = {}; + for (const [id, n] of Object.entries(counts)) out[id] = new Array<string>(n); + return out; +} + +export function laneFocusSummary( + view: PriorityView, + counts: Record<string, number>, +): FocusSummary { + return focusSummary( + view.model, + view.focusSlugs, + pendingByLeafFromCounts(counts), + ); +} diff --git a/editor/app/operations/channelPriorityView.ts b/editor/app/operations/channelPriorityView.ts @@ -3,14 +3,11 @@ import { getSettings } from "yt-dlp-transcript-common/lib/settings"; import { listSites, siteChannelIndex } from "yt-dlp-transcript-common/lib/site"; import { listChannelConfigs } from "yt-dlp-transcript-common/controller/channels"; import { - type ChannelPriority, - type FocusSummary, - type PendingByLeaf, type SiteChannelIndex, - focusSummary, isDefaultChannelPriority, resolveFocusSlugs, } from "yt-dlp-transcript-common/lib/channelPriority"; +import type { PriorityView } from "yt-dlp-transcript-common/views/channelPriority"; // THE STATUS PAYLOAD'S HALF OF CHANNEL PRIORITY — the tree the lane actually // dispatches from, and the focus banner's numbers. @@ -35,20 +32,15 @@ import { // `listChannelConfigs` nor `listSites` is read — which is what keeps a corpus // with no priorities set on exactly today's payload and today's cost. -export type PriorityView = { - model: ChannelPriority; - // The model says something, so the four trees are compiled rather than - // stored. The lane consoles read this to go read-only on the tree. - compiled: boolean; - // Every channel slug, the population `compileLaneRoot` filters per lane. - // Empty when nothing is compiled — it is never read in that case. - slugs: string[]; - focusSlugs: string[]; - // A display name for the focus: the site's title for a site focus, the slugs - // for a channel focus. Resolved here because the model stores a siteId and - // the banner does no I/O. - name: string | null; -}; +// The pure half — the `PriorityView` type and the two folds the payload applies +// to it — is `common/views/channelPriority.ts`. What is left here is the READ +// that resolves one, which is the half that cannot move: it lists channels and, +// for a site focus, the sites directory. +export type { PriorityView }; +export { + pendingByLeafFromCounts, + laneFocusSummary, +} from "yt-dlp-transcript-common/views/channelPriority"; // "a, b and c" for a short channel focus; "a, b and 4 more" past three, because // this lands mid-sentence in a banner and a 30-slug list is not a name. @@ -92,28 +84,3 @@ export async function readPriorityView( if (model.focus.kind === "channels") name = channelFocusName(focusSlugs); return { model, compiled: true, slugs, focusSlugs, name }; } - -// `focusSummary` reads nothing but each leaf's COUNT, and `computeLeafPending` -// throws the arrays away before this layer sees them (it returns counts plus a -// truncated head). Rather than widen that return type — which would put the -// whole pending set of every leaf on a three-second poll's heap — the counts are -// re-presented as arrays of the right length. `new Array(n)` allocates no -// elements; only `.length` is ever read. -export function pendingByLeafFromCounts( - counts: Record<string, number>, -): PendingByLeaf { - const out: Record<string, readonly string[]> = {}; - for (const [id, n] of Object.entries(counts)) out[id] = new Array<string>(n); - return out; -} - -export function laneFocusSummary( - view: PriorityView, - counts: Record<string, number>, -): FocusSummary { - return focusSummary( - view.model, - view.focusSlugs, - pendingByLeafFromCounts(counts), - ); -} diff --git a/editor/app/operations/lanes.ts b/editor/app/operations/lanes.ts @@ -1,131 +1,11 @@ -import { getPaths } from "yt-dlp-transcript-common/lib/paths"; -import { getSettings } from "yt-dlp-transcript-common/lib/settings"; -import { LANES, type AutoQueueKind } from "yt-dlp-transcript-common/lib/autoQueueTypes"; -import { - allOperations, - operationLabel, - operationsForLane, -} from "yt-dlp-transcript-common/lib/operations"; -import { getWorkerPool } from "yt-dlp-transcript-common/jobs/workerPool"; -import { workerMatches } from "yt-dlp-transcript-common/lib/workers"; -import { getChannelBriefs } from "../lib/requestCache"; -import { - buildOperationBands, - type OperationBand, -} from "yt-dlp-transcript-common/views/pipeline/buildBands"; - -// WHAT THE CONSOLE KNOWS ABOUT A LANE THAT IS NOT ITS RUNNER'S STATE. +// The lane payload itself is `common/views/autoQueueLanes.ts`; the shell that +// reads its inputs is `./status.ts`, which needs the same channel briefs. // -// This file used to be the two SWEEP-fed lanes: a state, a reason, a scope, a -// plan and an arm button apiece, for the two pipelines the auto-queue runner -// could not dispatch. It can dispatch them now (slice 1.2) and the sweeps are -// gone (1.3), so a lane's state, its work list and its controls all come from -// `/api/auto-queue/status` like the other two lanes'. -// -// TWO THINGS ARE LEFT, and both are properties of the REGISTRY rather than of a -// runner, which is why they are still built here: -// -// * the BANDS — one per catalogued operation, its five populations off the -// channel snapshots. The rail draws them, and they are per OPERATION where -// everything else on this page is per LANE. -// * WHICH OPERATIONS EACH LANE DISPATCHES, with their labels and the delegate -// workers that can take them. Settings-dependent (`digest.remoteEnabled` -// moves the digest operation between two queue keys that are both the -// digest lane), so it is a server read, and the console is a client -// component. - -export type LaneOperationView = { - id: string; - label: string; -}; - -export type LaneWorker = { - // The configured worker's id (slot-expansion suffixes folded back together). - id: string; - name: string; - kind: "remote" | "llm"; - slots: number; - // Slots currently enabled and not degraded — 0 across every delegate is the - // "configured but nothing can take it right now" warning state. - available: number; -}; - -export type AutoQueueLanesPayload = { - // One band per pipeline, in rail order. See components/pipelines/buildBands.ts. - bands: OperationBand[]; - // The operations each lane dispatches, in dependency order. Empty for the two - // bucket lanes, whose work list is a snapshot bucket rather than an operation. - operations: Record<AutoQueueKind, LaneOperationView[]>; - // WHERE A LANE'S WORK CAN RUN besides this machine, folded to one row per - // configured worker with its slot count. READ-ONLY and never part of any - // policy — delegation is a worker-pool concern decided per item at dispatch, - // and the console only reports it. Derived through the SAME workerMatches rule - // the dispatchers grant by, so this row and the actual routing cannot - // disagree. - runsOn: Record<AutoQueueKind, LaneWorker[]>; -}; - -export async function buildAutoQueueLanes(): Promise<AutoQueueLanesPayload> { - const paths = getPaths(); - const settings = getSettings(); - const briefs = await getChannelBriefs(paths); - // allOperations, not one lane's members: the rail is the CATALOG view — every - // operation that is switched on, whatever queue it runs on. - const kinds = allOperations(settings); - - const bands = buildOperationBands({ - snapshots: briefs.map((c) => c.snapshot ?? null), - operationIds: kinds.map((k) => k.id), - }); - - const workerSummary = getWorkerPool().summary(); - const contendsForOf = new Map( - kinds.map((k) => [k.id, (k.laneFor?.(settings) ?? k.lane).contendsFor]), - ); - // An llm worker serves an operation its tags (or the llm default set) name; a - // remote worker takes units only when TAGGED (the unit-dispatch opt-in) and - // its tags intersect [operation, contended resource]. Local workers are not - // listed — the lane always runs here, and saying so would be noise. - // ONE PASS PER WORKER SLOT, not per (operation, slot) pair. A worker serving - // two of the lane's operations is one row with one slot count — the slots are - // the machine's, not the operation's, and counting them twice would advertise - // capacity that does not exist. - const runsOnFor = (ids: readonly string[]): LaneWorker[] => { - const byBase = new Map<string, LaneWorker>(); - for (const w of workerSummary) { - const takesAny = ids.some((id) => - w.kind === "llm" - ? workerMatches({ kind: "llm", tags: w.tags }, [id]) - : w.kind === "remote" && - (w.tags?.length ?? 0) > 0 && - workerMatches({ kind: "remote", tags: w.tags }, [ - id, - contendsForOf.get(id) ?? "cpu", - ]), - ); - if (!takesAny) continue; - const base = w.id.split("#")[0]; - const entry = byBase.get(base) ?? { - id: base, - name: w.name.replace(/ #\d+$/, ""), - kind: w.kind as "remote" | "llm", - slots: 0, - available: 0, - }; - entry.slots++; - if (w.state === "enabled" && !w.degraded) entry.available++; - byBase.set(base, entry); - } - return [...byBase.values()]; - }; - - const operations = {} as Record<AutoQueueKind, LaneOperationView[]>; - const runsOn = {} as Record<AutoQueueKind, LaneWorker[]>; - for (const lane of LANES) { - const ids = operationsForLane(lane, settings).map((op) => op.id); - operations[lane] = ids.map((id) => ({ id, label: operationLabel(id) })); - runsOn[lane] = runsOnFor(ids); - } - - return { bands, operations, runsOn }; -} +// What is left here is one type re-export, for the two client components that +// still name it through this path (`LaneHeader`, `RunnerOperationView`). Slice +// 1's sub-slice D repoints those and deletes this file. +export type { + LaneOperationView, + LaneWorker, + AutoQueueLanesPayload, +} from "yt-dlp-transcript-common/views/autoQueueLanes"; diff --git a/editor/app/operations/status.ts b/editor/app/operations/status.ts @@ -1,179 +1,72 @@ import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { getSettings } from "yt-dlp-transcript-common/lib/settings"; import { - type AutoRunnerStatus, - type NextUpView, - type RecencyKeyView, computeLeafPending, getAutoRunnerStatus, - laneDispatchRoot, + type AutoRunnerStatus, + type LeafPending, } from "yt-dlp-transcript-common/controller/autoRunner"; import { - type AutoQueueKind, - type AutoQueuePick, readAutoQueueState, + type AutoQueueKind, } from "yt-dlp-transcript-common/jobs/autoQueueState"; import { LANES } from "yt-dlp-transcript-common/lib/autoQueueTypes"; -import type { AutoQueuePolicy } from "yt-dlp-transcript-common/jobs/autoQueuePolicy"; import { getWorkerPool } from "yt-dlp-transcript-common/jobs/workerPool"; -import { isGateHeld } from "yt-dlp-transcript-common/lib/pauseGates"; -import type { FocusSummary } from "yt-dlp-transcript-common/lib/channelPriority"; -import { buildAutoQueueLanes, type AutoQueueLanesPayload } from "./lanes"; +import { buildAutoQueueLanes } from "yt-dlp-transcript-common/views/autoQueueLanes"; import { - type PriorityView, - laneFocusSummary, - readPriorityView, -} from "./channelPriorityView"; - -// Read-only payload for the Auto-Queue panel: per-kind runner status (running?, -// what's in flight, in-flight counts per tree node), the effective policy, the -// recent pick log, and snapshot-derived pending counts per leaf — so the UI can -// show "cornbreadman: 12 pending" and which leaf is currently being serviced. -// Backs a passive UI poll, like /api/scheduler/status. - -// A platform currently in a rate-limit (429) / network cooldown. `untilMs` is -// the epoch ms the cooldown lapses; the client renders remaining time itself so -// it counts down live. Only the download kind ever has these. -export type PlatformCooldownView = { - platform: string; - untilMs: number; - fails: number; -}; + buildAutoQueueStatusPayload as build, + type AutoQueueKindStatus, + type AutoQueueStatusPayload, + type PlatformCooldownView, +} from "yt-dlp-transcript-common/views/autoQueueStatus"; +import { getChannelBriefs } from "../lib/requestCache"; +import { readPriorityView } from "./channelPriorityView"; -export type AutoQueueKindStatus = { - kind: AutoQueueKind; - // THE POLICY AS THE LANE DISPATCHES IT, not as it is stored. Every field is - // the stored one except `root`, which is the COMPILED tree whenever - // `settings.channelPriority` says anything (see channelPriorityView.ts). The - // claim ladder is drawn from this, so a rung is a rule the runner actually - // has — and with a priority model set the stored tree is not one. - policy: AutoQueuePolicy; - runner: AutoRunnerStatus; - pendingByLeaf: Record<string, number>; - // The first few pending ids of each rule, in the order the runner would hand - // them out — the drill-down behind each rung's count, and the only way to - // check an ordering change by eye. - pendingHeadByLeaf: Record<string, string[]>; - // videoId -> owning channel slug, for the drill-down's links. Covers exactly - // the ids in pendingHeadByLeaf (plus nextUp's). - ownerByVideo: Record<string, string>; - // videoId -> recency sort key, for the same ids, when the policy orders by - // recency. Absent entries just render without a date. - recencyByVideo: Record<string, RecencyKeyView>; - // What the policy would dispatch next, and which rules it passed over. - nextUp: NextUpView | null; - picks: AutoQueuePick[]; - // Platforms paused by a rate-limit/network backoff (download kind only). - cooldowns: PlatformCooldownView[]; - // This runner's lane gate — what the page's pause button draws, and what lets - // the operations rail show a runner HOLDING for the first time. - // - // THE TWO KINDS ANSWER IT FROM DIFFERENT PLACES, on purpose. Transcription's - // hold is LIVE, on the worker pool; settings.transcriptionsPaused is only what - // the boot hook re-applies after a restart, and the e2e harness rewrites - // settings wholesale between tests while the pool keeps its pausedSnapshot. - // Download has no live counterpart: its flag IS the gate, read at dispatch. - held: boolean; - // THE TREE ABOVE IS GENERATED, so the editor for it is read-only and the - // channel priorities on /channels are where it is edited. False for a corpus - // with no priorities set, which is every corpus until one is. - policyCompiled: boolean; - // THE FOCUS BANNER'S NUMBERS, FOR THIS LANE. Always present and inert when - // `active` is false, so the banner is one component with one early return - // rather than a conditional on the payload. The counts are this lane's own — - // `focusPending`/`otherPending` differ per lane by construction. - focus: FocusSummary; - // The focus's display name, resolved on the server (the model stores a - // siteId, and the banner does no I/O). Null when nothing is focused. - focusName: string | null; +export type { + AutoQueueKindStatus, + AutoQueueStatusPayload, + PlatformCooldownView, }; -export type AutoQueueStatusPayload = Record< - AutoQueueKind, - AutoQueueKindStatus -> & { - // The sweep half of the digest and backfill pipelines, plus the comparison - // rail's bands. On the SAME payload as the runners rather than a second - // endpoint, because the rail's whole purpose is that four lanes are read - // together — two polls would let the rail and the focused lane disagree about - // the same moment. - // - // Every lane now has a `status[lane]` entry too. The two are not duplicates: - // this is what the SWEEPS are doing (armed, scoped, planned), and the entry - // is what the lane's TREE says. Slice 1.3 deletes this half. - lanes: AutoQueueLanesPayload; -}; - -async function buildKind( - kind: AutoQueueKind, - priority: PriorityView, -): Promise<AutoQueueKindStatus> { - const paths = getPaths(); - const stored = getSettings().autoQueue[kind]; - // The stored policy, with the DISPATCHED root in place of the stored one. - // Spread rather than rebuilt, so `held`, `snoozeUntil`, `enabled`, `order` - // and `maxWorkers` come through untouched — the same rule every writer of a - // policy in this repo follows. - const policy: AutoQueuePolicy = { - ...stored, - root: laneDispatchRoot(kind, stored, priority, priority.slugs), - }; - const runner = getAutoRunnerStatus(kind); - const state = await readAutoQueueState(paths); - const pending = await computeLeafPending(kind, paths); - const now = Date.now(); - // Surface platforms still inside their cooldown window so the operator can see - // why an otherwise-pending platform isn't being serviced (and a manual sync on - // it would be refused). A manual sync/download 429 writes into this same map. - const cooldowns: PlatformCooldownView[] = Object.entries( - state[kind].platformBackoff, - ) - .filter(([, e]) => e.until > now) - .map(([platform, e]) => ({ platform, untilMs: e.until, fails: e.fails })) - .sort((a, b) => b.untilMs - a.untilMs); - return { - kind, - policy, - runner, - pendingByLeaf: pending.counts, - pendingHeadByLeaf: pending.head, - ownerByVideo: pending.owner, - recencyByVideo: pending.recency, - nextUp: pending.nextUp, - picks: state[kind].picks, - cooldowns, - // THE ASYMMETRY IS TRANSCRIPTION'S ALONE, and it is not a special case for - // "the first lane": its hold is LIVE on the worker pool, while every other - // lane's flag IS its gate. See lib/pauseGates.ts. - held: - kind === "transcription" - ? getWorkerPool().isPaused() - : isGateHeld(getSettings(), kind), - policyCompiled: priority.compiled, - // Keyed by COMPILED leaf id (`prio-focus-<slug>`), which is why this is - // computed here and not on the client: it is only meaningful against the - // counts of the tree the lane dispatches from. - focus: laneFocusSummary(priority, pending.counts), - focusName: priority.name, - }; -} - -// ONE ENTRY PER LANE, off LANES rather than off two hand-written names — so a -// lane added to the model appears on this payload with no edit here, which is -// the whole point of the widened type. +// THE SHELL. The payload is `common/views/autoQueueStatus.ts` (and the rail's +// bands are `views/autoQueueLanes.ts`); everything below is the reading. +// +// ONE RESOLUTION FOR ALL FOUR LANES, in three places rather than one. The focus +// set costs a channel listing and, for a site focus, a sites read; the +// auto-queue state document is one JSON parse — it was read once PER LANE, four +// times per poll, because `buildKind` did its own reading; and the channel +// briefs are shared with the lanes builder through the per-request cache. export async function buildAutoQueueStatusPayload(): Promise<AutoQueueStatusPayload> { - // ONE RESOLUTION FOR ALL FOUR LANES. The focus set costs a channel listing - // and, for a site focus, a sites read; the four lanes compile from the same - // one, so resolving per lane would pay for it four times on a 3 s poll. - const priority = await readPriorityView(); - const [kinds, lanes] = await Promise.all([ - Promise.all(LANES.map((lane) => buildKind(lane, priority))), - buildAutoQueueLanes(), + const paths = getPaths(); + const settings = getSettings(); + const pool = getWorkerPool(); + const [priority, state, briefs] = await Promise.all([ + readPriorityView(), + readAutoQueueState(paths), + getChannelBriefs(paths), ]); - return { - ...(Object.fromEntries( - LANES.map((lane, i) => [lane, kinds[i]]), - ) as Record<AutoQueueKind, AutoQueueKindStatus>), - lanes, - }; + const pendingByKind = await Promise.all( + LANES.map((lane) => computeLeafPending(lane, paths)), + ); + const byLane = <T>(values: readonly T[]): Record<AutoQueueKind, T> => + Object.fromEntries(LANES.map((lane, i) => [lane, values[i]])) as Record< + AutoQueueKind, + T + >; + return build({ + settings, + pool, + now: Date.now, + priority, + // Constructs the runner if it is not already there, which is exactly why + // the view cannot ask for it itself. + runner: byLane<AutoRunnerStatus>(LANES.map((lane) => getAutoRunnerStatus(lane))), + state, + pending: byLane<LeafPending>(pendingByKind), + lanes: buildAutoQueueLanes({ + settings, + briefs, + workerSummary: pool.summary(), + }), + }); }