Archilyzer · Source

archilyzer

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

commit 79031ffc55b487302964943a8f3d20fa3bca3ddb
parent 6bcce6ed122ec024bab6ba93b1e857a68fa6f2c0
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Mon,  7 Sep 2026 22:21:58 -0400

common: the sweeps and the arbiter retire, and ten settings fields go with them

Slice 1.3, second commit — the backend half. `digestSweep.ts` (368),
`backfillSweep.ts` (519), `sweepPreview.ts` (167), `sweepRecency.ts` (141),
`lib/sweepPlan.ts` (163), `arbiter.ts` (418) and `arbiterWork.ts` (67) are
deleted, with `sweepPreview.test.ts` and `arbiter.test.ts`; so are the
`backfill-sweep` and `operations-arbiter` job kinds and both boot-resume hooks.
Nothing replaces them: the lane runner shipped in 1.2 already does this work,
and gate A watched it do a production pass.

The arbiter never ran — zero `operations-arbiter` records in 1,305 production
jobs, no state, no boot resume, and it could not dispatch the two lanes carrying
100% of live dispatch. The sweeps ran for months and were the right answer
before there was a second dispatcher; what they never had was a tree, a claim
ladder, a next-up or a pick log.

THE SCOPE IS NOT LOST, it is migrated. `getSettings` runs
`migrateSweepsToLanes` on read, so `digest.sweepEnabled/sweepChannels` and
`backfill.sweepEnabled/sweepKinds/sweepChannels/order` become
`autoQueue.digest` / `.backfill` — enabled, and a strict root of channel leaves
— when those blocks are absent from the file. It never enables a lane the flag
did not, which on this corpus is the difference between a disabled lane and
~66,540 audio re-downloads.

`digest.recencyOrder` is NOT migrated onto the lane order, and that is the one
decision in here worth arguing with. The digest order is a COMPOSITION —
shortest inside a day, newest day first — which the lane already spells
`cheapest`; `recencyOrder` was only its date half. Migrating it would have
replaced the composition with one of its two terms. So the date half is fixed at
newest-first, the live value, and a digest lane wanting pure recency sets
`order: "newest"` and gets no duration term at all. `digest.recencyReach` and
`backfill.reach` are dropped rather than migrated: a runner lane has no reach
axis (OrderReach.tsx renders `reach={null}` for exactly this reason), because a
rule already orders every video it claims across every channel, and which rule
goes first is the tree.

`backfill.weight` retires too, and it is a real behaviour change rather than a
rename. One 0..1 scalar answered two different questions — "does this run
contend with transcription" and "how much of the machine may it have" — and its
default of 0 parked a network-bound attribution run behind a GPU transcription
it competes with for nothing. The yield is the operation's declared
`contendsFor` (already the rule since 1.2); the share is `concurrency` and the
lane policy's `maxWorkers`. `backfillLimit` takes `idleOnly` now.

`laneBlockedReason` was the last reader of a sweep flag and goes with them,
along with the `lane-blocked` idle reason. `startAutoRunnerBlockedReason` still
answers — there is one reason left, a switched-off policy, and an operator who
clicks Start and sees nothing still deserves to be told which switch.

The architecture allow-list SHRANK by one: `lib/sweepPlan.ts ->
controller/planOrder` died with the file. The two `jobs/* -> controller` entries
are unrelated to the sweeps (a snapshot builder and a capacity probe) and are
still real back-edges, so they stay; the guard fails on a stale entry, which is
how that was checked rather than assumed.

The editor half lands next; this commit leaves it referencing deleted modules.
common 888 tests (897 before, +8 laneMigration, -17: seven sweepPreview, seven
arbiter, three weight), tsc clean in common, mcp, export, homepage and umtool.

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

Diffstat:
Mcommon/architecture.test.ts | 6------
Dcommon/controller/arbiter.test.ts | 177-------------------------------------------------------------------------------
Dcommon/controller/arbiter.ts | 418-------------------------------------------------------------------------------
Dcommon/controller/arbiterWork.ts | 67-------------------------------------------------------------------
Mcommon/controller/autoRunner.ts | 46+++++++++++++++++-----------------------------
Dcommon/controller/backfillSweep.ts | 519-------------------------------------------------------------------------------
Dcommon/controller/digestSweep.ts | 368-------------------------------------------------------------------------------
Mcommon/controller/laneForOperation.test.ts | 24++++++++++++------------
Mcommon/controller/laneGuards.ts | 4++--
Mcommon/controller/noCorpusWalkInRenderPaths.test.ts | 23++++++++++-------------
Mcommon/controller/operationBatch.test.ts | 94++++++++++++++++++++++++++-----------------------------------------------------
Mcommon/controller/operationBatch.ts | 70++++++++++++++++++++++++++++++++++++++++------------------------------
Mcommon/controller/operationJobs.ts | 45++++++++++++++++++++++-----------------------
Mcommon/controller/operationLane.ts | 21+++++++++++----------
Mcommon/controller/remoteCapacity.ts | 2+-
Dcommon/controller/sweepPreview.test.ts | 321-------------------------------------------------------------------------------
Dcommon/controller/sweepPreview.ts | 167-------------------------------------------------------------------------------
Dcommon/controller/sweepRecency.ts | 0
Mcommon/jobs/autoQueuePolicy.ts | 15++++++++-------
Mcommon/jobs/jobKinds.ts | 23-----------------------
Mcommon/jobs/snapshotScheduler.ts | 5++---
Mcommon/lib/autoQueueTypes.ts | 4++--
Mcommon/lib/operations.ts | 28++++++++++++++--------------
Mcommon/lib/pauseGates.test.ts | 17++++++++---------
Mcommon/lib/pauseGates.ts | 32++------------------------------
Mcommon/lib/settings.ts | 126++++++++++++-------------------------------------------------------------------
Dcommon/lib/sweepPlan.ts | 163-------------------------------------------------------------------------------
Meditor/instrumentation.ts | 42+++++++-----------------------------------
28 files changed, 208 insertions(+), 2619 deletions(-)

diff --git a/common/architecture.test.ts b/common/architecture.test.ts @@ -71,12 +71,6 @@ const ALLOWED: Record<string, string> = { "lib/operations.test.ts -> controller/normalizeTranscript": "test of the above; moves with it", - // One planner calling the recency ordering, which is cached against LMDB. - // Phase 1 step 3 folds sweepPreview.ts and sweepRecency.ts into sweepPlan.ts - // and decides which side of the line the result lands on. - "lib/sweepPlan.ts -> controller/planOrder": - "orderPlanByRecency; resolved when the four sweep planners become one (phase 1 step 3)", - // Two schedulers reaching into the controller for the thing they schedule. // Real work, not a type: the snapshot builder and the remote-capacity probe. "jobs/snapshotScheduler.ts -> controller/channelSnapshot": diff --git a/common/controller/arbiter.test.ts b/common/controller/arbiter.test.ts @@ -1,177 +0,0 @@ -import { test } from "node:test"; -import assert from "node:assert/strict"; -import type { - AutoQueueGroup, - ChannelWork, -} from "../jobs/autoQueuePolicy"; -import type { Lane } from "../lib/operations"; -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: Lane = { queueKey: "digest:local", contendsFor: "gpu" }; -const CPU: Lane = { queueKey: "backfill", contendsFor: "cpu" }; - -const LANES: Record<string, Lane> = { - digest: GPU, - diarization: CPU, - "attribution-text": CPU, -}; -const laneOf = (op: string): Lane | 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 @@ -1,418 +0,0 @@ -// 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 { DIGEST_REMOTE_QUEUE } from "../lib/queueKeys"; -import { laneBlockedReason } from "../lib/pauseGates"; -import { getRegistry } from "../jobs/registry"; -import { runManagedFunction } from "../jobs/streamCommand"; -import { drainStream } from "../jobs/drainStream"; -import { - buildPendingByLeaf, - flattenLeaves, - type AutoQueueGroup, - type ChannelWork, -} from "../jobs/autoQueuePolicy"; -import { - DIGEST_OPERATION_ID, - allOperations, - operationLabel, - type Lane, -} from "../lib/operations"; -import { laneForOperation } from "./operationLane"; -import { runOperationChannelJob } from "./operationJobs"; -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: Lane; - // The leaf that claimed it, for the log. An operator asking "why that - // channel" gets the rule back rather than a shrug. - leafId: string; -}; - -// 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) => Lane | 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. - // - // No second filtering pass: buildArbiterChannelWork projects `buckets: {}`, - // so a bucket leaf finds no list here whatever it names, and the loop below - // skips every leaf with no `match.operation` anyway. retainLeaves was doing - // both of those a third time. - const pending = buildPendingByLeaf(root, channels, []); - 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 REGISTRY does not know has no lane, so it is not - // dispatched. Silently skipping is right: a snapshot may name a kind that - // has since been renamed, and download/transcription are catalogued - // (digest.dependsOn needs them to be) but dispatched by their own runners. - // - // A switched-off kind does not reach here at all, and not through this - // branch: `enabled` below projects ids only for kinds allOperations - // returns, so a disabled feature arrives with an empty id list and is - // skipped one line up. laneForOperation itself does not consult - // `enabled` — see controller/laneForOperation.test.ts. - 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. -// THE RULE MOVED to lib/pauseGates.laneBlockedReason, per LANE — because a -// runner is per lane where the arbiter was one dispatcher for both. This asks -// it about both, which is what the arbiter has always meant, and keeps the -// arbiter working until 1.3 retires it. -export function arbiterBlockedReason(): string | null { - const settings = getSettings(); - return ( - laneBlockedReason(settings, "digest") ?? - laneBlockedReason(settings, "backfill") - ); -} - -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. - // - // ONE CALL, whichever operation it is. This used to be a ternary on - // DIGEST_OPERATION_ID with two hand-written option sets, which meant the arbiter - // had to be edited in step with the registry and re-derived the digest lane - // from settings a second time — while `unit.lane` beside it was already the - // answer. It is now: laneForOperation asks the kind's own laneFor, so - // `unit.lane.queueKey` IS the live lane, and the choice below reads it rather - // than re-asking getSettings. - const outcome = await runOperationChannelJob({ - paths, - channelSlug: unit.channelSlug, - operation: 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 - // pins. - queueKey: unit.lane.queueKey, - digest: { - lane: unit.lane.queueKey === DIGEST_REMOTE_QUEUE ? "remote" : "local", - }, - // Behind anything an operator clicks by hand. - 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; - } - // WAIT FOR IT. The arbiter dispatches ONE unit per lane per pass and re-plans - // every pass, so a pass that returned before its job finished would re-plan - // against work still in flight and hand the next pass the same channel. The - // runners used to drain internally; they no longer do — see - // controller/operationJobs.ts for why that had to move to the callers. - await drainStream(outcome.stream); - 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([ - ...allOperations(settings).map((k) => k.id), - DIGEST_OPERATION_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 @@ -1,67 +0,0 @@ -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 OperationSnapshotEntry. 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 @@ -29,7 +29,6 @@ import { selectableBucketsForKind, } from "../jobs/autoQueuePolicy"; import { operationsForLane, type Operation } from "../lib/operations"; -import { laneBlockedReason } from "../lib/pauseGates"; import { cheapestComparator, classifyOperationUnit, @@ -172,10 +171,6 @@ export type AutoRunnerIdleReason = // no-workers: pauseAll disables every worker, so by slot count the two look // identical, and only one of them is fixed by enabling a worker. | "workers-paused" - // Sweep X is armed for this lane, so its runner stands down rather than - // starting the same channel twice. Digest and backfill only, and only until - // the sweeps retire in 1.3. - | "lane-blocked" // The lane's own gate is shut (digestsPaused, or backfill.enabled false). // DISTINCT from downloads-paused, which names one specific flag, and from // capped, which is contention rather than intent — the fix for this one is a @@ -474,16 +469,17 @@ async function recencyOrdering( // two, and because a composition that only exists as a pair of adjacent // sort() calls is exactly the thing a migration drops half of. // - // The date half's DIRECTION is `settings.digest.recencyOrder` — "newest" on - // the live corpus. It is a parameter rather than a constant because slice - // 1.3 retires that field into the lane's tree and fixes this at newest-first - // (see the as-shipped note); a digest lane wanting pure recency then sets - // `order: newest` and gets no duration term at all. + // The date half's DIRECTION IS FIXED AT NEWEST-FIRST, which is the value + // `settings.digest.recencyOrder` carried on the live corpus before slice 1.3 + // retired it. It is a constant rather than a policy field because `order` + // already holds one enum: a digest lane that wants pure recency sets + // `order: "newest"` and gets no duration term at all, and one that wants + // `cheapest` is asking for the composition this comparator IS. const durations = await loadDurations(paths, owner, candidateIds); return { compare: cheapestComparator({ duration: (id) => durations.get(id), - recency: makeRecencyComparator(keys, getSettings().digest.recencyOrder), + recency: makeRecencyComparator(keys, "newest"), }), keys, }; @@ -1016,17 +1012,6 @@ async function runLoop( live.idleReason = "downloads-paused"; return null; } - // TWO DISPATCHERS ON ONE LANE WOULD START THE SAME CHANNEL TWICE. A HOLD, - // not a stop — disarm the sweep and the runner resumes within one poll, - // where a stop would need the operator to notice and press Start. Says - // WHICH sweep, because "blocked" alone sends nobody anywhere. Retires with - // the sweeps in 1.3. - const blocked = laneBlockedReason(settings, kind); - if (blocked) { - if (live.idleReason !== "lane-blocked") onLog(`Auto-${kind} idle: ${blocked}`); - live.idleReason = "lane-blocked"; - return null; - } // The operation lanes need their run context before they can classify // anything. limit() keeps it warm; this is the path that WAITS for it, so a // first tick does not have to come back empty. @@ -1592,11 +1577,6 @@ export async function startAutoRunner( const settings = getSettings(); const policy = settings.autoQueue[kind]; if (!policy.enabled) return null; - // REFUSED, not started and then immediately idle: an operator clicking Start - // beside an armed sweep should be told which sweep, and startAutoRunnerReason - // is what the API hands back. The loop holds on the same rule if the sweep is - // armed while it runs. - if (laneBlockedReason(settings, kind)) return null; const singleton = getSingleton(); const existing = singleton.runners.get(kind); if (existing && getRegistry().get(existing.jobId)?.status === "running") { @@ -1633,9 +1613,17 @@ export async function startAutoRunner( } // Why this lane's runner will not start, or null when it will. The console asks -// the same question the starter asks and gets the same sentence back. +// the same question the starter asks and gets the same sentence back — an +// operator who clicks Start and sees nothing happen deserves the reason. +// +// ONE REASON LEFT. Until slice 1.3 there were two: a disabled policy, and an +// armed sweep on the same lane (two dispatchers would have started the same +// channel twice). The sweeps are gone, so the tree's own switch is the whole +// answer. export function startAutoRunnerBlockedReason(kind: AutoQueueKind): string | null { - return laneBlockedReason(getSettings(), kind); + return getSettings().autoQueue[kind].enabled + ? null + : "This lane's policy is switched off. Enable it in the rules below, then start the runner."; } // Start every enabled runner. Called from the editor instrumentation hook at diff --git a/common/controller/backfillSweep.ts b/common/controller/backfillSweep.ts @@ -1,519 +0,0 @@ -// The corpus-wide backfill sweep: one launcher for the catch-up that runs for -// days, and survives the restarts it will certainly outlive. -// -// Cloned from digestSweep.ts, which is the complete worked example of -// armed-and-resumable in this repo. Four of its shapes are load-bearing and are -// kept verbatim, each because getting it wrong has already cost something here: -// -// 1. IT LAUNCHES THE EXISTING PER-CHANNEL JOB, in sequence. Not one job per -// video: the registry keeps 100 records and the log 500, so 77,000 jobs -// would evict the history of the very run recording them. -// 2. IT STORES NO CURSOR. The work-list is recomputed each pass and -// eligibility is re-derived from disk inside every batch, so a restart, a -// newly-finished transcription and a settings change that invalidates the -// corpus are all just visible on the next pass. -// 3. IT PERSISTS INTENT AND SCOPE TOGETHER. The boot hook re-launches from -// settings alone, so arming without recording the scope resurrects a -// deliberately-bounded run as a corpus-wide one — and a stop that does not -// CLEAR the scope leaves a stale one to silently narrow the next sweep. -// Both were real bugs in the digest sweep. -// 4. THE ORCHESTRATOR RUNS ON queueKey "". It waits on the per-channel jobs -// that serialize on BACKFILL_QUEUE; putting it on that queue would deadlock -// — it would hold the only slot while waiting for a job that needs it. -// -// The termination guard is measured against REMAINING WORK, not against jobs -// starting. An engine that is broken makes every channel job start normally and -// then do nothing, which from out here looks exactly like progress. - -import type { Paths } from "../lib/paths"; -import { getPaths } from "../lib/paths"; -import { getSettings, writeSettings } from "../lib/settings"; -import { getRegistry } from "../jobs/registry"; -import { runManagedFunction } from "../jobs/streamCommand"; -import { drainStream } from "../jobs/drainStream"; -import { BACKFILL_QUEUE } from "../lib/queueKeys"; -import { - reachableOperationWork, - resolveBackfillLaneOperations, - type OperationSnapshotEntry, -} from "../lib/operations"; -import { listChannelStatsFromDisk, readChannelSnapshot } from "./channels"; -import { countOperationWork } from "./operationBatch"; -import { runBackfillChannelJob } from "./operationJobs"; -import type { AutoQueueOrder } from "../jobs/autoQueuePolicy"; -import { buildRecencyKeys } from "./recencyIndex"; -import { orderPlanByRecency } from "./planOrder"; -import type { SweepPlanEntry } from "../lib/sweepPlan"; - -export const BACKFILL_SWEEP_KIND = "backfill-sweep"; -// The per-channel runner and its job kind live in controller/operationJobs.ts -// now — one runner for a hand-clicked run, a swept one and an arbiter-dispatched -// one. This module is the SWEEP: the plan, the loop, and the drain. - -// How long the loop waits before re-planning when a pass did no work. Generous: -// nothing here is latency-sensitive and planning reads every channel's videos. -const IDLE_POLL_MS = 60_000; - -// Give up after this many consecutive passes that move no work. 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 -// the sweep cannot make progress on, and both want a human. -const MAX_BARREN_PASSES = 3; - -type SweepLive = { - jobId: string; - startedAt: number; - // The per-channel job the sweep is currently waiting on. Tracked so a stop can - // drain THAT too — without it, "stop" means "after the current channel", which - // on a large channel is hours. - channelJobId: string | null; -}; -type SweepSingleton = { live: SweepLive | null }; - -declare global { - // eslint-disable-next-line no-var - var __yttBackfillSweep__: SweepSingleton | undefined; -} - -function getSingleton(): SweepSingleton { - if (!globalThis.__yttBackfillSweep__) { - globalThis.__yttBackfillSweep__ = { live: null }; - } - return globalThis.__yttBackfillSweep__; -} - -export function getBackfillSweepJobId(): string | null { - const live = getSingleton().live; - if (!live) return null; - return getRegistry().get(live.jobId)?.status === "running" ? live.jobId : null; -} - -export type BackfillSweepOptions = { - paths?: Paths; - // Restrict to these kinds / channels. Absent = every enabled lane kind, whole - // corpus. - kindIds?: string[]; - channelSlugs?: string[]; -}; - -// One channel's row in the plan: how much is left, and the newest / oldest -// upload date (YYYYMMDD) among its REACHABLE ids — the dates present only when -// the plan was asked to order by recency. "" is "no dated reachable video": -// last under "newest", first under "oldest", the same rule makeRecencyComparator -// applies to a single video. -// -// ONE TYPE, TWO NAMES. The definition lives in lib/sweepPlan.ts because the -// console's preview folds the same shape in the browser, and a second -// declaration here is how the two would drift into disagreeing about a channel. -export type BackfillPlanEntry = SweepPlanEntry; - -// What is left, per channel, ordered heaviest-first by REACHABLE work. -// -// Reachable only, deliberately: ordering by total remaining would sort the -// corpus by how much of it is unreachable, which on the measured numbers (835 -// reachable vs ~76,270 needing a re-download) means sorting by noise. -// -// PLANNED OFF SNAPSHOTS, not off a corpus walk. countOperationWork() is a -// readdir + readVideoFiles + sidecar read for EVERY video — ~77,000 directory -// reads — and this function runs once per pass of a sweep that runs for days. -// The channel snapshots already hold these exact numbers in -// snapshot.backfill[kind], computed by a report the system regenerates anyway -// (runBackfillChannelJob calls requestChannelSnapshot when it finishes), so the -// walk was re-deriving, every hour, a number sitting in a file. -// -// The fallback is not a nicety: a channel that has never been snapshotted has no -// entry, and planning it as zero would silently exclude it from the sweep -// forever. So an ABSENT snapshot walks; a present one is trusted. -// -// Trusting it costs a bounded staleness — a snapshot regenerates on a ~1s -// debounce after mutations, and worst case a channel is planned against -// last-pass numbers. That is harmless here for the same reason the sweep stores -// no cursor: the batch re-derives eligibility from disk on every pull, so a -// wrong plan costs at most one channel visited in the wrong order, never work -// done twice or work missed. -export async function buildBackfillSweepPlan(opts: { - paths: Paths; - kindIds?: string[]; - channelSlugs?: string[]; - // Order the CHANNEL list by each channel's freshest (or oldest) reachable - // video instead of heaviest-first. Absent/"listed" leaves today's order - // untouched and costs nothing — no dates are read at all. - order?: AutoQueueOrder; - onLog?: (msg: string) => void; -}): Promise<BackfillPlanEntry[]> { - const wanted = new Set(opts.channelSlugs ?? []); - const channels = (await listChannelStatsFromDisk(opts.paths)).filter( - (ch) => wanted.size === 0 || wanted.has(ch.slug), - ); - // Resolved once: the same filtering resolveBackfillLaneOperations does for the batch, - // so the plan counts exactly the kinds the run will act on. A snapshot holds - // an entry per kind that was enabled when it was WRITTEN, which may be more - // than are in scope now. - const kinds = resolveBackfillLaneOperations(getSettings(), opts.kindIds).map( - (k) => k.id, - ); - const order = opts.order ?? "listed"; - const plan: BackfillPlanEntry[] = []; - // Reachable ids per channel, collected only when they will be used. The - // snapshot's `ids` IS the reachable set (see OperationSnapshotEntry), so this - // is a read of something already in hand — but on this corpus it is ~78,000 - // strings, so it is not built when nothing will sort by it. - const idsBySlug = order === "listed" ? null : new Map<string, string[]>(); - let walked = 0; - for (const ch of channels) { - const snapshot = await readChannelSnapshot(opts.paths, ch.slug); - const counts = snapshot?.backfill - ? countFromSnapshot(snapshot.backfill, kinds) - : ((walked++, - await countOperationWork("backfill", opts.paths, ch.slug, { - operationIds: opts.kindIds, - }))); - if (idsBySlug && snapshot?.backfill) { - const ids = new Set<string>(); - for (const id of kinds) { - for (const videoId of snapshot.backfill[id]?.ids ?? []) ids.add(videoId); - } - idsBySlug.set(ch.slug, [...ids]); - } - plan.push({ - channelSlug: ch.slug, - reachable: counts.reachable, - missingInput: counts.missingInput, - newestPending: "", - oldestPending: "", - }); - } - if (walked > 0) { - opts.onLog?.( - `${walked} channel(s) had no snapshot yet and were counted by walking their videos.`, - ); - } - - // ONE recency build for the whole plan. The digest plan gets its dates free - // (build:stats already put uploadDate on the row it was reading anyway); a - // backfill snapshot carries ids and no dates, so this is the one place that - // pays. It is affordable because a sweep pass is HOURS apart, and because - // recencyIndex memoizes a video's date process-wide once read. - if (idsBySlug) { - const candidateIds = new Set<string>(); - const owner = new Map<string, string>(); - for (const [slug, ids] of idsBySlug) { - for (const id of ids) { - candidateIds.add(id); - if (!owner.has(id)) owner.set(id, slug); - } - } - try { - const keys = await buildRecencyKeys({ - paths: opts.paths, - meta: channels.map((c) => ({ slug: c.slug })), - candidateIds, - owner, - // Every reachable id is a directory on disk, so there is nothing for - // playlist interpolation to estimate and the playlist reads would be - // pure cost. - interpolate: false, - }); - for (const entry of plan) { - for (const id of idsBySlug.get(entry.channelSlug) ?? []) { - const key = keys.get(id)?.key ?? ""; - if (!key) continue; - if (key > entry.newestPending) entry.newestPending = key; - if (!entry.oldestPending || key < entry.oldestPending) { - entry.oldestPending = key; - } - } - } - } catch { - // Dates are an ordering nicety; failing to read them must not stop a - // sweep. Every key stays "", which sorts the plan exactly as it does - // today via the weight tiebreak below. - } - } - - return orderPlanByRecency( - plan, - order, - (a, b) => - b.reachable - a.reachable || a.channelSlug.localeCompare(b.channelSlug), - ); -} - -// The snapshot's per-kind counts, summed over the kinds in scope. `missing` and -// `stale` are the reachable half; `missingInput` stays its own number and is -// never added to them — see lib/operations.ts's header. -// -// EXPORTED so the console's preview counts with the RUN'S OWN ARITHMETIC rather -// than a second copy of it. A preview that disagrees with the run about how much -// work a channel holds is worse than no preview: it is the number an operator -// arms a multi-day commitment against. See ./sweepPreview.ts. -export function countFromSnapshot( - backfill: Record<string, OperationSnapshotEntry>, - kinds: string[], -): { reachable: number; missingInput: number } { - let reachable = 0; - let missingInput = 0; - for (const id of kinds) { - const entry = backfill[id]; - if (!entry) continue; - reachable += reachableOperationWork(entry); - missingInput += entry.missingInput; - } - return { reachable, missingInput }; -} - -async function runSweepLoop( - paths: Paths, - kindIds: string[] | undefined, - channelSlugs: string[] | undefined, - live: SweepLive, - onLog: (msg: string) => void, - signal: AbortSignal, - drainSignal: AbortSignal, -): Promise<void> { - let pass = 0; - // Consecutive passes that left the remaining work UNCHANGED — see the header - // for why this measures the plan and not whether jobs started. - let barrenPasses = 0; - let lastReachable: number | null = null; - for (;;) { - if (signal.aborted || drainSignal.aborted) return; - // The operator turned it off: stop CLEANLY rather than being cancelled, so - // the job ends "done" and the queue is released. - if (!getSettings().backfill.sweepEnabled) { - onLog("Backfill sweep disarmed in settings — stopping."); - return; - } - // The feature itself may have been switched off underneath us. Same clean - // stop: a sweep for a disabled backfill has nothing to do and should not - // sit there looking wedged. - if (resolveBackfillLaneOperations(getSettings(), kindIds).length === 0) { - onLog( - "No enabled backfill kind is in scope — stopping. (Enable the feature in Settings, then re-arm.)", - ); - return; - } - - pass++; - // Resolved per pass, not at launch: switching to newest-first mid-sweep - // takes effect on the next plan rather than only after stopping and - // re-arming. Only "corpus" reaches the CHANNEL order — under "channel" the - // plan stays heaviest-first and each per-channel batch applies the order - // itself from settings.backfill.order. - const backfillNow = getSettings().backfill; - const planOrder = - backfillNow.reach === "corpus" ? backfillNow.order : "listed"; - const plan = await buildBackfillSweepPlan({ - paths, - kindIds, - channelSlugs, - order: planOrder, - onLog, - }); - const work = plan.filter((c) => c.reachable > 0); - const reachable = plan.reduce((n, c) => n + c.reachable, 0); - const missingInput = plan.reduce((n, c) => n + c.missingInput, 0); - if (work.length === 0) { - onLog( - `Pass ${pass}: nothing reachable left to backfill. ` + - `${missingInput.toLocaleString()} video(s) would need their media re-acquired first.`, - ); - return; - } - // BOTH NUMBERS, never their sum. See lib/operations.ts's header. - onLog( - `Pass ${pass}: ${work.length} channel(s), ${reachable.toLocaleString()} video(s) ` + - `reachable now, ${missingInput.toLocaleString()} needing their media re-acquired` + - (planOrder === "listed" - ? ", heaviest channel first." - : `, ${planOrder} channel first (${work[0].channelSlug} @ ${ - planOrder === "newest" - ? work[0].newestPending || "undated" - : work[0].oldestPending || "undated" - }).`), - ); - - let didWork = false; - for (const channel of work) { - if (signal.aborted || drainSignal.aborted) return; - if (!getSettings().backfill.sweepEnabled) { - onLog("Backfill sweep disarmed in settings — stopping."); - return; - } - onLog( - `→ ${channel.channelSlug}: ${channel.reachable.toLocaleString()} reachable.`, - ); - const outcome = await runBackfillChannelJob({ - paths, - channelSlug: channel.channelSlug, - kindIds, - // Behind anything an operator clicks by hand: a sweep is days long and - // must never make a deliberate single-channel run wait for it. - background: true, - onStarted: (jobId) => { - live.channelJobId = jobId; - }, - }); - if (!outcome.ok) { - live.channelJobId = null; - // One channel failing to START is not the sweep failing. - onLog(`!! ${channel.channelSlug}: ${outcome.error ?? "failed to start"}`); - continue; - } - // DRAIN HERE, not in the runner. The sweep is sequential by design and the - // queue would serialize these anyway — awaiting makes that explicit and - // lets the loop re-plan against real results. It is the CALLER's property, - // which is why the runner returns the stream instead of consuming it: a - // stage card renders that same stream as a live log, and a runner that - // drained would leave every card's log dead. - await drainStream(outcome.stream); - live.channelJobId = null; - didWork = true; - } - - const moved = lastReachable === null || lastReachable > reachable; - barrenPasses = moved ? 0 : barrenPasses + 1; - lastReachable = reachable; - - if (barrenPasses >= MAX_BARREN_PASSES) { - onLog( - `Pass ${pass} made no progress for ${MAX_BARREN_PASSES} passes in a row ` + - `(${reachable.toLocaleString()} still outstanding). Stopping rather than ` + - `spinning — check the engine and the job logs, then restart the sweep.`, - ); - return; - } - - if (!didWork || !moved) { - onLog( - `Pass ${pass} made no progress; 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 }, - ); - }); -} - -// Arm and start the corpus-wide sweep. Persists `sweepEnabled` WITH its scope so -// a restart resumes exactly the run that was armed (see -// resumeBackfillSweepIfEnabled). -export async function startBackfillSweep( - opts: BackfillSweepOptions = {}, -): Promise<string | null> { - const paths = opts.paths ?? getPaths(); - const running = getBackfillSweepJobId(); - if (running) return running; - - const settings = getSettings(); - const channelScope = opts.channelSlugs ?? settings.backfill.sweepChannels; - const kindScope = opts.kindIds ?? settings.backfill.sweepKinds; - const same = (a: string[], b: string[]) => - a.join("\u0000") === b.join("\u0000"); - if ( - !settings.backfill.sweepEnabled || - !same(settings.backfill.sweepChannels, channelScope) || - !same(settings.backfill.sweepKinds, kindScope) - ) { - // AWAITED, and it matters: the loop reads `sweepEnabled` at the top of its - // first pass, so an un-awaited write races it and the sweep quits - // immediately with "disarmed in settings" — refusing to start at all. - await writeSettings({ - ...settings, - backfill: { - ...settings.backfill, - sweepEnabled: true, - sweepChannels: channelScope, - sweepKinds: kindScope, - }, - }); - } - - const live: SweepLive = { - jobId: "", - startedAt: Date.now(), - channelJobId: null, - }; - const result = await runManagedFunction({ - kind: BACKFILL_SWEEP_KIND, - // Empty key — see the header. The orchestrator does no work itself; it waits - // on per-channel jobs that serialize on BACKFILL_QUEUE. - queueKey: "", - paths, - fn: async (onLog, signal, _setProgress, ctx) => { - live.jobId = ctx.jobId; - try { - await runSweepLoop( - paths, - kindScope.length > 0 ? kindScope : undefined, - channelScope.length > 0 ? channelScope : undefined, - live, - onLog, - signal, - ctx.drainSignal, - ); - } finally { - if (getSingleton().live === live) getSingleton().live = null; - } - }, - }); - if (!result.ok) return null; - live.jobId = result.jobId; - getSingleton().live = live; - return result.jobId; -} - -// Boot hook. Re-launching is safe rather than merely convenient: the sweep -// stores no cursor and the batch re-derives eligibility from disk on every pull, -// so a resumed sweep re-does exactly zero work. "Is a sweep running" is process -// state, and process state is what a restart destroys. -export async function resumeBackfillSweepIfEnabled( - paths: Paths = getPaths(), -): Promise<void> { - if (!getSettings().backfill.sweepEnabled) return; - await startBackfillSweep({ paths }); -} - -// Disarm and stop. -export async function stopBackfillSweep(): Promise<boolean> { - const settings = getSettings(); - // Note the second and third clauses: guarding on `sweepEnabled` alone means a - // stop on an already-disarmed sweep skips the write entirely and leaves the - // SCOPE behind — which is how a stale scope survives to silently narrow the - // next sweep. That was measured on the digest sweep, not theorised. - if ( - settings.backfill.sweepEnabled || - settings.backfill.sweepChannels.length > 0 || - settings.backfill.sweepKinds.length > 0 - ) { - await writeSettings({ - ...settings, - backfill: { - ...settings.backfill, - sweepEnabled: false, - sweepChannels: [], - sweepKinds: [], - }, - }); - } - const live = getSingleton().live; - if (!live) return false; - // Drain BOTH: the orchestrator, and the per-channel job it is waiting on. - // Draining only the orchestrator makes "stop" mean "after the current - // channel", while the button sits there looking wedged. A drain is still - // graceful — the video in flight completes and nothing part-done is discarded. - if (live.channelJobId) getRegistry().requestDrain(live.channelJobId); - return getRegistry().requestDrain(live.jobId); -} diff --git a/common/controller/digestSweep.ts b/common/controller/digestSweep.ts @@ -1,368 +0,0 @@ -// The corpus-wide digest sweep: one launcher for the thing that runs for weeks. -// -// Before this, "sweep the corpus" meant clicking Digest on 63 channel pages by -// hand and losing the whole thing to a server restart. Everything underneath — -// the per-channel batch, resume-by-re-deriving-from-disk, the pause, the GPU -// yield — already existed and was already correct. What was missing was -// something to walk the channels. -// -// Three shapes are deliberate: -// -// 1. IT LAUNCHES THE EXISTING PER-CHANNEL JOB, in sequence. It does not -// introduce a batch over videos. operationBatch.ts is explicit that one job -// per video would be 119,600 jobs against a 100-record registry and a -// 500-record log, evicting the history of the very run it is recording. -// 63 sequential channel jobs is the same contract, driven. -// -// 2. IT STORES NO CURSOR. Ordering is recomputed from -// `buildDigestSweepPlan` each pass and eligibility is re-derived from disk -// inside every batch, so a restart, a newly finished transcription, a -// human confirming a duplicate cluster, and a settings change that -// invalidates the corpus are all just visible on the next pass. A frozen -// work-list would be wrong within hours of a multi-week run starting. -// -// 3. IT LOOPS UNTIL THE PLAN IS EMPTY, rather than passing over the channel -// list once. A batch can legitimately leave work behind — the spend cap, -// a drain, a video whose transcript landed mid-pass — and a one-pass sweep -// would report itself finished with the corpus unfinished. -// -// Heaviest channel first, by remaining AUDIO-HOURS. Cost is audio, not videos: -// HasanAbiVODs3 is 8,331 audio-hours and outweighs every duplicate mirror in -// the corpus combined, so leaving it for last is how a sweep spends 70 days -// looking nearly done. - -import type { Paths } from "../lib/paths"; -import { getPaths } from "../lib/paths"; -import { getSettings, writeSettings } from "../lib/settings"; -import { getRegistry } from "../jobs/registry"; -import { runManagedFunction } from "../jobs/streamCommand"; -import { drainStream } from "../jobs/drainStream"; -import { - audioHours, - buildDigestSweepPlan, - chunksPerAudioHour, - sweepDays, - type DigestSweepPlan, -} from "./digestPlan"; -import type { DigestLaneChoice } from "./digestTarget"; - -export const DIGEST_SWEEP_KIND = "digest-sweep"; - -// How long the loop waits before recomputing the plan when a pass did no work. -// Generous: nothing here is latency-sensitive and recomputing the plan reads -// every channel's sidecars. -const IDLE_POLL_MS = 60_000; - -// Give up after this many consecutive passes that move no work. Three, not one: -// a single barren pass is legitimate (every remaining video failed a guard this -// time round and may not next time), but three in a row is a broken engine or a -// corpus the sweep cannot make progress on, and both want a human. -const MAX_BARREN_PASSES = 3; - -type SweepLive = { - jobId: string; - startedAt: number; - // The per-channel job the sweep is currently waiting on. Tracked so a stop - // can drain THAT too — without it, "stop" means "after the current channel - // finishes", and the heaviest channel is 8.7 days long. - channelJobId: string | null; -}; -type SweepSingleton = { live: SweepLive | null }; - -declare global { - // eslint-disable-next-line no-var - var __yttDigestSweep__: SweepSingleton | undefined; -} - -function getSingleton(): SweepSingleton { - if (!globalThis.__yttDigestSweep__) { - globalThis.__yttDigestSweep__ = { live: null }; - } - return globalThis.__yttDigestSweep__; -} - -export function getDigestSweepJobId(): string | null { - const live = getSingleton().live; - if (!live) return null; - return getRegistry().get(live.jobId)?.status === "running" - ? live.jobId - : null; -} - -export type DigestSweepOptions = { - paths?: Paths; - lane?: DigestLaneChoice; - // Restrict the sweep to these channels. Absent = the whole corpus. - channelSlugs?: string[]; -}; - -// The per-channel runner lives in controller/operationJobs.ts now, beside the -// backfill one. It no longer drains — see the call site below for where that -// moved to, and why. -import { runDigestChannelJob } from "./operationJobs"; - -async function runSweepLoop( - paths: Paths, - lane: DigestLaneChoice, - channelSlugs: string[] | undefined, - live: SweepLive, - onLog: (msg: string) => void, - signal: AbortSignal, - drainSignal: AbortSignal, -): Promise<void> { - let pass = 0; - // Consecutive passes that left the remaining work UNCHANGED. This is the real - // termination guard, and it has to measure the plan rather than whether jobs - // started: an unreachable engine makes every channel job start normally and - // then throw inside its body, which looks exactly like progress from out here. - // Without this the loop would spin through all 63 channels as fast as they can - // fail, forever, evicting the job registry with its own wreckage. - let barrenPasses = 0; - let lastGenerateSeconds: number | null = null; - for (;;) { - if (signal.aborted || drainSignal.aborted) return; - // The operator turned the sweep off: stop cleanly rather than being - // cancelled, so the job ends "done" and the queue is released. - if (!getSettings().digest.sweepEnabled) { - onLog("Sweep disarmed in settings — stopping."); - return; - } - - pass++; - // Reach, resolved per pass rather than at launch: an operator who switches - // to newest-first mid-sweep should see it take effect on the next plan, not - // only after stopping and re-arming a multi-week run. - // - // Only "corpus" reaches the CHANNEL order; "channel" leaves the plan - // heaviest-first and lets each per-channel batch apply the order itself - // (it reads settings.digest.recencyOrder on its own). - const digestNow = getSettings().digest; - const planOrder = - digestNow.recencyReach === "corpus" ? digestNow.recencyOrder : "listed"; - const plan: DigestSweepPlan = await buildDigestSweepPlan({ - paths, - lane, - channelSlugs, - order: planOrder, - onLog, - }); - const work = plan.channels.filter((c) => c.generateSeconds > 0); - if (work.length === 0) { - onLog( - `Pass ${pass}: nothing left to generate — ${plan.fresh.videos.toLocaleString()} video(s) already digested at the current identity.`, - ); - return; - } - onLog( - `Pass ${pass}: ${work.length} channel(s), ` + - `${plan.generateChunks.toLocaleString()} chunks / ` + - `${audioHours(plan.generateSeconds).toFixed(0)} audio-hours to generate ` + - `(~${sweepDays(plan.generateChunks).toFixed(1)} days at the measured rate, ` + - `${chunksPerAudioHour(plan.generateChunks, plan.generateSeconds).toFixed(2)} chunks/audio-h), ` + - `${audioHours(plan.sharedSeconds).toFixed(0)} audio-hours covered by cluster sharing` + - (planOrder === "listed" - ? ", heaviest channel first." - : `, ${planOrder} channel first (${work[0].channelSlug} @ ${ - planOrder === "newest" - ? work[0].newestPending - : work[0].oldestPending - }).`), - ); - - let didWork = false; - for (const channel of work) { - if (signal.aborted || drainSignal.aborted) return; - if (!getSettings().digest.sweepEnabled) { - onLog("Sweep disarmed in settings — stopping."); - return; - } - onLog( - `→ ${channel.channelSlug}: ${channel.generateChunks.toLocaleString()} chunks / ` + - `${audioHours(channel.generateSeconds).toFixed(0)} audio-hours ` + - `(~${sweepDays(channel.generateChunks).toFixed(1)} days).`, - ); - const outcome = await runDigestChannelJob({ - paths, - channelSlug: channel.channelSlug, - lane, - // Behind anything an operator clicks by hand: a sweep is weeks long and - // must never make a deliberate single-channel run wait for it. - background: true, - onStarted: (jobId) => { - live.channelJobId = jobId; - }, - }); - if (!outcome.ok) { - live.channelJobId = null; - // One channel failing to START is not the sweep failing. - onLog(`!! ${channel.channelSlug}: ${outcome.error ?? "failed to start"}`); - continue; - } - // DRAIN HERE, not in the runner. The sweep is sequential by design (one - // GPU) and the queue would serialize these anyway — awaiting makes that - // explicit and lets the loop re-plan against real results. It belongs to - // the CALLER: the editor's digest card renders this same stream as a live - // log, and a runner that drained would have consumed it first. - await drainStream(outcome.stream); - live.channelJobId = null; - didWork = true; - } - - // Did the pass actually MOVE anything? Compared against the previous pass's - // remaining work, because that is the only measure a failing engine cannot - // fake. Some tolerance: a pass that shifts less than a minute of audio is - // noise, not progress. - const moved = - lastGenerateSeconds === null || - lastGenerateSeconds - plan.generateSeconds > 60; - barrenPasses = moved ? 0 : barrenPasses + 1; - lastGenerateSeconds = plan.generateSeconds; - - if (barrenPasses >= MAX_BARREN_PASSES) { - onLog( - `Pass ${pass} made no progress for ${MAX_BARREN_PASSES} passes in a row ` + - `(${audioHours(plan.generateSeconds).toFixed(0)} audio-hours still outstanding). ` + - `Stopping rather than spinning — check the engine and the job logs, then restart the sweep.`, - ); - return; - } - - if (!didWork || !moved) { - // Nothing started, or nothing moved. Back off before re-planning: a tight - // retry against a down engine helps nobody and buries its own diagnosis. - onLog( - `Pass ${pass} made no progress; 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 }, - ); - }); -} - -// Arm and start the corpus-wide sweep. Persists `sweepEnabled` so a restart -// resumes it (see resumeDigestSweepIfEnabled). -export async function startDigestSweep( - opts: DigestSweepOptions = {}, -): Promise<string | null> { - const paths = opts.paths ?? getPaths(); - const running = getDigestSweepJobId(); - if (running) return running; - - const settings = getSettings(); - // The SCOPE is persisted with the flag, not just the flag. The boot hook - // re-launches from settings alone, so arming without recording the scope - // would resurrect a deliberately-bounded run as a corpus-wide one. - const scope = opts.channelSlugs ?? settings.digest.sweepChannels; - if ( - !settings.digest.sweepEnabled || - settings.digest.sweepChannels.join("\u0000") !== scope.join("\u0000") - ) { - // AWAITED, and it matters: the loop reads `sweepEnabled` at the top of its - // very first pass, so an un-awaited write races it and the sweep quits - // immediately with "disarmed in settings" — refusing to start at all. - await writeSettings({ - ...settings, - digest: { - ...settings.digest, - sweepEnabled: true, - sweepChannels: scope, - }, - }); - } - - const live: SweepLive = { - jobId: "", - startedAt: Date.now(), - channelJobId: null, - }; - const lane = opts.lane ?? "local"; - const result = await runManagedFunction({ - kind: DIGEST_SWEEP_KIND, - // Empty key: the orchestrator itself does no work, it waits on per-channel - // jobs that serialize on the digest queue. Putting it ON that queue would - // deadlock — it would hold the only slot while waiting for a job that needs - // the same slot. - queueKey: "", - paths, - fn: async (onLog, signal, _setProgress, ctx) => { - live.jobId = ctx.jobId; - try { - await runSweepLoop( - paths, - lane, - scope.length > 0 ? scope : undefined, - live, - onLog, - signal, - ctx.drainSignal, - ); - } finally { - if (getSingleton().live === live) getSingleton().live = null; - } - }, - }); - if (!result.ok) return null; - live.jobId = result.jobId; - getSingleton().live = live; - return result.jobId; -} - -// Boot hook. The auto-transcribe/auto-download runners are restarted at server -// start the same way; digests were not, so a restart silently ended a sweep -// that had been running for days and nothing said so. -export async function resumeDigestSweepIfEnabled( - paths: Paths = getPaths(), -): Promise<void> { - if (!getSettings().digest.sweepEnabled) return; - await startDigestSweep({ paths }); -} - -// Disarm and stop. Persisting the flag is the point: without it a restart would -// resurrect a sweep the operator had just stopped. -export async function stopDigestSweep(): Promise<boolean> { - const settings = getSettings(); - // Note the second clause: guarding on `sweepEnabled` alone means a stop on an - // already-disarmed sweep skips the write entirely and leaves the SCOPE behind - // — which is how a stale scope survives to narrow the next sweep. Measured: - // after a stop, settings still read sweepChannels ["teamrcn"]. - if ( - settings.digest.sweepEnabled || - settings.digest.sweepChannels.length > 0 - ) { - // Awaited so a restart cannot resurrect a sweep the operator just stopped. - await writeSettings({ - ...settings, - digest: { - ...settings.digest, - sweepEnabled: false, - // Cleared too. A leftover scope is worse than no scope: the next - // "Start Digest Sweep" would silently cover only the channels of a - // sweep somebody stopped weeks ago, and report itself finished. - sweepChannels: [], - }, - }); - } - const live = getSingleton().live; - if (!live) return false; - // Drain BOTH: the orchestrator, and the per-channel job it is waiting on. - // Draining only the orchestrator would mean "stop" waits for the current - // channel to finish — up to 8.7 days on the heaviest one — while the button - // sat there looking wedged. A drain is still graceful: the video in flight - // completes, the batch stops taking new ones, and nothing part-generated is - // thrown away. - if (live.channelJobId) getRegistry().requestDrain(live.channelJobId); - return getRegistry().requestDrain(live.jobId); -} diff --git a/common/controller/laneForOperation.test.ts b/common/controller/laneForOperation.test.ts @@ -1,12 +1,12 @@ -// The REAL lane resolver, not the stub arbiter.test.ts plans units with. +// The REAL lane resolver, against real settings on disk. // // Run with: node_modules/.bin/tsx --test common/controller/laneForOperation.test.ts // // Its own file because it needs a settings seam: laneForOperation calls // getSettings(), and getPaths() memoizes its first answer at module scope, so -// the env must be set before anything imports the module under test. -// arbiter.test.ts is deliberately settings-free — it passes `laneOf` in — and -// these cases are exactly the ones a stub would answer vacuously. +// the env must be set before anything imports the module under test. Every other +// test of this machinery passes a `laneOf` stub in, and these cases are exactly +// the ones a stub would answer vacuously. // // Named for the function rather than for controller/operationLane.ts, which is // the one-function module it now lives in. @@ -43,7 +43,7 @@ function writeSettings(settings: Record<string, unknown>): void { test("digest's lane follows remoteEnabled, live", () => { // THE POINT OF laneFor. `digest.lane` is a DECLARATION and can only carry the // default (the local, GPU-bound queue). Which queue a run actually takes is a - // configuration choice, and the arbiter must reserve by the live answer or a + // configuration choice, and a dispatcher must reserve by the live answer or a // remote digest lands on the local key — where registry.ts's concurrency-1 // would serialize it behind GPU work it competes with for nothing. writeSettings({ digest: { remoteEnabled: false } }); @@ -61,7 +61,7 @@ test("digest's lane follows remoteEnabled, live", () => { test("diarization resolves through its laneFor too", () => { // Regression guard on the shape of the resolver rather than on diarization: // before this, laneForOperation special-cased digest and returned `.lane` for - // everything else, so diarization's laneFor was invisible to the arbiter. The + // everything else, so diarization's laneFor was invisible to callers. The // queue key is the same either way by design (two diarizations must never run // at once), so `contendsFor` is what shows the live answer got through. writeSettings({ @@ -80,11 +80,11 @@ test("diarization resolves through its laneFor too", () => { test("a KNOWN external operation still has no lane", () => { // download and transcription are in operationCatalog() — they have to be, or // digest.dependsOn = ["transcription"] names nothing — but they are dispatched - // by their OWN runners. laneForOperation asks getOperation, not the - // catalog, precisely so they come back null and planArbiterUnits skips them. + // by their OWN runners off their own queue keys. laneForOperation asks + // getOperation, not the catalog, precisely so they come back null. // - // A catalog lookup would hand them a lane and the arbiter would start a - // backfill channel job for work no backfill kind can do. That is a different + // A catalog lookup would hand them a lane, and a caller would reserve the + // backfill queue for work no backfill operation can do. That is a different // failure from an unknown id, and only this case would catch it. writeSettings({}); assert.equal(laneForOperation("download"), null); @@ -97,8 +97,8 @@ test("a switched-off kind still resolves a lane — the off-switch is elsewhere" // Pinning the BOUNDARY, because it is the one an obvious-looking "fix" would // move. This resolver answers "which lane would this run on", not "should it // run": getOperation is a registry lookup and has never consulted - // `enabled`. What keeps a disabled feature out of the arbiter is the `enabled` - // set in runArbiterPass, which projects ids only for kinds allOperations + // `enabled`. What keeps a disabled feature out of dispatch is + // `operationsForLane`, which projects only the operations `allOperations` // returns — so a switched-off kind arrives with an EMPTY id list and produces // no unit, well before a lane is asked for. // diff --git a/common/controller/laneGuards.ts b/common/controller/laneGuards.ts @@ -17,9 +17,9 @@ import { transcriptionActivity } from "./digestYield"; // remoteEnabled fail-fast, the engine probe() fail-fast … and not one of // them is expressible as backfillLimit()'s single scalar." // -// That is the whole reason there are two schedulers for one kind of work. Every +// That is the reason there used to be two schedulers for one kind of work. Every // one of those guards lived inside the digest batch's own closure, where nothing -// else could ask about it — so an arbiter that wanted to dispatch digest would +// else could ask about it — so a second dispatcher wanting to run digest would // have had to re-implement all six, correctly, from memory. They live here now, // and operationBatch calls them, so there is exactly one definition and no second // opinion to drift. diff --git a/common/controller/noCorpusWalkInRenderPaths.test.ts b/common/controller/noCorpusWalkInRenderPaths.test.ts @@ -29,25 +29,22 @@ const EDITOR_APP = path.resolve(HERE, "..", "..", "editor", "app"); // EVERY identifier that reaches the corpus walk, not just the walk itself. // -// The guard grepped for one name and that was not enough: `buildBackfillSweepPlan` -// CALLS listChannelStatsFromDisk, so importing it into a page would have walked -// 474,559 files on a 3-second poll and passed this test with room to spare. -// `buildDigestSweepPlan` is the same hazard from the other sweep — an LMDB scan -// plus a per-video freshness check, priced for a run that happens hourly. +// The guard grepped for one name and that was not enough: the retired +// `buildBackfillSweepPlan` CALLED listChannelStatsFromDisk, so importing it into +// a page would have walked 474,559 files on a 3-second poll and passed this test +// with room to spare. `buildDigestSweepPlan` is the same hazard and is still +// here — an LMDB scan plus a per-video freshness check, priced for the +// `bin/digest-plan` ETA rather than for a render. // -// If you are here because you want one of these on a screen: the snapshot-only -// preview is ../controller/sweepPreview.ts, which is the same counting off the -// channel snapshots and is what the /auto-queue console draws. +// If you are here because you want it on a screen: what a lane's console draws +// is `computeLeafPending` (controller/autoRunner.ts), which folds the channel +// SNAPSHOTS the payload has already read. // // THE MATCH IS TEXTUAL AND CONTEXT-BLIND, so a mere MENTION in a comment fails // it too. That is deliberate and not worth softening: a guard that skipped // comments and strings is a guard an offending call can hide behind, and the // cost of the false positive is one reworded comment. -const BANNED = [ - "listChannelStatsFromDisk", - "buildBackfillSweepPlan", - "buildDigestSweepPlan", -]; +const BANNED = ["listChannelStatsFromDisk", "buildDigestSweepPlan"]; async function walk(dir: string): Promise<string[]> { const out: string[] = []; diff --git a/common/controller/operationBatch.test.ts b/common/controller/operationBatch.test.ts @@ -26,80 +26,48 @@ import type { Operation, OperationClassification } from "../lib/operations"; // and a test that needs a GPU, a registry and a worker pool is a test nobody // runs before shipping. -test("weight 0 is idle-only: full slots when the primary lane is quiet", () => { - assert.equal( - backfillLimit({ weight: 0, slots: 4, primaryBusy: false }), - 4, - ); - assert.equal(backfillLimit({ weight: 0, slots: 1, primaryBusy: false }), 1); +test("idle-only means full slots when the primary lane is quiet", () => { + assert.equal(backfillLimit({ idleOnly: true, slots: 4, primaryBusy: false }), 4); + assert.equal(backfillLimit({ idleOnly: true, slots: 1, primaryBusy: false }), 1); }); -test("weight 0 stands aside completely while the primary lane works", () => { +test("idle-only stands aside completely while the primary lane works", () => { // Zero is a HOLD, not a stop: runPool idle-waits at a zero limit rather than // finishing, so the lane resumes the moment transcription is free without // re-deriving anything. That is the whole reason this is a limit and not a // scheduler. - assert.equal(backfillLimit({ weight: 0, slots: 4, primaryBusy: true }), 0); - assert.equal(backfillLimit({ weight: 0, slots: 16, primaryBusy: true }), 0); -}); - -test("a positive weight is a guaranteed share, busy or not", () => { - // The point of a non-zero weight: the primary lane being busy no longer parks - // the backfill. - assert.equal(backfillLimit({ weight: 0.5, slots: 4, primaryBusy: true }), 2); - assert.equal(backfillLimit({ weight: 0.5, slots: 4, primaryBusy: false }), 2); - assert.equal(backfillLimit({ weight: 1, slots: 4, primaryBusy: true }), 4); + assert.equal(backfillLimit({ idleOnly: true, slots: 4, primaryBusy: true }), 0); + assert.equal(backfillLimit({ idleOnly: true, slots: 16, primaryBusy: true }), 0); }); -test("a small weight is a SLOW lane, never a stopped one", () => { - // floor(4 * 0.25) is 1, and floor(1 * 0.25) is 0 — which without the floor - // would silently turn "a quarter of the machine" into "never runs", and look - // exactly like a wedge. - assert.equal(backfillLimit({ weight: 0.25, slots: 4, primaryBusy: true }), 1); - assert.equal(backfillLimit({ weight: 0.25, slots: 1, primaryBusy: true }), 1); - assert.equal(backfillLimit({ weight: 0.01, slots: 8, primaryBusy: true }), 1); +test("a run that is not idle-only keeps its slots, busy or not", () => { + // The `weight` scalar retired in slice 1.3 used to answer this with a share. + // The question it was really being asked is "does this run contend with + // transcription", which the operation's declared `contendsFor` answers — so a + // network-bound attribution run keeps its slots rather than parking behind a + // GPU it is not competing for. + assert.equal(backfillLimit({ idleOnly: false, slots: 4, primaryBusy: true }), 4); + assert.equal(backfillLimit({ idleOnly: false, slots: 4, primaryBusy: false }), 4); }); -test("no slots means no work, whatever the weight", () => { - assert.equal(backfillLimit({ weight: 0, slots: 0, primaryBusy: false }), 0); - assert.equal(backfillLimit({ weight: 1, slots: 0, primaryBusy: false }), 0); +test("no slots means no work, idle-only or not", () => { + assert.equal(backfillLimit({ idleOnly: true, slots: 0, primaryBusy: false }), 0); + assert.equal(backfillLimit({ idleOnly: false, slots: 0, primaryBusy: false }), 0); }); -test("the shipped default is idle-only", () => { - // Stated as a test because it is a promise the feature makes: catch-up work on - // a corpus that already exists must never slow down new arrivals. +test("the shipped default is off, and holds no disk", () => { + // Stated as a test because these are promises the feature makes: nothing runs + // until an operator says so, and nothing re-downloads media without an + // explicit opt-in. const d = defaultBackfill(); - assert.equal(d.weight, 0); assert.equal(d.enabled, false); assert.equal(d.allowRedownload, false); - assert.equal(backfillLimit({ ...d, slots: d.concurrency, primaryBusy: true }), 0); -}); - -test("a hand-edited weight is clamped rather than rejected", () => { - // settings.json is hand-editable. A 5 means "as much as possible", and reading - // it as the idle-only 0 would be the opposite of the intent. - assert.equal(sanitizeBackfill({ weight: 5 }).weight, 1); - assert.equal(sanitizeBackfill({ weight: -3 }).weight, 0); - assert.equal(sanitizeBackfill({ weight: "half" }).weight, 0); - assert.equal(sanitizeBackfill({ weight: 0.25 }).weight, 0.25); - // Concurrency shares clampPositiveInt's floor of 1, so a 0 cannot silently - // park the lane either. + assert.equal(d.concurrency, 1); + // Concurrency shares clampPositiveInt's floor of 1, so a hand-edited 0 cannot + // silently park the lane. assert.equal(sanitizeBackfill({ concurrency: 0 }).concurrency, 1); }); -test("a stale sweep scope survives sanitization as a list of slugs", () => { - // The scope is persisted WITH the flag, so it has to round-trip: junk entries - // are dropped, real ones kept. - assert.deepEqual( - sanitizeBackfill({ sweepChannels: ["a", "", 7, "b"], sweepKinds: ["diarization"] }), - { - ...defaultBackfill(), - sweepChannels: ["a", "b"], - sweepKinds: ["diarization"], - }, - ); -}); - // --------------------------------------------------------------------------- // The candidate pull's dispatch decision. // @@ -296,10 +264,10 @@ test("the backfill lane holds when the lane switch is off, and says so", () => { assert.equal(verdict.hold?.reason, "paused"); }); -test("an enabled backfill lane with a weight runs, primary busy or not", () => { +test("an enabled backfill lane with no GPU-bound operation runs at its slots", () => { const verdict = laneLimit( settingsWith({ - backfill: { ...defaultBackfill(), enabled: true, weight: 1, concurrency: 3 }, + backfill: { ...defaultBackfill(), enabled: true, concurrency: 3 }, }), { lane: "backfill", @@ -324,7 +292,7 @@ test("the GPU carve-out is keyed on contendsFor, not on laneFor existing", () => // `lane: { contendsFor: "gpu" }` and NO laneFor is caught, and one with a // laneFor resolving to CPU is not. const settings = settingsWith({ - backfill: { ...defaultBackfill(), enabled: true, weight: 1, concurrency: 4 }, + backfill: { ...defaultBackfill(), enabled: true, concurrency: 4 }, }); const cpuOnly = laneLimit(settings, { lane: "backfill", @@ -334,7 +302,7 @@ test("the GPU carve-out is keyed on contendsFor, not on laneFor existing", () => llmActive: 0, unitActive: 0, }); - // A weight is a CPU concept, so a CPU-bound run keeps its guaranteed share. + // Sharing is a CPU concept, so a CPU-bound run keeps its slots. assert.equal(cpuOnly.limit, 4); const gpuBound = laneLimit(settings, { @@ -351,7 +319,7 @@ test("the GPU carve-out is keyed on contendsFor, not on laneFor existing", () => assert.equal(gpuBound.limit, 4, "idle: nothing is transcribing in this process"); const weightless = laneLimit( settingsWith({ - backfill: { ...defaultBackfill(), enabled: true, weight: 1, concurrency: 4 }, + backfill: { ...defaultBackfill(), enabled: true, concurrency: 4 }, }), { lane: "backfill", @@ -362,8 +330,8 @@ test("the GPU carve-out is keyed on contendsFor, not on laneFor existing", () => unitActive: 0, }, ); - // Same shape as backfillLimit({weight: 0, primaryBusy: false}) — the weight - // was zeroed by the carve-out, and an idle primary still gives full slots. + // Same shape as backfillLimit({idleOnly: true, primaryBusy: false}) — the + // carve-out made it idle-only, and an idle primary still gives full slots. assert.equal(weightless.limit, 4); }); diff --git a/common/controller/operationBatch.ts b/common/controller/operationBatch.ts @@ -186,24 +186,28 @@ export function candidateAction( // established and for the same reason: this decides whether a multi-day lane // runs at all, and a test that needs the hardware is a test nobody runs. // -// weight 0 (the default) — IDLE-ONLY. Full slots when the primary lane is -// quiet, zero while it works. Catch-up is by definition work on a corpus -// that already exists, so it must never slow down new arrivals. Zero is a -// HOLD, not a stop: runPool idle-waits, so the lane resumes the moment the -// primary is free, having re-derived nothing. -// weight > 0 — a guaranteed share, floored at 1. The floor is the point: a -// small weight should mean a slow lane, not a stopped one, and -// floor(1 * 0.25) is 0. +// idleOnly — full slots when the primary transcription lane is quiet, zero +// while it works. Zero is a HOLD, not a stop: runPool idle-waits, so the +// lane resumes the moment the primary is free, having re-derived nothing. +// otherwise — the lane's own slots, whatever transcription is doing. +// +// THE INPUT USED TO BE `settings.backfill.weight`, a 0..1 scalar that meant +// "idle-only" at 0 and "a guaranteed share, floored at 1" above it. Slice 1.3 +// retired it: the yield is the OPERATION'S DECLARED `contendsFor` (a run that +// could dispatch GPU work stands aside; one that contends for the network does +// not), and the share is `concurrency` and the lane policy's `maxWorkers`. One +// scalar standing for both a resource question and a size question is how a +// network-bound attribution run came to park itself behind a GPU transcription +// it was not competing with. export function backfillLimit(opts: { - weight: number; + idleOnly: boolean; slots: number; primaryBusy: boolean; }): number { const slots = Math.max(0, Math.floor(opts.slots)); if (slots === 0) return 0; - if (opts.weight <= 0) return opts.primaryBusy ? 0 : slots; - const share = Math.floor(slots * Math.min(1, opts.weight)); - return Math.max(1, share); + if (opts.idleOnly) return opts.primaryBusy ? 0 : slots; + return slots; } // The live, per-RUN facts a limit cannot read from settings: leases this run @@ -341,14 +345,12 @@ function backfillLaneLimit( // parked. The cost is that a backfill and a digest can overlap on the CPU; the // benefit is that neither can deadlock the other. const activity = transcriptionActivity(); - // A GUARANTEED SHARE IS A CPU CONCEPT. `weight > 0` means "keep a slice of the - // cores running even while transcription works", which is reasonable when the - // contended resource is cores and divisible. It is not available for VRAM: - // sortformer on Vulkan holds ~4.4 GB of the same 8 GB card parakeet is using, - // so "a small share" is not a slower run, it is an out-of-memory failure of - // whichever lane allocates second. + // A RUN THAT COULD DISPATCH GPU WORK IS IDLE-ONLY. Sharing is a CPU concept: + // it is reasonable when the contended resource is cores and divisible, and it + // is not available for VRAM — sortformer on Vulkan holds ~4.4 GB of the same + // 8 GB card parakeet is using, so "a small share" is not a slower run, it is + // an out-of-memory failure of whichever lane allocates second. // - // So a run that could dispatch GPU work is idle-only whatever the weight says. // The cost is bluntness — one such operation makes the whole run idle-only, // including its cheap CPU-bound siblings — and that is the deliberate // direction to be wrong in, since the alternative is an OOM mid-sweep. @@ -366,7 +368,7 @@ function backfillLaneLimit( laneYieldsToTranscription(laneForOperation(op.id) ?? op.lane), ); const local = backfillLimit({ - weight: gpuBound ? 0 : cfg.weight, + idleOnly: gpuBound, slots: cfg.concurrency, primaryBusy: activity.busy, }); @@ -488,10 +490,10 @@ export function resetDurationMemo(): void { // reachable ids are what it draws, so it puts them behind everything it can // price instead of letting `null` sort to the front. // -// `direction` is a PARAMETER rather than a constant because today's live value -// comes from `settings.digest.recencyOrder` ("newest"), and slice 1.3 retires -// that field into the lane's tree. Until then the caller passes what settings -// say; after it, the lane's own `order` does. +// `direction` is a PARAMETER rather than a constant so this stays a pure +// comparator factory. Its one caller fixes it at "newest" — the value +// `settings.digest.recencyOrder` carried on the live corpus before slice 1.3 +// retired that field. export function cheapestComparator(opts: { duration: (id: string) => number | null | undefined; recency?: ((a: string, b: string) => number) | null; @@ -1859,13 +1861,21 @@ export async function runOperationBatch( return result; } -// The lane's configured recency order, resolved INSIDE the batch so an armed -// sweep, a hand-clicked channel run and an ids-scoped run all inherit the same -// setting instead of three call sites each remembering to pass it. +// The lane's configured recency order, resolved INSIDE the batch so a +// hand-clicked channel run, an ids-scoped run and the lane runner all inherit +// the same setting instead of three call sites each remembering to pass it. +// +// IT COMES OFF THE LANE'S POLICY NOW — `settings.digest.recencyOrder` and +// `settings.backfill.order` retired with the sweeps in slice 1.3 — so the +// per-channel button and the lane runner read one field. The digest lane's own +// default is `cheapest`, which is a composition (shortest inside a day, newest +// day first) that this batch expresses as its own `digestOrder` plus a date +// sort; the date half of it is newest-first, exactly as the retired +// `recencyOrder` was on the live corpus. function laneOrderFor(run: OperationRun): AutoQueueOrder { - return run.digest - ? run.settings.digest.recencyOrder - : run.settings.backfill.order; + const order = run.settings.autoQueue[run.digest ? "digest" : "backfill"].order; + if (order === "cheapest") return run.digest ? "newest" : "listed"; + return order ?? "listed"; } function batchHeadline( diff --git a/common/controller/operationJobs.ts b/common/controller/operationJobs.ts @@ -1,8 +1,8 @@ // ONE WAY TO RUN ONE OPERATION OVER ONE CHANNEL. // -// Before this file there were four, and no two agreed. The sweep had a -// per-channel runner; the arbiter called it for backfill and a different one for -// digest; the editor's stage cards each had a hand-written copy of the same +// Before this file there were four, and no two agreed. Each sweep had a +// per-channel runner; the arbiter called one for backfill and a different one +// for digest; the editor's stage cards each had a hand-written copy of the same // runManagedFunction call; the /channels group buttons called those copies. The // copies had drifted in the way copies do — the editor's backfill summary line // printed `deferred` and `blocked`, the sweep's printed `skipped`, and neither @@ -12,14 +12,13 @@ // THE SPLIT THAT MATTERS: these runners START a job and RETURN ITS STREAM. They // do not drain it. // -// Draining is what a SWEEP wants — it is sequential by design, and awaiting is -// what lets its loop re-plan against real results — and it is exactly what a -// stage card must not do: the editor consumes the returned stream client-side to -// draw a live log, and a runner that drained would have consumed it first and -// left every card's log dead. So the drain lives at the three call sites that -// want it (backfillSweep, digestSweep, arbiter), which is also where "sequential -// by design" is a true statement about the caller rather than a property -// silently imposed on everyone. +// Draining is what a SEQUENTIAL caller wants — awaiting is what lets a loop +// re-plan against real results — and it is exactly what a stage card must not +// do: the editor consumes the returned stream client-side to draw a live log, +// and a runner that drained would have consumed it first and left every card's +// log dead. So the drain lives at the call sites that want it, which is also +// where "sequential by design" is a true statement about the caller rather than +// a property silently imposed on everyone. // // AND `onDone`, because common/ CANNOT import next/cache. revalidatePath is the // editor's business; the hook is how a server action gets it to fire at job end @@ -59,8 +58,8 @@ function isDigestOrder(v: unknown): v is DigestOrder { export type BackfillChannelJobOptions = { paths: Paths; channelSlug: string; - // Absent = every enabled lane kind. The arbiter passes exactly one; a stage - // card's per-row "run only this" passes one; the group button passes none, + // Absent = every enabled lane kind. A stage card's per-row "run only this" + // passes one; the lane runner passes one; the group button passes none, // because that station IS the lane rather than one operation on it. kindIds?: string[]; // Only these videos (intersected with disk by the batch). A stage card passes @@ -69,8 +68,8 @@ export type BackfillChannelJobOptions = { // An operator's queue override from the card's QueueControl. `undefined` is // "no control wired up" and takes the lane's own key; "" means immediate. queueKey?: string; - // Behind anything an operator asks for by hand. The sweeps and the arbiter - // set it; a stage card does not. + // Behind anything an operator asks for by hand. A background dispatcher sets + // it; a stage card does not. background?: boolean; onStarted?: (jobId: string) => void; onDone?: () => void; @@ -182,8 +181,8 @@ export async function runDigestChannelJob( const settings = getSettings(); // FAIL FAST, before runManagedFunction, so an off metered lane produces an // error the caller can show rather than a started job that dies. Lives here - // rather than in the editor action because the arbiter can ask for the remote - // lane too, and a check only one caller performs is a check. + // rather than in the editor action because the lane runner can ask for the + // remote lane too, and a check only one caller performs is a check. if (lane === "remote" && !settings.digest.remoteEnabled) { return { ok: false, @@ -302,12 +301,12 @@ export type OperationChannelJobOptions = { // RUN ONE OPERATION OVER ONE CHANNEL, whichever operation it is. // -// This is the seam the arbiter and the /channels group buttons wanted: a caller -// that has an operation id and a channel slug should not also have to know which -// of two runners that id belongs to, nor which queue key to reserve. Both were -// open-coded before — the arbiter with a ternary on DIGEST_OPERATION_ID, the group -// buttons with a per-station map — and both had to be edited in step with the -// registry. +// This is the seam the /channels group buttons wanted: a caller that has an +// operation id and a channel slug should not also have to know which of two +// runners that id belongs to, nor which queue key to reserve. It was open-coded +// twice before — the retired arbiter with a ternary on DIGEST_OPERATION_ID, the +// group buttons with a per-station map — and both had to be edited in step with +// the registry. export function runOperationChannelJob( opts: OperationChannelJobOptions, ): Promise<StreamActionResult> { diff --git a/common/controller/operationLane.ts b/common/controller/operationLane.ts @@ -1,9 +1,11 @@ // WHICH LANE AN OPERATION RUNS ON. One function, one module. // -// Its own file because two things need it and neither may import the other: the -// arbiter plans units against it, and controller/operationJobs.ts reserves a -// queue key with it — while the arbiter calls operationJobs to start the jobs it -// plans. Left in arbiter.ts, that was an import cycle. +// Its own file because it started as the arbiter's, and two things needed it +// while neither could import the other: the arbiter planned units against it and +// controller/operationJobs.ts reserves a queue key with it, while the arbiter +// called operationJobs to start the jobs it planned. The arbiter retired in +// slice 1.3; the module stays where it is, because a queue-key resolver is not +// dispatch's private business. import { getSettings } from "../lib/settings"; import { getOperation, type Lane } from "../lib/operations"; @@ -16,14 +18,13 @@ import { getOperation, type Lane } from "../lib/operations"; // the GPU and contending for cores — and a kind without one has a fixed lane. // This used to special-case digest here, which meant the rule lived in two // places and only digest's copy was live: diarization's laneFor was invisible -// to the arbiter. +// to every caller. // // getOperation, NOT operationCatalog(), and that is deliberate. Today -// getOperation("download") is undefined, so download and transcription get -// no lane and planArbiterUnits skips them — which is correct, because they are -// dispatched by their own runners. A -// catalog lookup would hand them a lane and the arbiter would start a backfill -// channel job for work no backfill kind can do. +// getOperation("download") is undefined, so download and transcription get no +// lane — which is correct, because they are dispatched by their own runners off +// their own queue keys. A catalog lookup would hand them a lane and a caller +// would reserve the backfill queue for work no backfill operation can do. export function laneForOperation(operation: string): Lane | null { const kind = getOperation(operation); if (!kind) return null; diff --git a/common/controller/remoteCapacity.ts b/common/controller/remoteCapacity.ts @@ -7,7 +7,7 @@ // every remote as one slot. This module keeps that answer: the count of the // remote's enabled, non-degraded workers. // -// TTL-CACHED, THE sweepRecency WAY. reconfigure() is synchronous and runs on +// TTL-CACHED. reconfigure() is synchronous and runs on // every batch start and settings save; a network probe cannot live inside it. // So the cache is read synchronously (cachedRemoteSlots), staleness is asked // synchronously (remoteCapacityStale), and the probe itself is fired diff --git a/common/controller/sweepPreview.test.ts b/common/controller/sweepPreview.test.ts @@ -1,321 +0,0 @@ -// Run with: -// pnpm --filter yt-dlp-transcript-common exec tsx --test common/controller/sweepPreview.test.ts -// -// THE ONE THING WORTH TESTING HERE is that the preview and the run agree. The -// console is a commitment device: an operator reads the plan, sees "11,342 to do -// across 3 of 68 channels", and arms a run that may take days. A preview that -// counts differently from the planner is not a cosmetic defect — it is the -// number the decision was made on. -// -// So the first test builds a fixture corpus and asserts buildSweepPreview() -// equals buildBackfillSweepPlan() over it, entry for entry. The rest pin the -// four rules the fold has to get right on its own: an empty scope is EVERY kind -// (not none), an unknown kind id contributes nothing (sanitizeBackfill keeps -// them), a channel with no snapshot is unknown rather than zero, and the order -// is the one controller/planOrder documents. - -import { mkdtempSync, writeFileSync, mkdirSync } from "node:fs"; -import { rm } from "node:fs/promises"; -import os from "node:os"; -import path from "node:path"; -import { test, after } from "node:test"; -import assert from "node:assert/strict"; - -// Set BEFORE anything can call getPaths(), which memoizes its first answer at -// module scope. node:test runs each file in its own process, so this is scoped -// to this file alone. -const ROOT = mkdtempSync(path.join(os.tmpdir(), "sweep-preview-")); -process.env.TRANSCRIPTS_DIR = ROOT; -// SETTINGS_FILE is its own env key — it defaults to the MONOREPO root, not the -// corpus, so without this the fixture would silently resolve its kinds against -// the live install's settings.json. -process.env.SETTINGS_FILE = path.join(ROOT, "settings.json"); - -const { getPaths } = await import("../lib/paths"); -const { buildBackfillSweepPlan } = await import("./backfillSweep"); -const { buildSweepChannelCounts, buildSweepPreview } = await import( - "./sweepPreview" -); -const { foldSweepPlan, sweepPlanTotals } = await import("../lib/sweepPlan"); -const { listChannelBriefs } = await import("./channels"); - -after(() => rm(ROOT, { recursive: true, force: true })); - -const KINDS = ["attribution-diarized", "attribution-text"]; - -// A snapshot entry for one kind. `ids` is EXACTLY the reachable set, which is -// the invariant channelSnapshot.test.ts pins and this file relies on. -function entry(missing: number, missingInput: number, prefix: string) { - return { - missing, - stale: 0, - partial: 0, - missingInput, - deferred: 0, - blocked: 0, - eligible: missing + missingInput + 5, - ids: Array.from({ length: missing }, (_, i) => `${prefix}${i}`), - }; -} - -function writeChannel( - slug: string, - backfill: Record<string, ReturnType<typeof entry>> | null, -): void { - const dir = path.join(ROOT, "channels", slug); - mkdirSync(path.join(dir, "data"), { recursive: true }); - writeFileSync( - path.join(dir, "config.json"), - // `handling` is what parseChannelConfig requires; without it the config is - // null and the channel is invisible to every reader here. - JSON.stringify({ - handling: "transcribe", - name: slug, - url: `https://example.test/${slug}`, - }), - ); - if (backfill === null) return; - writeFileSync( - path.join(dir, "snapshot.json"), - JSON.stringify({ - generatedAt: new Date(0).toISOString(), - totals: { videos: 10, transcribed: 10, downloaded: 10 }, - backfill, - buckets: {}, - undownloadedIds: [], - }), - ); -} - -writeFileSync( - path.join(ROOT, "settings.json"), - JSON.stringify({ - attribution: { enabled: true, textOnlyEnabled: true, diarizedEnabled: true }, - }), -); -writeChannel("alpha", { - "attribution-diarized": entry(4, 1, "a"), - "attribution-text": entry(11, 0, "t"), -}); -writeChannel("beta", { - "attribution-diarized": entry(2, 7, "b"), - "attribution-text": entry(0, 0, "u"), -}); -writeChannel("gamma", { - "attribution-diarized": entry(0, 0, "g"), - "attribution-text": entry(0, 0, "h"), -}); -// No snapshot at all — the case the run WALKS and the preview must not silently -// report as zero. -writeChannel("delta", null); - -test("the preview counts what the run's own planner counts", async () => { - const paths = getPaths(); - const briefs = await listChannelBriefs(paths); - const plan = await buildBackfillSweepPlan({ paths }); - const preview = buildSweepPreview({ briefs, kindIds: KINDS }); - // `delta` has no snapshot, so the planner walks it (and finds an empty data - // dir); the preview reports it unknown. Both therefore say zero for it here — - // what the test pins is that every channel the snapshots DO cover agrees. - assert.deepEqual(preview, plan); - const totals = sweepPlanTotals(preview); - assert.equal(totals.reachable, 4 + 11 + 2); - assert.equal(totals.missingInput, 1 + 7); - assert.equal(totals.channels, 2); - assert.equal(totals.idle, 2); -}); - -test("scoping to one kind is the run's own scope, counted the same way", async () => { - const paths = getPaths(); - const briefs = await listChannelBriefs(paths); - const scope = ["attribution-diarized"]; - const plan = await buildBackfillSweepPlan({ paths, kindIds: scope }); - const preview = buildSweepPreview({ - briefs, - kindIds: KINDS, - scopeKindIds: scope, - }); - assert.deepEqual(preview, plan); - // And the point of the whole console: the expensive lane dropping out is - // visible as a smaller number, not as a different-looking screen. - assert.equal(sweepPlanTotals(preview).reachable, 4 + 2); -}); - -test("an empty scope is EVERY kind, never none", () => { - const counts = buildSweepChannelCounts({ - briefs: [ - { - slug: "alpha", - snapshot: { - backfill: { - "attribution-diarized": entry(4, 0, "a"), - "attribution-text": entry(11, 0, "t"), - }, - }, - } as never, - ], - kindIds: KINDS, - }); - assert.equal(sweepPlanTotals(foldSweepPlan({ counts })).reachable, 15); - assert.equal( - sweepPlanTotals(foldSweepPlan({ counts, kindIds: [] })).reachable, - 15, - ); -}); - -test("an unknown kind id contributes nothing — it does not match everything", () => { - // sanitizeBackfill keeps an id that no longer names a registered kind, so a - // stored scope can name one. Counting it as zero is the safe reading; the - // console says so out loud, and this pins that it cannot silently widen. - const counts = buildSweepChannelCounts({ - briefs: [ - { - slug: "alpha", - snapshot: { backfill: { "attribution-text": entry(11, 0, "t") } }, - } as never, - ], - kindIds: [...KINDS, "operation-that-was-removed"], - }); - assert.equal( - sweepPlanTotals( - foldSweepPlan({ counts, kindIds: ["operation-that-was-removed"] }), - ).reachable, - 0, - ); -}); - -test("a channel with no snapshot is unknown, not zero", () => { - const counts = buildSweepChannelCounts({ - briefs: [{ slug: "delta", snapshot: null }], - kindIds: KINDS, - }); - assert.equal(counts[0].unknown, true); - assert.equal(counts[0].byKind["attribution-text"].reachable, 0); - const reported = buildSweepChannelCounts({ - briefs: [ - { - slug: "gamma", - snapshot: { backfill: { "attribution-text": entry(0, 0, "h") } }, - } as never, - ], - kindIds: KINDS, - }); - // The distinction the flag exists for: both read 0, and only one of them is a - // finding. - assert.equal(reported[0].unknown, false); -}); - -test("dates come from the reachable ids, per kind, and fold as a max of maxima", () => { - const counts = buildSweepChannelCounts({ - briefs: [ - { - slug: "alpha", - snapshot: { - backfill: { - "attribution-diarized": entry(2, 0, "a"), - "attribution-text": entry(2, 0, "t"), - }, - }, - } as never, - ], - kindIds: KINDS, - dates: new Map([ - ["a0", "20240101"], - ["a1", "20240202"], - ["t0", "20260819"], - // t1 deliberately undated: an id nothing could date must not become the - // extreme, in either direction. - ]), - }); - const byKind = counts[0].byKind; - assert.deepEqual( - [byKind["attribution-diarized"].newestPending, byKind["attribution-diarized"].oldestPending], - ["20240202", "20240101"], - ); - assert.deepEqual( - [byKind["attribution-text"].newestPending, byKind["attribution-text"].oldestPending], - ["20260819", "20260819"], - ); - // Both kinds in scope: the union's extremes. - const both = foldSweepPlan({ counts })[0]; - assert.deepEqual( - [both.newestPending, both.oldestPending], - ["20260819", "20240101"], - ); - // Drop the expensive lane and the row re-dates itself — no second lookup. - const one = foldSweepPlan({ counts, kindIds: ["attribution-diarized"] })[0]; - assert.deepEqual( - [one.newestPending, one.oldestPending], - ["20240202", "20240101"], - ); -}); - -test("order is planOrder's rule: newest first, weight as the tiebreak", () => { - const counts = [ - { - channelSlug: "heavy-and-old", - unknown: false, - byKind: { - k: { - reachable: 900, - missingInput: 0, - blocked: 0, - deferred: 0, - eligible: null, - present: null, - newestPending: "20200101", - oldestPending: "20190101", - }, - }, - }, - { - channelSlug: "light-and-fresh", - unknown: false, - byKind: { - k: { - reachable: 3, - missingInput: 0, - blocked: 0, - deferred: 0, - eligible: null, - present: null, - newestPending: "20260819", - oldestPending: "20260819", - }, - }, - }, - { - channelSlug: "undated", - unknown: false, - byKind: { - k: { - reachable: 50, - missingInput: 0, - blocked: 0, - deferred: 0, - eligible: null, - present: null, - newestPending: "", - oldestPending: "", - }, - }, - }, - ]; - assert.deepEqual( - foldSweepPlan({ counts, order: "newest" }).map((e) => e.channelSlug), - // Undated sorts LAST under newest — the same thing makeRecencyComparator - // does to an undatable video. - ["light-and-fresh", "heavy-and-old", "undated"], - ); - assert.deepEqual( - foldSweepPlan({ counts, order: "oldest" }).map((e) => e.channelSlug), - // ...and FIRST under oldest. - ["undated", "heavy-and-old", "light-and-fresh"], - ); - assert.deepEqual( - foldSweepPlan({ counts, order: "listed" }).map((e) => e.channelSlug), - // "listed" consults no date at all: heaviest first, byte-for-byte the - // historical plan. - ["heavy-and-old", "undated", "light-and-fresh"], - ); -}); diff --git a/common/controller/sweepPreview.ts b/common/controller/sweepPreview.ts @@ -1,167 +0,0 @@ -// The sweep plan, built from SNAPSHOTS ONLY — the half of ./sweepPlan.ts that -// has to touch the corpus, kept apart from the half that has to run in a -// browser. -// -// WHY THIS IS NOT buildBackfillSweepPlan. That function is the run's planner and -// it is right for the run: an absent snapshot makes it WALK the channel's -// videos, because planning an unreported channel as zero would exclude it from -// the sweep forever. One walk is ~474,559 file touches and ~4 seconds, which is -// why common/controller/noCorpusWalkInRenderPaths.test.ts bans it from anything -// that renders — and this console re-plans on a 3-second poll. So the preview -// takes the other trade: snapshots only, and a channel with no snapshot is -// reported as UNKNOWN rather than as zero or as a walk. -// -// What it does NOT re-implement is the counting. countFromSnapshot is exported -// from ./backfillSweep and used verbatim here, because a preview that disagrees -// with the run about how much work a channel holds is worse than no preview: it -// is the number an operator arms a multi-day commitment against. - -import type { AutoQueueOrder } from "../jobs/autoQueuePolicy"; -import { presentOperationWork } from "../lib/operations"; -import { - foldSweepPlan, - type SweepChannelCounts, - type SweepKindCounts, - type SweepPlanEntry, -} from "../lib/sweepPlan"; -import { countFromSnapshot } from "./backfillSweep"; -import { type ChannelSnapshot } from "./channelSnapshot"; - -// Exactly what listChannelBriefs already returns, narrowed to the two fields -// this needs. Structural, so the editor's per-request brief cache feeds it with -// no mapping. -export type SweepPreviewBrief = { - slug: string; - snapshot: ChannelSnapshot | null; -}; - -export type BuildSweepChannelCountsInput = { - briefs: ReadonlyArray<SweepPreviewBrief>; - // Every operation the console can offer — the CATALOG, not the scope. Counts - // are carried for all of them so deselecting one re-plans in the browser with - // no round trip; which of them are actually in scope is the fold's business. - kindIds: ReadonlyArray<string>; - // videoId -> YYYYMMDD, for the reachable ids. Absent (or an id absent from - // it) leaves the date "", which controller/planOrder sorts last under - // "newest" and first under "oldest" — the same thing it does for a corpus - // with no dates at all, which is a documented fall-through to weight order - // rather than an arbitrary shuffle. See ./sweepRecency.ts for who fills it. - dates?: ReadonlyMap<string, string>; -}; - -function emptyKindCounts(): SweepKindCounts { - return { - reachable: 0, - missingInput: 0, - blocked: 0, - deferred: 0, - // 0, not null, for an operation with no entry at all: the channel is not - // withholding a number, there is simply no population. A snapshot that - // HAS an entry but predates `eligible` is the null case, below. - eligible: 0, - present: 0, - newestPending: "", - oldestPending: "", - }; -} - -// The reachable ids for one operation on one channel. -// -// snapshot.backfill[id].ids IS the reachable set and nothing else — never -// missing-input, deferred or blocked (a test in lib/operations keeps it equal -// to reachableOperationWork). So dating these ids dates exactly the work the -// sweep would do, which is the only population the order is about. -function reachableIdsFor( - snapshot: ChannelSnapshot, - id: string, -): ReadonlyArray<string> { - const entry = snapshot.backfill?.[id]; - if (entry) return entry.ids ?? []; - // No entry: no work KNOWN for this operation on this channel. True of digest - // exactly as it is of diarization — there is one work list per operation and - // it is the registry's. - return []; -} - -type WorkCounts = Omit<SweepKindCounts, "newestPending" | "oldestPending">; - -function countsFor(snapshot: ChannelSnapshot, id: string): WorkCounts | null { - const entry = snapshot.backfill?.[id]; - if (entry) { - // The RUN'S OWN arithmetic for the two figures the plan is about, on one - // kind at a time. - const { reachable, missingInput } = countFromSnapshot(snapshot.backfill!, [ - id, - ]); - return { - reachable, - missingInput, - // `?? 0` at every read: every snapshot on disk predates one or other of - // these fields, and undefined poisons a sum to NaN. - blocked: entry.blocked ?? 0, - deferred: entry.deferred ?? 0, - eligible: entry.eligible ?? null, - present: presentOperationWork(entry), - }; - } - // Null, never zeroes: a channel whose snapshot carries no entry for this - // operation cannot say how much work it holds, and the caller already handles - // that answer for every other kind. - return null; -} - -export function buildSweepChannelCounts({ - briefs, - kindIds, - dates, -}: BuildSweepChannelCountsInput): SweepChannelCounts[] { - return briefs.map((brief) => { - const snapshot = brief.snapshot; - const byKind: Record<string, SweepKindCounts> = {}; - for (const id of kindIds) { - const counts = emptyKindCounts(); - byKind[id] = counts; - if (!snapshot) continue; - const work = countsFor(snapshot, id); - if (work) Object.assign(counts, work); - if (!dates) continue; - for (const videoId of reachableIdsFor(snapshot, id)) { - const key = dates.get(videoId); - if (!key) continue; - if (key > counts.newestPending) counts.newestPending = key; - if (!counts.oldestPending || key < counts.oldestPending) { - counts.oldestPending = key; - } - } - } - return { - channelSlug: brief.slug, - // No snapshot at all — see the field's comment in lib/sweepPlan.ts. NOT - // the same as a snapshot that reports zero work, which is a real finding. - unknown: snapshot === null, - byKind, - }; - }); -} - -// Counts + fold in one call, for a server that wants the plan directly (the -// offline sanity script, and the tests that compare this against the run's own -// planner). -export function buildSweepPreview({ - briefs, - kindIds, - scopeKindIds, - order, - dates, -}: BuildSweepChannelCountsInput & { - // The scope, when it differs from the catalog. Absent = every catalog kind, - // which is what an unscoped sweep runs. - scopeKindIds?: ReadonlyArray<string>; - order?: AutoQueueOrder; -}): SweepPlanEntry[] { - return foldSweepPlan({ - counts: buildSweepChannelCounts({ briefs, kindIds, dates }), - kindIds: scopeKindIds, - order, - }); -} diff --git a/common/controller/sweepRecency.ts b/common/controller/sweepRecency.ts Binary files differ. diff --git a/common/jobs/autoQueuePolicy.ts b/common/jobs/autoQueuePolicy.ts @@ -73,10 +73,10 @@ export function sanitizeAutoQueueReach(value: unknown): AutoQueueReach { // Coerce a stored/raw value to a legal order. Anything unrecognised — including // a missing field on a settings file written before the field existed — means -// "listed", i.e. today's behaviour. One sanitizer, because the same enum is now -// stored in three places (the two auto-queue policies, digest.recencyOrder and -// backfill.order) and three copies of this line would eventually disagree about -// what an absent field means. +// "listed", i.e. today's behaviour. One sanitizer, because the same enum is +// stored on all four lane policies — it was stored in three MORE places before +// slice 1.3 folded digest.recencyOrder and backfill.order into them — and copies +// of this line would eventually disagree about what an absent field means. export function sanitizeAutoQueueOrder(value: unknown): AutoQueueOrder { return value === "newest" || value === "oldest" || value === "cheapest" ? value @@ -222,9 +222,10 @@ export type ChannelWork = { // its deletion is the point of the four-lane model rather than a tidy-up. // // It existed because one tree could hold leaves belonging to two dispatchers — -// a bucket leaf the runner owned and an operation leaf the arbiter owned — and -// each had to zero the other's before selecting, or auto-transcribe would hand -// a digest candidate to whisper. That was a filter applied AFTER the draw. +// a bucket leaf the runner owned and an operation leaf the (now retired) +// arbiter owned — and each had to zero the other's before selecting, or +// auto-transcribe would hand a digest candidate to whisper. That was a filter +// applied AFTER the draw. // // Now the LANE is the dispatcher, and the draw itself is lane-scoped: a lane // projects only the buckets bucketsForKind gives it and only the operations diff --git a/common/jobs/jobKinds.ts b/common/jobs/jobKinds.ts @@ -240,29 +240,6 @@ const JOB_KINDS: Record<string, JobKindMeta> = { replayable: true, queueKeyStrategy: "custom", }, - // The corpus-wide backfill sweep's orchestrator. Drainable, because a "stop" - // has to reach the channel job it is waiting on rather than meaning "after the - // 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", - drainable: true, - replayable: false, - queueKeyStrategy: "parallel", - }, // One channel through the backfill lane. Drainable (the batch stops taking new // videos and lets the in-flight one finish) and replayable, because it // re-derives its work-list from disk on every run and so replays correctly diff --git a/common/jobs/snapshotScheduler.ts b/common/jobs/snapshotScheduler.ts @@ -44,9 +44,8 @@ const NO_REGEN_KINDS = new Set<string>([ // sweep is tens of thousands of sub-operations, ctx.recordTaskDone arms a // regen after every one, and each regen is a full 16-way per-video walk of the // channel. Leaving them in would spend more machine time regenerating reports - // than doing the backfill. backfillSweep's per-channel job regenerates ONCE, - // at job end. - "backfill-sweep", + // than doing the backfill. The per-channel backfill job regenerates ONCE, at + // job end. "backfill-channel", ]); diff --git a/common/lib/autoQueueTypes.ts b/common/lib/autoQueueTypes.ts @@ -3,8 +3,8 @@ // // These types live in lib/ (the model layer) rather than beside the engine in // jobs/autoQueuePolicy.ts because lib/ may not import jobs/ — settings.ts, -// pauseGates.ts, operations.ts and sweepPlan.ts all need to NAME an auto-queue -// policy, and only the runner needs to run one. The engine re-exports every +// pauseGates.ts, operations.ts and laneMigration.ts all need to NAME an +// auto-queue policy, and only the runner needs to run one. The engine re-exports every // name below, so no import site had to change. // // Enforced by ../architecture.test.ts. diff --git a/common/lib/operations.ts b/common/lib/operations.ts @@ -17,12 +17,11 @@ // and the indicator. // // WHAT "BACKFILL" STILL MEANS IN THIS FILE. The backfill LANE — BACKFILL_QUEUE, -// backfillLaneOperations, backfillLaneEntriesOf, controller/operationBatch.ts and -// controller/backfillSweep.ts — is ONE QUEUE that several operations share: one -// pause, one sweep, one share. It is the only thing this file still calls -// "backfill", and its persisted contracts (the `backfill` key on the snapshot, -// `settings.backfill`, the `backfill-channel` and `backfill-sweep` job kinds) -// keep the word because they are on disk. An operation is not a backfill; a +// backfillLaneOperations, backfillLaneEntriesOf and controller/operationBatch.ts +// — is ONE QUEUE that several operations share: one pause, one runner, one +// share. It is the only thing this file still calls "backfill", and its +// persisted contracts (the `backfill` key on the snapshot, `settings.backfill`, +// the `backfill-channel` job kind) keep the word because they are on disk. An operation is not a backfill; a // backfill is what the lane does to the corpus that predates an operation. // Anything named `*Kind*` outside this file's exports means "one entry of this // registry", which stays true. @@ -390,7 +389,7 @@ export type Operation = { // per video", "~1 model call per transcript chunk". // // "UNIT" IS NOT THE WORD FOR THIS, and the rule is one-sided. A UNIT is one - // item of dispatchable work: what ArbiterUnit, startWorkerUnit, + // item of dispatchable work: what runOperationUnit, startWorkerUnit, // runUnitViaRemote and /api/worker/unit move around, and three of those are // persisted contracts (the route, the `worker-unit` and `auto-download-unit` // job kinds). What one video of an operation costs is its COST BASIS, and no @@ -1057,8 +1056,9 @@ function toBackfillOutcome(outcome: AttributeOneOutcome): OperationRunOutcome { // implementation of this same idea: their own sweep, their own per-channel // batch, their own yield probe, their own cost planner, their own queue keys, // their own settings block with its own pause, and their own snapshot counter. -// backfillSweep.ts is, in its own words, a clone of digestSweep.ts. Registering -// the operation is how that convergence starts. +// The backfill sweep was, in its own words, a clone of the digest one — both +// retired in slice 1.3. Registering the operation is how that convergence +// started. // // NOTHING ABOUT DIGEST GENERATION IS REWRITTEN HERE. Every function this entry // calls is the one the digest controller already calls — resolveDigestTarget for @@ -1104,9 +1104,9 @@ const digest: Operation = { lane: digestLaneFor("local-gpu"), // The live answer, for the same reason diarization has one: which lane this // runs on is a CONFIGURATION CHOICE, not a fact about the operation, and - // `lane` above can only carry the default. laneForOperation asks this, so the - // arbiter dispatches a remote digest onto DIGEST_REMOTE_QUEUE instead of the - // local key it declares. + // `lane` above can only carry the default. laneForOperation asks this, so a + // remote digest is dispatched onto DIGEST_REMOTE_QUEUE instead of the local + // key it declares. // // Declaring it does NOT put digest into the backfill lane's idle-only rule. // That used to need saying twice, because the rule keyed off laneFor's @@ -1410,8 +1410,8 @@ export type OperationDescriptor = { // and an external operation with none is a legal entry (transcode was one, // 2026-08-26 → 08-30). A console that guessed from the id would hand such an // entry the transcription runner's controls — a live Start button over the - // wrong lane. Absent for every backfill kind: those are dispatched by the - // sweep and the arbiter. + // wrong lane. Absent for every registry operation: since slice 1.2 the console + // chooses a lane's console off `pauseLaneFor`, which reads the queue key. // // Typed to the auto-queue kinds on purpose; the sync heartbeat is a runner in // the IA doc's sense but not one of these, which is why /operations/sync is diff --git a/common/lib/pauseGates.test.ts b/common/lib/pauseGates.test.ts @@ -55,26 +55,25 @@ test("backfill's polarity is pinned: enabled false IS held", () => { assert.equal(isGateHeld(running, "backfill"), false); }); -test("holding the backfill lane leaves the sweep's scope alone", () => { - // The bug this forbids: a rebuilt literal instead of a spread would disarm a - // multi-week sweep, or drop its scope, on a pause click. +test("holding the backfill lane leaves the rest of its block alone", () => { + // The bug this forbids: a rebuilt literal instead of a spread would drop the + // lane's other settings on a pause click. It used to guard the sweep's + // persisted scope; the scope is the lane's TREE now, in `autoQueue.backfill`, + // which the same rule protects — `withGateHeld` must touch one field. const base: SiteSettings = { ...defaultSiteSettings(), backfill: { ...defaultSiteSettings().backfill, enabled: true, - sweepEnabled: true, - sweepKinds: ["diarization"], - sweepChannels: ["a-channel"], + concurrency: 3, allowRedownload: true, }, }; const held = withGateHeld(base, "backfill", true); assert.equal(held.backfill.enabled, false); - assert.equal(held.backfill.sweepEnabled, true); - assert.deepEqual(held.backfill.sweepKinds, ["diarization"]); - assert.deepEqual(held.backfill.sweepChannels, ["a-channel"]); + assert.equal(held.backfill.concurrency, 3); assert.equal(held.backfill.allowRedownload, true); + assert.deepEqual(held.autoQueue, base.autoQueue); }); test("isGateHeld touches only its own lane's field", () => { diff --git a/common/lib/pauseGates.ts b/common/lib/pauseGates.ts @@ -93,34 +93,6 @@ export function isGateHeld(settings: SiteSettings, lane: PauseLane): boolean { } } -// WHY THIS LANE'S RUNNER WILL NOT START, or null when it will. -// -// Stated as a function so a 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. -// -// One rule, and a TEMPORARY one: while a lane still has a sweep, arming both is -// two dispatchers on one lane and they would start the same channel twice. It -// was `arbiter.ts`'s `arbiterBlockedReason` (which asked about BOTH sweeps at -// once, because the arbiter was one dispatcher for both lanes); it is per lane -// now, because a runner is. It goes when the sweeps do, in slice 1.3, and this -// whole function goes with them. -// -// A HOLD, not a stop, at the runner: the loop idles with this as its reason -// rather than ending, so disarming the sweep resumes dispatch within one poll. -export function laneBlockedReason( - settings: SiteSettings, - lane: PauseLane, -): string | null { - if (lane === "digest" && 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 (lane === "backfill" && 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; -} - // Set a lane's gate, returning a NEW settings object. Pure — no I/O; the caller // writes it. // @@ -130,8 +102,8 @@ export function laneBlockedReason( // drift. // // SPREAD-AND-OVERRIDE, never a rebuilt literal: `backfill` also carries -// sweepEnabled, sweepKinds, sweepChannels, weight and allowRedownload, and a -// literal here would disarm a multi-week sweep on a pause click. +// `concurrency` and `allowRedownload`, and a literal here would silently reset +// a disk-holding opt-in on a pause click. export function withGateHeld( settings: SiteSettings, lane: PauseLane, diff --git a/common/lib/settings.ts b/common/lib/settings.ts @@ -19,11 +19,8 @@ import { sanitizeWorkers, validateWorkers, } from "./workers"; -import type { - AutoQueueOrder, - AutoQueueReach, - AutoQueueSettings, -} from "./autoQueueTypes"; +import type { AutoQueueSettings } from "./autoQueueTypes"; +import { migrateSweepsToLanes } from "./laneMigration"; // The four SANITIZERS still come from the engine. They are the auto-queue's // half of the settings schema and belong in lib/ with the rest of it, but that // move is phase 3 slice 4 (one schema, one writer) — not a rename. Recorded in @@ -31,8 +28,6 @@ import type { import { defaultAutoQueue, sanitizeAutoQueue, - sanitizeAutoQueueOrder, - sanitizeAutoQueueReach, } from "../jobs/autoQueuePolicy"; import { DEFAULT_DIARIZATION_ENGINE, @@ -285,49 +280,23 @@ export type AttributionSettings = { // Configuration for the backfill lane — the generic answer to "a derived-data // feature landed and 77,000 existing videos do not have it". // -// The one knob that matters is `weight`, and it is a SHARE, not a priority: the -// registry submits every named queue at concurrency 1 and SchedulerTier only -// orders work within a single key, so there is no priority system to join. What -// the lane actually gets is its own queueKey (concurrency with transcription) -// plus a limit() that returns 0 to stand aside — the same mechanism the digest -// yield uses, which fails OPEN so a bad read costs contention rather than a -// deadlock. +// WHAT THE LANE GETS is its own queueKey (concurrency with transcription) plus a +// limit() that returns 0 to stand aside — the same mechanism the digest yield +// uses, which fails OPEN so a bad read costs contention rather than a deadlock. +// There is no priority system to join: the registry submits every named queue at +// concurrency 1 and SchedulerTier only orders work within a single key. +// +// THE SHARE IS `concurrency` AND THE LANE'S `autoQueue.backfill.maxWorkers`; the +// yield is the operation's declared `contendsFor`. Slice 1.3 retired the +// `weight` scalar that used to mean both — see backfillLimit(). export type BackfillSettings = { // Master switch for the lane. Off means the registry still REPORTS what is // missing (that is the indicator's whole job) but nothing runs. enabled: boolean; - // The resource share, 0..1. - // - // 0 (default) — idle-only: run only while the primary transcription lane is - // doing nothing. A backfill is by definition catch-up work on a corpus - // that already exists, so it must never slow down new arrivals. - // >0 — a guaranteed share of the lane's slots, floored at 1 so a small - // weight is a slow lane rather than a stopped one. - // - // Clamped to [0, 1]. See backfillLimit() in controller/operationBatch.ts. - weight: number; // Slots the lane may use when it is not standing aside. Kept at 1 by default // for the same reason diarization.concurrency is: this is CPU-bound work // competing with GPU feeding and the digest sweep for the same 8 threads. concurrency: number; - // A corpus-wide sweep is armed. Persisted so a restart resumes it, exactly as - // digest.sweepEnabled is (editor/instrumentation.ts). - sweepEnabled: boolean; - // The sweep's SCOPE, persisted alongside the flag rather than only in the - // launching call. The boot hook re-launches from settings alone, so arming - // without recording the scope resurrects a deliberately-bounded run as a - // corpus-wide one — the bug digestSweep.ts documents at its start/stop pair. - // Empty = every registered lane-tier kind / every channel. - sweepKinds: string[]; - sweepChannels: string[]; - // The order the backfill batch hands out a channel's videos. Same three - // values and the same "listed means don't sort" contract as the auto-queue - // policies' `order`, so an operator learns the control once. Default - // "listed", i.e. today's behaviour. - order: AutoQueueOrder; - // How far `order` reaches (see AutoQueueReach). Same two values and the same - // meaning as digest.recencyReach. - reach: AutoQueueReach; // Re-acquire media for videos whose input is GONE (audio deleted after // transcription). OFF by default and deliberately so: measured on this corpus, // 836 videos still have media and ~76,270 would need a re-download — 91x the @@ -498,19 +467,6 @@ export type DigestSettings = { // Composes with `yieldToTranscription`: that is the master switch, this only // narrows which workers it reacts to. yieldToCpuWorkers: boolean; - // The corpus-wide sweep is armed. Read at boot by the editor's instrumentation - // hook, the same way the auto-transcribe/auto-download runners are, so a sweep - // survives a server restart. It is persisted INTENT, not a cursor: the batch - // re-derives eligibility from disk on every pull, so a resumed sweep does zero - // rework and needs nothing else remembered. - sweepEnabled: boolean; - // Channels the armed sweep covers. Empty = the whole corpus. - // - // Persisted alongside `sweepEnabled` because the boot hook re-launches from - // settings alone: without it, a sweep deliberately scoped to two channels - // would come back after a restart as an unscoped corpus-wide run — silently - // widening GPU-weeks of work that an operator had bounded on purpose. - sweepChannels: string[]; // Hard ceiling on cumulative metered spend per job, USD. 0 = no cap. Only ever // consulted for a metered app. spendCapUsd: number; @@ -531,20 +487,6 @@ export type DigestSettings = { // Was a scored variable in the bake-off rather than a pre-applied fix; the // measurement is in and "chunk-local" is now the shipped default. timestampMode: DigestTimestampMode; - // The order the digest batch hands out a channel's videos: newest upload - // first, oldest first, or the listed order (today's behaviour, and the - // default). Deliberately a SEPARATE field from DigestBatchOptions.order — - // shortest-first/longest-first is a duration axis and this is a date axis, - // and they compose rather than exclude: the batch sorts by duration and then - // stable-sorts by date, so "newest first" reads as newest day first, shortest - // video within a day. One enum carrying both would make them look mutually - // exclusive, which they are not. - recencyOrder: AutoQueueOrder; - // How far `recencyOrder` reaches (see AutoQueueReach). "channel" (the - // default) orders each channel's own candidates and leaves the sweep visiting - // channels heaviest-first; "corpus" additionally orders the CHANNELS by their - // freshest pending video, so the channel holding the newest work goes first. - recencyReach: AutoQueueReach; // A free-text label for a non-default prompt shape, folded into the recorded // provenance by digestPromptVariant(). Setting it invalidates every digest // generated under a different label, which is exactly what makes a bake-off @@ -901,16 +843,9 @@ export function defaultDigest(): DigestSettings { // OFF. A CPU-pinned worker is not GPU contention, and treating it as such // stalled the digest lane for nothing. See DigestSettings.yieldToCpuWorkers. yieldToCpuWorkers: false, - // OFF. A corpus-wide sweep is GPU-weeks of work and is never armed by - // default — an operator starts it. - sweepEnabled: false, - sweepChannels: [], spendCapUsd: 0, sections: ["chapters"], timestampMode: DEFAULT_DIGEST_TIMESTAMP_MODE, - // "listed" — today's behaviour exactly. Ordering is opt-in. - recencyOrder: "listed", - recencyReach: "channel", promptVariant: "", }; } @@ -986,12 +921,6 @@ export function sanitizeDigest(value: unknown): DigestSettings { // transcription. `yieldToTranscription` still gates the whole thing, so the // GPU-safe default is untouched. yieldToCpuWorkers: r.yieldToCpuWorkers === true, - sweepEnabled: r.sweepEnabled === true, - sweepChannels: Array.isArray(r.sweepChannels) - ? r.sweepChannels.filter( - (v): v is string => typeof v === "string" && v.trim().length > 0, - ) - : [], spendCapUsd: typeof r.spendCapUsd === "number" && r.spendCapUsd > 0 ? Math.round(r.spendCapUsd * 100) / 100 @@ -1002,8 +931,6 @@ export function sanitizeDigest(value: unknown): DigestSettings { timestampMode: isDigestTimestampMode(r.timestampMode) ? r.timestampMode : d.timestampMode, - recencyOrder: sanitizeAutoQueueOrder(r.recencyOrder), - recencyReach: sanitizeAutoQueueReach(r.recencyReach), // Trimmed and length-capped: it goes into provenance on every record, and a // runaway value would bloat 119k sidecars. promptVariant: @@ -1109,15 +1036,7 @@ export function defaultSiteSettings(): SiteSettings { export function defaultBackfill(): BackfillSettings { return { enabled: false, - // Idle-only. See BackfillSettings.weight. - weight: 0, concurrency: 1, - sweepEnabled: false, - sweepKinds: [], - sweepChannels: [], - // "listed" — today's behaviour exactly. Ordering is opt-in. - order: "listed", - reach: "channel", // See BackfillSettings.allowRedownload — this one holds disk. allowRedownload: false, }; @@ -1127,24 +1046,11 @@ export function sanitizeBackfill(value: unknown): BackfillSettings { const d = defaultBackfill(); if (!value || typeof value !== "object") return d; const r = value as Record<string, unknown>; - const slugs = (v: unknown): string[] => - Array.isArray(v) - ? v.filter((s): s is string => typeof s === "string" && s.trim() !== "") - : []; return { enabled: r.enabled === true, // Clamped rather than rejected: a hand-edited 5 means "as much as possible", - // and reading it as the idle-only 0 would be the opposite of the intent. - weight: - typeof r.weight === "number" && Number.isFinite(r.weight) - ? Math.min(1, Math.max(0, r.weight)) - : d.weight, + // and reading it as 0 would be the opposite of the intent. concurrency: clampPositiveInt(r.concurrency, d.concurrency, 16), - sweepEnabled: r.sweepEnabled === true, - sweepKinds: slugs(r.sweepKinds), - sweepChannels: slugs(r.sweepChannels), - order: sanitizeAutoQueueOrder(r.order), - reach: sanitizeAutoQueueReach(r.reach), allowRedownload: r.allowRedownload === true, }; } @@ -1485,7 +1391,13 @@ export function getSettings(): SiteSettings { merged.autoRefreshIntervalSeconds, ); merged.syncScheduler = sanitizeSyncScheduler(merged.syncScheduler); - merged.autoQueue = sanitizeAutoQueue(merged.autoQueue); + // THE SWEEPS' SCOPE, ON READ. `migrateSweepsToLanes` fills in + // `autoQueue.digest` / `.backfill` from the retired sweep fields when — and + // only when — the FILE does not already spell them, which is why it is handed + // `parsed` rather than `merged`: `merged` has had defaults folded in and can + // no longer tell "absent" from "default". It never enables a lane the sweep + // flag did not. See lib/laneMigration.ts. + merged.autoQueue = sanitizeAutoQueue(migrateSweepsToLanes(parsed)); merged.socialLinks = parseSocialLinks(merged.socialLinks); merged.homepageUrl = normalizeHomepageUrl(merged.homepageUrl); merged.savedVideoBackup = sanitizeSavedVideoBackup(merged.savedVideoBackup); diff --git a/common/lib/sweepPlan.ts b/common/lib/sweepPlan.ts @@ -1,163 +0,0 @@ -// THE PLAN: the ordered channel itinerary a sweep will actually walk. -// -// A sweep is a commitment, not a toggle — on this corpus arming one can mean -// ~194,000 model calls — and until now the only way to read what was being -// committed to was to arm it and watch the log. This module is the shape that -// makes the commitment legible BEFORE the click: the same counting the run -// does, folded over whichever operations are in scope, ordered by the same -// comparator, so the itinerary on screen is the itinerary that runs. -// -// PURE, AND IN lib/ RATHER THAN controller/ ON PURPOSE. The console re-plans on -// every checkbox — deselect an operation and the totals fall, the rows re-sort -// and the button re-labels itself — and a round trip per keystroke would make -// that feel like a form rather than an instrument. So the fold has to run in -// the browser, which means it may not import anything that touches disk: -// controller/channelSnapshot reaches runYtdlp reaches execa, and a client -// component importing that fails `next build` on `node:child_process`. This is -// the same split components/pipelines/band.ts already documents, for the same -// reason and in the same direction. -// -// COUNTS ARE CARRIED PER OPERATION AND NEVER PRE-SUMMED. The server hands over -// one entry per kind per channel; the client sums only what is selected. That -// is not merely a convenience: the reason the scope control exists at all is -// that "11,337 reachable" is the same shape of number whether one unit is a -// single audio pass or ~1 model call per transcript CHUNK, so a payload that -// arrived pre-summed would have already thrown away the distinction the console -// is being built to show. - -import { orderPlanByRecency } from "../controller/planOrder"; -import type { AutoQueueOrder } from "./autoQueueTypes"; - -// One operation's outstanding work on one channel. -// -// `reachable` and `missingInput` stay apart here as they do everywhere else — -// see lib/operations.ts's header. On the measured corpus they are four orders -// of magnitude apart on diarization, and one "remaining" number would say the -// same thing about a lane that is finished and a lane that cannot start. -export type SweepKindCounts = { - reachable: number; - missingInput: number; - // The band's other two work populations, carried so a plan row can draw the - // same five-fill instrument /channels and the comparison rail draw — and - // NEVER added to `reachable`. Waiting on an upstream operation and held by a - // gate are not work this run can do. - blocked: number; - deferred: number; - // The band's coverage halves. `null` is "this snapshot cannot say", and it - // poisons a sum deliberately: a partial denominator is smaller than its own - // numerator, which is a worse lie than admitting the number is not knowable. - eligible: number | null; - present: number | null; - // Newest / oldest upload date (YYYYMMDD) among THIS operation's reachable - // ids, or "" when none of them could be dated. "" sorts last under "newest" - // and first under "oldest", exactly as controller/planOrder documents for a - // channel and makeRecencyComparator does for a single video. - // - // Per KIND rather than per channel, because the fold below is a max over the - // selected kinds and a max of maxima is the max of the union — so toggling an - // operation off re-dates the row correctly with no second lookup. - newestPending: string; - oldestPending: string; -}; - -export type SweepChannelCounts = { - channelSlug: string; - // This channel has never been reported on, so every count below is 0 BECAUSE - // NOTHING HAS LOOKED — not because there is nothing to do. The run's own - // planner walks such a channel's videos rather than planning it as zero - // (see buildBackfillSweepPlan), which a 3-second poll cannot do; so the - // console says "not reported yet" where it would otherwise print a confident - // zero. Unknown is not zero, here as everywhere else on these surfaces. - unknown: boolean; - // Keyed by operation id. An operation with no entry contributes nothing, - // which is also what an UNKNOWN id does — sanitizeBackfill keeps an id that - // no longer names a registered kind, and a stored scope naming one must - // count as zero rather than as everything. - byKind: Record<string, SweepKindCounts>; -}; - -// One channel's row in the plan. The shape controller/backfillSweep's own -// planner returns, so preview and run cannot describe a channel differently. -export type SweepPlanEntry = { - channelSlug: string; - reachable: number; - missingInput: number; - newestPending: string; - oldestPending: string; -}; - -// Fold the per-kind counts down to one row per channel, ordered. -// -// `kindIds` empty means EVERY kind on the payload — the same rule -// resolveBackfillLaneOperations applies to an absent scope, and the same rule -// startBackfillSweep persists. It is deliberately not "no kinds": an unscoped -// sweep is the corpus-wide one, and rendering it as an empty plan would say the -// opposite of what arming it would do. -export function foldSweepPlan({ - counts, - kindIds, - order = "listed", -}: { - counts: ReadonlyArray<SweepChannelCounts>; - kindIds?: ReadonlyArray<string>; - order?: AutoQueueOrder; -}): SweepPlanEntry[] { - const scope = kindIds && kindIds.length > 0 ? new Set(kindIds) : null; - const plan: SweepPlanEntry[] = counts.map((channel) => { - const entry: SweepPlanEntry = { - channelSlug: channel.channelSlug, - reachable: 0, - missingInput: 0, - newestPending: "", - oldestPending: "", - }; - for (const [id, kind] of Object.entries(channel.byKind)) { - if (scope && !scope.has(id)) continue; - entry.reachable += kind.reachable; - entry.missingInput += kind.missingInput; - if (kind.newestPending > entry.newestPending) { - entry.newestPending = kind.newestPending; - } - if ( - kind.oldestPending && - (!entry.oldestPending || kind.oldestPending < entry.oldestPending) - ) { - entry.oldestPending = kind.oldestPending; - } - } - return entry; - }); - // The SAME comparator and the SAME weight tiebreak buildBackfillSweepPlan - // passes. Not a copy of the rule: the rule itself. - return orderPlanByRecency( - plan, - order, - (a, b) => - b.reachable - a.reachable || a.channelSlug.localeCompare(b.channelSlug), - ); -} - -export type SweepPlanTotals = { - reachable: number; - missingInput: number; - // Channels the sweep would actually visit — those holding reachable work. - channels: number; - // Channels in the corpus that hold none. Stated separately rather than - // subtracted at the call site so "65 channels with nothing to do" and "3 of 68" - // come from one place and cannot disagree. - idle: number; -}; - -export function sweepPlanTotals( - plan: ReadonlyArray<SweepPlanEntry>, -): SweepPlanTotals { - let reachable = 0; - let missingInput = 0; - let channels = 0; - for (const entry of plan) { - reachable += entry.reachable; - missingInput += entry.missingInput; - if (entry.reachable > 0) channels++; - } - return { reachable, missingInput, channels, idle: plan.length - channels }; -} diff --git a/editor/instrumentation.ts b/editor/instrumentation.ts @@ -81,8 +81,8 @@ export async function register() { // Everything past here STARTS work. On an idle boot, nothing does. if (idle) return; - // Start the automatic priority-queue runners (auto-transcribe / auto-download) - // if their policies are enabled. Each is a self-managed registry job; this only + // Start the automatic priority-queue runners — all four lanes — if their + // policies are enabled. Each is a self-managed registry job; this only // kicks them off and returns. Best-effort — a failure here must not stop the // server from starting. try { @@ -94,37 +94,9 @@ export async function register() { /* a runner that fails to start must not block server readiness */ } - // Resume the corpus-wide digest sweep, if one is armed. A sweep is GPU-WEEKS - // long, so it will outlive several restarts by construction — and before this - // hook a restart silently ended one that had been running for days, with - // nothing to say so. - // - // Re-launching is safe rather than merely convenient: the sweep stores no - // cursor and the batch re-derives eligibility from disk on every pull, so a - // resumed sweep re-does exactly zero work. The digest PAUSE needs no hook - // here for the same reason downloadsPaused doesn't — it is read at dispatch - // time — but "is a sweep running" is process state, and process state is what - // a restart destroys. - try { - const { resumeDigestSweepIfEnabled } = await import( - "yt-dlp-transcript-common/controller/digestSweep" - ); - await resumeDigestSweepIfEnabled(); - } catch { - /* a sweep that fails to resume must not block server readiness */ - } - - // Resume the corpus-wide BACKFILL sweep, if one is armed. Same reasoning as - // the digest sweep above, and safe for the same reason: the sweep stores no - // cursor and the batch re-derives eligibility from disk on every pull, so a - // resumed sweep re-does exactly zero work. What a restart destroys is the - // process state — "is a sweep running" — and this is what restores it. - try { - const { resumeBackfillSweepIfEnabled } = await import( - "yt-dlp-transcript-common/controller/backfillSweep" - ); - await resumeBackfillSweepIfEnabled(); - } catch { - /* a sweep that fails to resume must not block server readiness */ - } + // NO SECOND RESUME HOOK. There used to be two more here, one per sweep, and + // each re-launched a corpus-wide pass from its own persisted flag. Both lanes + // are runner lanes now, so `startAutoRunnersIfEnabled` above resumes all four + // from one switch — `autoQueue[lane].enabled` — which is what a lane being + // armed has always meant on the transcription and download lanes. }