Archilyzer · Source

archilyzer

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

commit 1f7c28b9bc26ffa64201ad7c287f6ceb02790406
parent d6f60ed9e7e5ac26078880c42170940e9549090d
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Mon,  7 Sep 2026 20:20:47 -0400

common: the runner runs operations, and the digest lane's order is one comparator

Slice 1.2, second commit. `autoRunner.ts` branches on the LANE: transcription
and download are byte-for-byte what they were, digest and backfill drive
`operationBatch`. `startAutoRunnersIfEnabled` iterates LANES instead of naming
two, which is the sentence slice 1.1's comment promised would change here.

`buildChannelWork` stops intersecting the projection with what the tree NAMES
and projects what the lane DISPATCHES. It has to: a digest-lane leaf naming no
operation draws the lane's whole union, the way a bucket-less leaf draws the
default bucket union, so the union must be there whether a leaf asks for it or
not. `buildPendingByLeaf` gains `defaultOperations` — the operation half of
`defaultBuckets` — and the 1.1 early-return guard goes with the reason it was
added. The safety rail is intact and cheaper to state: `operationsForLane` is
[] for the two bucket lanes, so a `digest` leaf in the auto-transcribe tree
still finds no list, with no branch anywhere saying so.

THE PICK IS A VIDEO AND THE UNIT IS AN OPERATION, so `run()` decides which. The
lane's operations are asked in dependency order and the first with work runs —
which is also how a video gets diarized on one pick and attributed on the next.
That forced the session's completed-set to be keyed by (operation, video)
rather than by video: retiring the whole video after one unit would break the
chain the lane exists to walk. A video leaves a leaf's list only once every
operation THAT LEAF draws is done for it.

**The `cheapest` order, named.** `digestBatch` sorted by duration and then
sorted AGAIN, stably, by upload date — so the DATE is primary and the duration
is the tiebreak within a day: "newest day first, shortest video within a day".
Two adjacent sort() calls are exactly the composition a migration drops half
of, so it is one comparator now (`cheapestComparator`), with the date half's
direction a PARAMETER — today `settings.digest.recencyOrder` ("newest" live),
fixed at newest-first in 1.3 when that field retires into the tree. Unknown
durations sort last rather than first. On a lane that cannot price its work
`cheapest` falls back to LISTED, not to id order and not to a date order nobody
asked for.

**`laneBlockedReason(lane, settings)`** moves out of `arbiter.ts` into
`lib/pauseGates.ts`, per lane rather than for both at once (a runner is per
lane; the arbiter was not). The runner refuses to start beside an armed sweep
and HOLDS if one is armed under it, naming which. The arbiter keeps working
through it until 1.3.

**The status poll.** Two lanes went from reading nothing to folding 56,223 and
80,193 pending ids over 68 snapshots, on a three-second poll, four lanes at
once. `readChannelSnapshotShared` memoizes the parse on (mtime, size) for that
one read path — deliberately not `readChannelSnapshot` itself, which the
snapshot GENERATOR calls and where a shared object is a much larger claim.
Measured on the live corpus, warm: 256 ms with the memo against 386-413 ms
without, from a 86-102 ms baseline where two lanes read nothing at all. The
cold first pass is 4.8 s, which is the duration index filling; a duration never
changes, so it is paid once per process.

Numbers: none move — the live `settings.json` carries no `autoQueue.digest` or
`.backfill` block, so the script's lane sections are unchanged and the diff is
empty. common 896, tsc clean in six packages.

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

Diffstat:
Mcommon/controller/arbiter.ts | 16+++++++++-------
Mcommon/controller/autoRunner.ts | 536+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------------
Mcommon/controller/channels.ts | 57+++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/autoQueuePolicy.ts | 37++++++++++++++++++++++++++++---------
Mcommon/jobs/jobKinds.test.ts | 6++++++
Mcommon/jobs/jobKinds.ts | 22++++++++++++++++++++++
Mcommon/lib/pauseGates.ts | 28++++++++++++++++++++++++++++
Meditor/app/jobs/active/buildActiveJobs.ts | 6++++++
Meditor/app/operations/components/dispatch.ts | 12+++++++++++-
Aplans/tools/phase1-poll-timing.ts | 46++++++++++++++++++++++++++++++++++++++++++++++
10 files changed, 673 insertions(+), 93 deletions(-)

diff --git a/common/controller/arbiter.ts b/common/controller/arbiter.ts @@ -35,6 +35,7 @@ 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"; @@ -191,15 +192,16 @@ function channelOf( // 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(); - if (settings.digest.sweepEnabled) { - return "The digest sweep is armed. Two dispatchers on one lane would start the same channel twice — stop the sweep first."; - } - if (settings.backfill.sweepEnabled) { - return "The backfill sweep is armed. Two dispatchers on one lane would start the same channel twice — stop the sweep first."; - } - return null; + return ( + laneBlockedReason(settings, "digest") ?? + laneBlockedReason(settings, "backfill") + ); } async function runOneUnit( diff --git a/common/controller/autoRunner.ts b/common/controller/autoRunner.ts @@ -13,9 +13,11 @@ import { runPool } from "../jobs/concurrentRunner"; import { requestChannelSnapshot } from "../jobs/snapshotScheduler"; import { makeTaskTracker } from "../jobs/taskHooks"; import type { JobRunContext } from "../jobs/streamCommand"; +import type { JobProgress } from "../jobs/registry"; import { type ActiveCounts, type AutoQueueGroup, + type AutoQueueLeaf, type AutoQueueOrder, type AutoQueuePolicy, type ChannelWork, @@ -26,7 +28,20 @@ import { defaultBucketsForPolicy, selectableBucketsForKind, } from "../jobs/autoQueuePolicy"; -import { operationsForLane } from "../lib/operations"; +import { operationsForLane, type Operation } from "../lib/operations"; +import { laneBlockedReason } from "../lib/pauseGates"; +import { + cheapestComparator, + classifyOperationUnit, + laneLimit, + openOperationRun, + operationLaneLive, + preflightOperationRun, + runOperationUnit, + videoDurationSeconds, + type OperationLane, + type OperationRun, +} from "./operationBatch"; import { type RecencyKey, buildRecencyKeys, @@ -36,6 +51,7 @@ import { // Re-exported so the editor's status payload can name the type without reaching // past the runner into the index module it is an implementation detail of. export type RecencyKeyView = RecencyKey; +import { LANES } from "../lib/autoQueueTypes"; import { type AutoQueueKind, readAutoQueueState, @@ -58,7 +74,7 @@ import { isGateHeld } from "../lib/pauseGates"; import { listChannelConfigs, readChannelConfig, - readChannelSnapshot, + readChannelSnapshotShared, } from "./channels"; import { transcribeOneFromQueue } from "./transcribeOneFromQueue"; import { findVideoSourceUrl } from "./undownloadedVideos"; @@ -86,9 +102,30 @@ export const AUTO_DOWNLOAD_KIND = "auto-download"; // against a manual sync. Marked background so a clicked sync preempts queued // units. NOT in DRAINABLE_KINDS — a single-video download has nothing to drain. export const AUTO_DOWNLOAD_UNIT_KIND = "auto-download-unit"; +// The two operation lanes' runner jobs. Their UNITS make no job record at all — +// they run in-process inside this loop — which is the difference that lets a +// lane work per video where the manual verbs must work per channel: the +// registry keeps 100 records and the log 500, so 55,956 digest jobs would evict +// the history of the run that made them. +export const AUTO_DIGEST_KIND = "auto-digest"; +export const AUTO_BACKFILL_KIND = "auto-backfill"; + +const RUNNER_JOB_KINDS: Record<AutoQueueKind, string> = { + transcription: AUTO_TRANSCRIBE_KIND, + download: AUTO_DOWNLOAD_KIND, + digest: AUTO_DIGEST_KIND, + backfill: AUTO_BACKFILL_KIND, +}; export function autoRunnerJobKind(kind: AutoQueueKind): string { - return kind === "transcription" ? AUTO_TRANSCRIBE_KIND : AUTO_DOWNLOAD_KIND; + return RUNNER_JOB_KINDS[kind]; +} + +// Which lanes dispatch OPERATIONS rather than buckets — the two whose executor +// is controller/operationBatch.ts. A type guard rather than a set, so a branch +// on it narrows. +export function isOperationLane(kind: AutoQueueKind): kind is OperationLane { + return kind === "digest" || kind === "backfill"; } // How long the loop sleeps when there's no work and nothing in flight, before @@ -134,6 +171,19 @@ 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 + // click on the Resume button beside it. + | "lane-held" + // The lane's engine did not answer its probe (a dead ollama). An idle, not a + // stop: it is a configuration problem that an operator fixes without + // restarting anything, and the probe is re-asked on a TTL. + | "engine-unreachable" // There is genuinely nothing to do. | "no-pending" // Work exists but every path to it is at a maxWorkers ceiling (a node cap, or @@ -211,23 +261,15 @@ async function buildChannelWork( paths: Paths, kind: AutoQueueKind, meta: ReadonlyArray<ChannelMeta>, - // Operation id spaces to project alongside the buckets — exactly the ones the - // policy tree actually names, never the whole catalog. A snapshot's - // backfill[op].ids is up to 11,329 strings on one channel, so projecting an - // operation nothing asks for would put ~78,000 strings a tick behind a - // feature that is off. Every tree written before operations existed names - // none, so this is empty and the projection is byte-identical to today's. + // The operation id spaces to project alongside the buckets: exactly the ones + // THIS LANE DISPATCHES, never the whole catalog. See laneOperationIds. operations: ReadonlyArray<string> = [], ): Promise<{ channels: ChannelWork[]; owner: Map<string, string> }> { - // NOTHING TO PROJECT MEANS NOTHING TO READ. The status payload folds all four - // lanes on a three-second poll, and the digest and backfill lanes have no - // buckets and — while their tree is the default catch-all leaf — no - // operations either, so without this they would read and JSON.parse all 68 - // snapshots (~6.5 MB) twice a tick to build a map every leaf draws [] from. - // - // DELETE THIS the moment a lane's default draw is its operation union (slice - // 1.2): from then on `operations` is non-empty for every live lane and the - // guard only fires for a lane that is genuinely switched off. + // NOTHING TO PROJECT MEANS NOTHING TO READ. A lane with no buckets and no + // operations draws [] from every leaf by construction, so reading 68 + // snapshots to build the map would be pure cost. This no longer fires for the + // digest or backfill lanes — their default draw is their operation union — + // and remains for a lane whose operations are all switched off. if (selectableBucketsForKind(kind).length === 0 && operations.length === 0) { return { channels: [], owner: new Map() }; } @@ -235,9 +277,12 @@ async function buildChannelWork( const owner = new Map<string, string>(); // Read every snapshot concurrently, then fold sequentially — `owner` is // first-writer-wins, so the fold must stay in `meta` order to stay - // deterministic. + // deterministic. Through the SHARED reader: this exact projection runs on the + // runner's scheduling tick and on a three-second status poll folding four + // lanes, and re-parsing ~6.5 MB of unchanged JSON eight times a tick is the + // cost that made the memo worth having. const snaps = await mapConcurrent(meta, SNAPSHOT_READ_CONCURRENCY, (m) => - readChannelSnapshot(paths, m.slug), + readChannelSnapshotShared(paths, m.slug), ); for (const [i, { slug, platform }] of meta.entries()) { const snap = snaps[i]; @@ -286,27 +331,42 @@ export function operationsNamedBy(root: AutoQueueGroup): string[] { return [...out]; } -// WHAT THIS LANE MAY PROJECT: the operations its tree names, INTERSECTED with -// the operations the lane actually dispatches. +// WHAT THIS LANE MAY PROJECT: the operations the lane actually dispatches, and +// nothing else. // -// The intersection is the safety rail retainLeaves used to be. A leaf naming -// `digest` inside the auto-transcribe tree is legal to write and always has -// been; before the lanes existed it was projected and then zeroed by a second -// pass, and forgetting that pass handed a digest candidate to whisper. Now it -// is never projected at all, so the leaf draws nothing and the wrong engine -// cannot be reached from the wrong tree. +// This is the safety rail retainLeaves used to be. A leaf naming `digest` +// inside the auto-transcribe tree is legal to write and always has been; before +// the lanes existed it was projected and then zeroed by a second pass, and +// forgetting that pass handed a digest candidate to whisper. It is never +// projected at all now, so the leaf draws nothing and the wrong engine cannot be +// reached from the wrong tree. `operationsForLane` is [] for the two bucket +// lanes, which is what makes that true with no branch here. // -// It stays an INTERSECTION rather than "everything the lane dispatches" because -// projection costs real memory: one channel's `digest.ids` is up to 11,329 -// strings, and this runs on a three-second status poll. Slice 1.2's `next()` -// widens it to the lane's default draw for the lanes that have an executor. -function laneOperationIds(kind: AutoQueueKind, root: AutoQueueGroup): string[] { - const named = operationsNamedBy(root); - if (named.length === 0) return []; - const dispatched = new Set( - operationsForLane(kind, getSettings()).map((op) => op.id), - ); - return named.filter((id) => dispatched.has(id)); +// SLICE 1.1 INTERSECTED THIS WITH `operationsNamedBy(root)`, to keep the +// projection free while nothing dispatched off it. It cannot any more: a +// digest-lane leaf that names no operation draws the lane's whole union — the +// way a bucket-less leaf draws the default bucket union — so the union has to +// be projected whether a leaf names it or not. What that costs is one channel's +// `digest.ids` (up to 11,329 strings) per channel per tick, which is why +// readChannelSnapshotShared exists. +function laneOperationIds(kind: AutoQueueKind): string[] { + return operationsForLane(kind, getSettings()).map((op) => op.id); +} + +// The operations one leaf draws: the one it NAMES, or the lane's whole union. +// The exact mirror of `defaultBucketsForPolicy` for the operation lanes, and +// the reason a catch-all leaf on the digest lane is not inert. +function leafOperationIds( + leaf: AutoQueueLeaf, + laneOperations: ReadonlyArray<string>, +): string[] { + if (laneOperations.length === 0) return []; + if (!leaf.match.operation) return [...laneOperations]; + // A leaf naming something this lane does not dispatch draws nothing — the + // same answer the projection gives it. + return laneOperations.includes(leaf.match.operation) + ? [leaf.match.operation] + : []; } // The policy's `order`, defaulted for settings files written before the field @@ -335,6 +395,17 @@ async function recencyOrdering( }> { const order = orderOf(policy); if (order === "listed") return { compare: null, keys: new Map() }; + // `cheapest` IS NOT A DATE ORDER, and it is answerable only where durations + // exist. The digest lane's candidates all have a normalized transcript (that + // is what makes them reachable), so every one of them can be priced; a + // download candidate has nothing on disk at all. The sanitizer accepts the + // value on any lane — one enum, one sanitizer — so a lane that cannot price + // its work falls back to LISTED, which is what "no comparator" means + // everywhere else here. Not to id order, and not to a date order nobody asked + // for: those would both be a real reordering dressed up as a fallback. + if (order === "cheapest" && kind !== "digest") { + return { compare: null, keys: new Map() }; + } // Candidates AND their owning channel, derived in one pass from the same // projection buildPendingByLeaf claims from — first-writer-wins, exactly as // buildChannelWork builds its own owner map, so the two cannot disagree. @@ -363,12 +434,59 @@ async function recencyOrdering( // would be pure cost. interpolate: kind === "download", }); - return { compare: makeRecencyComparator(keys, order), keys }; + if (order !== "cheapest") { + return { compare: makeRecencyComparator(keys, order), keys }; + } + // THE DIGEST LANE'S ORDER IS TWO RULES, and this is where they meet. + // + // What the per-channel batch does is sort by duration and then sort AGAIN, + // stably, by upload date — so the date is PRIMARY and the duration is the + // tiebreak within a day: "newest day first, shortest video within a day". + // Written out as one comparator here because the runner has one sort, not + // 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. + const durations = await loadDurations(paths, owner, candidateIds); + return { + compare: cheapestComparator({ + duration: (id) => durations.get(id), + recency: makeRecencyComparator(keys, getSettings().digest.recencyOrder), + }), + keys, + }; } catch { return { compare: null, keys: new Map() }; } } +// Duration per candidate, for the `cheapest` order. Memoized process-wide by +// operationBatch (a duration never changes), so only the first tick of a +// process pays — and it pays it concurrently rather than 55,956 opens in +// series. A video with no normalized transcript reads null and sorts last. +async function loadDurations( + paths: Paths, + owner: ReadonlyMap<string, string>, + candidateIds: ReadonlySet<string>, +): Promise<Map<string, number | null>> { + const ids = [...candidateIds]; + const out = new Map<string, number | null>(); + const seconds = await mapConcurrent(ids, DURATION_READ_CONCURRENCY, (id) => { + const slug = owner.get(id); + return slug + ? videoDurationSeconds(paths, slug, id) + : Promise.resolve(null); + }); + ids.forEach((id, i) => out.set(id, seconds[i])); + return out; +} + +const DURATION_READ_CONCURRENCY = 64; + // How many drill-down ids the status payload carries per rule. Enough to see // whether an ordering change did what you asked; small enough that a 9,000-video // catch-all doesn't ship 9,000 strings on a 3-second poll. @@ -411,11 +529,12 @@ export async function computeLeafPending( ): Promise<LeafPending> { const policy = getSettings().autoQueue[kind]; const meta = await listChannelMeta(paths); + const laneOperations = laneOperationIds(kind); const { channels, owner } = await buildChannelWork( paths, kind, meta, - laneOperationIds(kind, policy.root), + laneOperations, ); const { compare, keys } = await recencyOrdering( kind, @@ -432,7 +551,7 @@ export async function computeLeafPending( policy.root, channels, defaultBucketsForPolicy(kind, policy), - compare ? { compare } : undefined, + { ...(compare ? { compare } : {}), defaultOperations: laneOperations }, ); // A video already in flight is not "next up" — drop the live runner's set // before asking the policy, exactly as next() does. @@ -538,6 +657,7 @@ async function runLoop( paths: Paths, live: RunnerLive, onLog: (line: string) => void, + setProgress: (snap: JobProgress) => void, signal: AbortSignal, ctx: JobRunContext, ): Promise<void> { @@ -548,24 +668,76 @@ async function runLoop( // Drop long-lapsed platform cooldowns on boot; entries still in (or recently // out of) their window are kept so an Odysee 429 cooldown survives a restart. pruneExpired(kindState.platformBackoff, Date.now()); + // The session's "already done, do not re-pick" set, KEYED BY (operation, + // video) rather than by video. + // + // On a bucket lane the operation is "" and this is exactly what it always + // was. On an operation lane it has to be a pair: the same video is + // legitimately pending for diarization AND for attribution-diarized — the + // second one depends on the first's output — so retiring the video after one + // unit would break the dependency chain the lane exists to walk. A video + // leaves the pending list only once EVERY operation its leaf draws is done + // for it (see dropCompleted). const completed = new Set<string>(); const completedOrder: string[] = []; const counters = { transcribed: 0, failed: 0, skipped: 0 }; - const markCompleted = (videoId: string) => { - if (completed.has(videoId)) return; - completed.add(videoId); - completedOrder.push(videoId); + const completedKey = (operation: string, videoId: string) => + `${operation}\u0000${videoId}`; + + const markCompleted = (videoId: string, operation = "") => { + const key = completedKey(operation, videoId); + if (completed.has(key)) return; + completed.add(key); + completedOrder.push(key); if (completedOrder.length > COMPLETED_CAP) { const old = completedOrder.shift(); if (old) completed.delete(old); } }; + // Drop what this session has already finished, PER LEAF — because which + // operations a leaf draws is a per-leaf fact (the one it names, or the lane's + // whole union). A bucket lane's leaves all draw [] and the rule collapses to + // "this video is done", which is what it has always been. + const dropCompleted = ( + pending: Record<string, string[]>, + root: AutoQueueGroup, + laneOperations: ReadonlyArray<string>, + ): void => { + for (const leaf of flattenLeaves(root)) { + const ops = leafOperationIds(leaf, laneOperations); + const list = pending[leaf.id]; + if (!list) continue; + pending[leaf.id] = + ops.length === 0 + ? list.filter((id) => !completed.has(completedKey("", id))) + : list.filter( + (id) => !ops.every((op) => completed.has(completedKey(op, id))), + ); + } + }; + const persist = () => { void writeAutoQueueState(paths, state).catch(() => {}); }; + // Units this loop has finished, and what was still pending when it last + // looked. Only the operation lanes draw a bar: the two bucket lanes have had + // none since they were written, and adding one to them here would be a + // visible change to a surface this slice is not touching. + let lastPending = 0; + const reportRunnerProgress = (pending = lastPending): void => { + lastPending = pending; + const done = counters.transcribed + counters.failed + counters.skipped; + setProgress({ + metric: kind === "digest" ? "digests" : "backfills", + initial: 0, + current: done, + target: done + pending, + }); + }; + // Downloads have no worker pool; concurrency politeness is per-platform (one // download per source at a time, matching the per-platform serial job queue). // Track in-flight downloads per platform so the loop never launches a second @@ -618,14 +790,109 @@ async function runLoop( if (signal.aborted) onHardAbort(); else signal.addEventListener("abort", onHardAbort, { once: true }); + // ── The operation lanes' run context ──────────────────────────────────── + // + // `openOperationRun` resolves what a unit needs ONCE — the engine, the + // duplicate-cluster plan, the fan-out envelope — because resolving it per + // video is how a counter and a runner end up disagreeing about what is + // stale. A per-channel batch opens one per job; a runner is long-lived, so it + // re-opens when the operation settings change or the TTL lapses, which is + // also what re-asks the engine probe. + // + // REFRESHED FROM limit(), not only from next(), and that is load-bearing: + // when the probe fails the limit is 0, runPool idle-waits and NEVER CALLS + // next() — so a next()-only refresh would wedge on the first bad probe and + // never re-probe. limit() is synchronous, so it starts the refresh and uses + // whatever is cached; the next poll picks up the answer. + const laneRun: { + run: OperationRun | null; + key: string; + at: number; + error: string | null; + refreshing: boolean; + } = { run: null, key: "", at: 0, error: null, refreshing: false }; + const LANE_RUN_TTL_MS = 60_000; + + const laneRunKey = (settings: ReturnType<typeof getSettings>): string => + JSON.stringify([ + settings.digest, + settings.backfill, + settings.diarization, + settings.attribution, + ]); + + const refreshLaneRun = async (force = false): Promise<void> => { + if (!isOperationLane(kind) || laneRun.refreshing) return; + const settings = getSettings(); + const key = laneRunKey(settings); + if ( + !force && + laneRun.run && + laneRun.key === key && + Date.now() - laneRun.at < LANE_RUN_TTL_MS + ) { + return; + } + laneRun.refreshing = true; + try { + const opened = await openOperationRun({ + lane: kind, + paths, + onLog, + }); + // The engine fail-fast, asked ONCE per context rather than per video: an + // unreachable ollama would otherwise produce one failure per candidate + // over 55,956 of them. A refusal is an IDLE, not a stop — it is a + // configuration problem an operator fixes without restarting anything. + await preflightOperationRun(opened); + if (laneRun.error) onLog(`Auto-${kind}: engine reachable again.`); + laneRun.run = opened; + laneRun.error = null; + } catch (err) { + const message = (err as Error)?.message ?? String(err); + if (laneRun.error !== message) onLog(`Auto-${kind} idle: ${message}`); + laneRun.run = null; + laneRun.error = message; + } finally { + laneRun.key = key; + laneRun.at = Date.now(); + laneRun.refreshing = false; + } + }; + // Transcription is throttled by the worker pool's slot count; download has no // pool, so its real limiter is the per-platform gate in next() — leave the // numeric cap at the policy max (or unbounded). runPool drops the target to 0 // itself while draining, so this never needs to check the drain signal. const limit = (): number => { - const policy = getSettings().autoQueue[kind]; + const settings = getSettings(); + const policy = settings.autoQueue[kind]; + const ceiling = policy.maxWorkers ?? Number.POSITIVE_INFINITY; + if (isOperationLane(kind)) { + // Kicked, not awaited — see the note on laneRun. A zero limit is a HOLD: + // runPool idle-waits, the job keeps its place, and the lane resumes on + // the poll after whatever it is holding for clears. + void refreshLaneRun(); + if (!laneRun.run) { + if (live.inFlight.size === 0 && laneRun.error) { + live.idleReason = "engine-unreachable"; + } + return 0; + } + const verdict = laneLimit( + settings, + operationLaneLive(laneRun.run, policy.maxWorkers ?? undefined), + ); + if (verdict.hold && live.inFlight.size === 0) { + // The lane's own gate, in the runner's vocabulary. `paused` is the + // operator's own hold; a yield or a spend cap is contention, and both + // resolve without anyone touching a switch. + live.idleReason = verdict.hold.reason === "paused" ? "lane-held" : "capped"; + } + return Math.min(verdict.limit, ceiling); + } if (kind !== "transcription") { - return policy.maxWorkers ?? Number.POSITIVE_INFINITY; + return ceiling; } const slots = eligibleSlots(); // A zero limit means runPool never calls next(), so this is the only place @@ -643,7 +910,7 @@ async function runLoop( ? "workers-paused" : "no-workers"; } - return Math.min(slots, policy.maxWorkers ?? Number.POSITIVE_INFINITY); + return Math.min(slots, ceiling); }; // Select the next unit and RESERVE its slot (so the next selection sees it), @@ -689,6 +956,27 @@ 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. + if (isOperationLane(kind)) { + await refreshLaneRun(); + if (!laneRun.run) { + live.idleReason = "engine-unreachable"; + return null; + } + } // THE UNATTENDED DISK GATE. This runner is the one path that dispatches // downloads for days with nobody watching, and until now it was the only // byte-writing path with no disk check at all — every manual action had a @@ -726,11 +1014,12 @@ async function runLoop( const slugToPlatform = new Map( metaCache.map((m) => [m.slug, m.platform ?? "unknown"]), ); + const laneOperations = laneOperationIds(kind); const { channels, owner } = await buildChannelWork( paths, kind, metaCache, - laneOperationIds(kind, policy.root), + laneOperations, ); // Re-derived each iteration from the freshly-read policy, like `enabled`, so // toggling the replace-auto-captions lane (or the ordering) takes effect @@ -750,11 +1039,23 @@ async function runLoop( policy.root, channels, defaultBucketsForPolicy(kind, policy), - compare ? { compare } : undefined, + { ...(compare ? { compare } : {}), defaultOperations: laneOperations }, ); - const exclude = new Set<string>([...live.inFlight.keys(), ...completed]); - removeIds(pending, exclude); + // A video already in flight is dropped whatever lane this is: two units on + // one video would race the same directory, and on an operation lane it is + // also how the dependency order is kept — diarization finishes before + // attribution-diarized is offered the same video. + removeIds(pending, new Set(live.inFlight.keys())); + dropCompleted(pending, policy.root, laneOperations); const pendingBeforeGates = countPending(pending); + // THE RUNNER JOB'S OWN BAR, on the metrics the per-channel jobs already + // use — so a lane's runner row reads like the manual verb's row rather than + // like an opaque long-lived loop. `target` moves as the corpus does (this + // is a loop, not a fixed work list), which is honest: it is "what is left + // right now", re-derived from the same projection the pick uses. + if (isOperationLane(kind)) { + reportRunnerProgress(pendingBeforeGates); + } // Download only: drop videos whose platform already has an in-flight // download (busy) OR is in a rate-limit/network backoff window (cooling // down), so a busy/throttled platform yields to the next-priority free one. @@ -845,23 +1146,89 @@ async function runLoop( return { pick, channelSlug, unitPlatform }; }; + // ONE UNIT ON AN OPERATION LANE. + // + // The pick is a VIDEO; which operation it needs is decided here, from disk, + // in dependency order. That is the price of `pending` being a list of ids: a + // leaf drawing the lane's whole union holds a video that may need diarization + // OR attribution, and the snapshot's per-operation lists cannot say which + // once they are merged into one leaf. Re-deriving is cheap (a readdir and a + // sidecar read per operation) and it is the invariant anyway — every pull + // re-derives. + // + // FIRST ONE WITH WORK WINS, and `run.operations` is already dependency + // ordered, so a video that needs diarization gets diarized now and attributed + // on a later pick — after the snapshot regeneration this run requests puts it + // in the attribution list. That is the dependency chain, walked one unit at a + // time. + const runOperationPick = async ( + picked: Picked, + runSignal: AbortSignal, + ): Promise<UnitResult> => { + const laneOperations = laneOperationIds(kind); + const opened = laneRun.run; + if (!opened) return { outcome: "skipped" }; + const policy = getSettings().autoQueue[kind]; + const leaf = flattenLeaves(policy.root).find( + (l) => l.id === picked.pick.leafId, + ); + const wanted = new Set( + leaf ? leafOperationIds(leaf, laneOperations) : laneOperations, + ); + const ops = opened.operations.filter((op) => wanted.has(op.id)); + const videoId = picked.pick.videoId; + const where = { channelSlug: picked.channelSlug, videoId }; + for (const op of ops) { + const plan = await classifyOperationUnit(opened, { ...where, op }); + if (plan.action !== "dispatch") { + markCompleted(videoId, op.id); + continue; + } + const task = tracker.start({ + id: videoId, + label: `${op.id} ${picked.channelSlug}/${videoId}`, + kind: kind === "digest" ? "digest" : "backfill", + }); + try { + const report = await runOperationUnit( + opened, + { ...where, op }, + { onLog: task.onLog, signal: runSignal }, + ); + markCompleted(videoId, op.id); + if (report.outcome === "done") return { outcome: "transcribed" }; + if (report.outcome === "failed") return { outcome: "failed" }; + return { outcome: "skipped" }; + } finally { + task.end(); + } + } + // Nothing on this lane has work for this video — a snapshot the regen has + // not caught up with, or a concurrent pass got there first. Retire it from + // the session's list, or the loop offers it again on every tick. + for (const op of ops) markCompleted(videoId, op.id); + return { outcome: "skipped" }; + }; + // Run one reserved unit, then release its reservation and fold its outcome // into counters/backoff. Mirrors the old `.finally` accounting exactly. const run = async (picked: Picked, runSignal: AbortSignal): Promise<void> => { const { pick, channelSlug, unitPlatform } = picked; let result: UnitResult = { outcome: "failed" }; try { - result = await launchUnit({ - kind, - paths, - pick, - channelSlug, - tracker, - onLog, - signal: runSignal, - drainSignal: ctx.drainSignal, - onChildJob: (jid) => childJobIds.set(pick.videoId, jid), - }); + result = isOperationLane(kind) + ? await runOperationPick(picked, runSignal) + : await launchUnit({ + kind, + paths, + pick, + channelSlug, + tracker, + onLog, + signal: runSignal, + drainSignal: ctx.drainSignal, + onChildJob: (jid) => childJobIds.set(pick.videoId, jid), + }); } catch (err) { result = { outcome: "failed" }; onLog(`Auto-${kind} unit error on ${pick.videoId}: ${String(err)}`); @@ -881,6 +1248,9 @@ async function runLoop( if (result.outcome === "transcribed") counters.transcribed++; else if (result.outcome === "failed") counters.failed++; else counters.skipped++; + if (isOperationLane(kind)) { + reportRunnerProgress(Math.max(0, lastPending - 1)); + } // Per-platform backoff: a rate-limit / network failure pauses the whole // platform with an exponential cooldown, and the video is deliberately NOT @@ -905,7 +1275,10 @@ async function runLoop( if (unitPlatform && result.outcome === "transcribed") { clearBackoff(kindState.platformBackoff, unitPlatform); } - markCompleted(pick.videoId); + // An operation lane retires (operation, video) pairs inside + // runOperationPick — it is the only thing that knows WHICH operation + // ran — so a bare id here would be a key nothing ever reads. + if (!isOperationLane(kind)) markCompleted(pick.videoId); } // Persist so the cooldown (and reset) survive a restart. persist(); @@ -1156,8 +1529,14 @@ export async function startAutoRunner( kind: AutoQueueKind, paths: Paths = getPaths(), ): Promise<string | null> { - const policy = getSettings().autoQueue[kind]; + 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") { @@ -1176,10 +1555,10 @@ export async function startAutoRunner( kind: autoRunnerJobKind(kind), queueKey: "", // run in parallel to everything; not serialized paths, - fn: async (onLog, signal, _setProgress, ctx) => { + fn: async (onLog, signal, setProgress, ctx) => { live.jobId = ctx.jobId; try { - await runLoop(kind, paths, live, onLog, signal, ctx); + await runLoop(kind, paths, live, onLog, setProgress, signal, ctx); } finally { if (singleton.runners.get(kind) === live) { singleton.runners.delete(kind); @@ -1193,20 +1572,25 @@ export async function startAutoRunner( return result.jobId; } +// 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. +export function startAutoRunnerBlockedReason(kind: AutoQueueKind): string | null { + return laneBlockedReason(getSettings(), kind); +} + // Start every enabled runner. Called from the editor instrumentation hook at // server startup. // -// NAMED, NOT ITERATED OVER LANES, and that is the whole safety story of slice -// 1.1. The digest and backfill lanes exist in the model — a policy, a tree, a -// gate, a status entry — and have no executor: `run()` below knows how to -// transcribe and how to download and nothing else. Iterating LANES here would -// start two runners that would pick a video and then not know what to do with -// it. They join this list in slice 1.2, when the executor does. +// ITERATED OVER LANES since slice 1.2, and that is the whole point of the +// slice: every lane now has an executor. `startAutoRunner` returns null for a +// lane whose policy is disabled (all four ship disabled but the two the +// operator has armed) and for one whose sweep is armed, so this list is safe to +// widen — it was named rather than iterated in 1.1 precisely because two of +// them would have picked a video and not known what to do with it. export async function startAutoRunnersIfEnabled( paths: Paths = getPaths(), ): Promise<void> { - await startAutoRunner("transcription", paths); - await startAutoRunner("download", paths); + for (const lane of LANES) await startAutoRunner(lane, paths); } // Stop a runner (hard cancel its job). In-flight units are aborted. Returns true diff --git a/common/controller/channels.ts b/common/controller/channels.ts @@ -291,6 +291,63 @@ export async function readChannelSnapshot( } } +// THE SAME READ, MEMOIZED ON (mtime, size) — for the ONE path that does it on a +// timer. +// +// The auto-queue runner projects every channel's snapshot on each scheduling +// tick, and `/api/auto-queue/status` folds FOUR lanes doing the same on a +// three-second poll. Since slice 1.2 the digest and backfill lanes draw their +// operation work lists from `snapshot.backfill[op].ids` — up to 11,329 strings +// on one channel — so that poll went from parsing nothing for two lanes to +// parsing all 68 snapshots (~6.5 MB) twice a tick, to produce an object +// identical to the one it produced three seconds earlier. +// +// Keyed on mtimeMs AND size, not on time: a snapshot regenerated by +// snapshotScheduler must be seen on the next tick, and it always is — the +// generator writes through a rename, so the mtime moves. A stat is ~5 µs +// against a parse of ~100 KB. +// +// DELIBERATELY NOT `readChannelSnapshot` ITSELF. That one is called by the +// snapshot GENERATOR and by half a dozen editor render paths, and a shared +// parsed object is only safe for callers that treat it as read-only. This +// function's callers do (they copy id arrays out and never write back); making +// the general reader share objects would be a much larger claim about code that +// has never had to honour it. +const snapshotMemo = new Map< + string, + { mtimeMs: number; size: number; snapshot: ChannelSnapshot } +>(); +// One entry per channel; the corpus has 68. The cap only bounds a pathological +// corpus, and dropping the whole map is correct because every entry is a cache. +const SNAPSHOT_MEMO_CAP = 512; + +export async function readChannelSnapshotShared( + paths: Paths, + slug: string, +): Promise<ChannelSnapshot | null> { + const file = snapshotPath(paths, slug); + try { + const st = await stat(file); + const hit = snapshotMemo.get(file); + if (hit && hit.mtimeMs === st.mtimeMs && hit.size === st.size) { + return hit.snapshot; + } + const snapshot = JSON.parse(await readFile(file, "utf8")) as ChannelSnapshot; + if (snapshotMemo.size >= SNAPSHOT_MEMO_CAP) snapshotMemo.clear(); + snapshotMemo.set(file, { mtimeMs: st.mtimeMs, size: st.size, snapshot }); + return snapshot; + } catch { + return null; + } +} + +// Test seam. Nothing in production drops the memo — an mtime change already +// invalidates an entry — but a test that writes a snapshot twice inside one +// filesystem timestamp tick needs to. +export function resetChannelSnapshotMemo(): void { + snapshotMemo.clear(); +} + // Slug + config + the precomputed snapshot for every channel. One readdir plus // two small file reads per channel, all concurrent: ~60 ms against the corpus // the walk above needs 4.4 s for. diff --git a/common/jobs/autoQueuePolicy.ts b/common/jobs/autoQueuePolicy.ts @@ -278,6 +278,17 @@ export function buildPendingByLeaf( // index (common/controller/recencyIndex.ts) when policy.order != "listed"; // omitting it reproduces the historical order byte-for-byte. compare?: (a: string, b: string) => number; + // THE OPERATION HALF OF `defaultBuckets`: what a leaf naming neither a + // bucket nor an operation draws from `ch.operations`. Empty for the two + // bucket lanes, which is what keeps every tree written before operations + // existed byte-identical; the digest and backfill lanes pass their whole + // dispatch set, so a catch-all leaf on them claims their work rather than + // being inert. + // + // Drawn AFTER the default buckets, so on a lane that somehow had both the + // bucket work would still be claimed first — a priority statement, not an + // implementation detail. + defaultOperations?: ReadonlyArray<string>; }, ): Record<string, string[]> { const leaves = flattenLeaves(root); @@ -292,22 +303,30 @@ export function buildPendingByLeaf( // existed has no operation anywhere, so every key is "\0" + id and the // dedup behaves exactly as it always has. const claimed = new Set<string>(); + const defaultOperations = opts?.defaultOperations ?? []; for (const leaf of leaves) { - const operation = leaf.match.operation ?? ""; // An operation leaf draws from its own id space and ignores buckets - // entirely; sanitizeMatch guarantees a leaf never names both. - const lists: ReadonlyArray<string> = operation - ? [operation] + // entirely; sanitizeMatch guarantees a leaf never names both. A leaf naming + // NEITHER draws both defaults — which on any given lane means exactly one + // of them, since a lane has buckets or operations and never both. + const draws: ReadonlyArray<{ list: string; operation: boolean }> = leaf.match + .operation + ? [{ list: leaf.match.operation, operation: true }] : leaf.match.bucket - ? [leaf.match.bucket] - : defaultBuckets; - for (const list of lists) { + ? [{ list: leaf.match.bucket, operation: false }] + : [ + ...defaultBuckets.map((b) => ({ list: b, operation: false })), + ...defaultOperations.map((o) => ({ list: o, operation: true })), + ]; + for (const draw of draws) { for (const ch of channels) { if (!matchesChannel(leaf.match, ch)) continue; - const ids = operation ? ch.operations?.[list] : ch.buckets[list]; + const ids = draw.operation + ? ch.operations?.[draw.list] + : ch.buckets[draw.list]; if (!ids) continue; for (const id of ids) { - const claim = `${operation}\u0000${id}`; + const claim = `${draw.operation ? draw.list : ""}\u0000${id}`; if (claimed.has(claim)) continue; claimed.add(claim); pending[leaf.id].push(id); diff --git a/common/jobs/jobKinds.test.ts b/common/jobs/jobKinds.test.ts @@ -58,6 +58,12 @@ const ADDED_KINDS: Record<string, { label: string; drainable: boolean }> = { label: "Purge superseded auto-captions", drainable: false, }, + // The two operation lanes' runners (phase 1 slice 1.2). Drainable like the + // two runners they join — a drain stops picking and lets the in-flight unit + // finish — and never replayable: a long-lived loop is not a unit of work, and + // "retry" on one would mean "start a second dispatcher". + "auto-digest": { label: "Auto-digest runner", drainable: true }, + "auto-backfill": { label: "Auto-backfill runner", drainable: true }, }; test("added kinds carry their pinned label and drainability", () => { diff --git a/common/jobs/jobKinds.ts b/common/jobs/jobKinds.ts @@ -63,6 +63,28 @@ const JOB_KINDS: Record<string, JobKindMeta> = { replayable: false, queueKeyStrategy: "platform", }, + // The two OPERATION lanes' runners. Same shape as the two above — one + // long-lived job per lane on queueKey "", drainable, never replayable — with + // one difference worth stating: their units make NO job record. A digest or a + // diarization dispatched by the runner runs in-process inside the loop, so + // there is no "auto-digest-unit" here to match auto-download-unit. That is + // what lets a lane work per video where the manual verbs must work per + // channel: the registry keeps 100 records and the log 500, and 55,956 digest + // jobs would evict the history of the run that made them. + "auto-digest": { + kind: "auto-digest", + label: "Auto-digest runner", + drainable: true, + replayable: false, + queueKeyStrategy: "parallel", + }, + "auto-backfill": { + kind: "auto-backfill", + label: "Auto-backfill runner", + drainable: true, + replayable: false, + queueKeyStrategy: "parallel", + }, "whisper-all": { kind: "whisper-all", label: "Transcribe all", diff --git a/common/lib/pauseGates.ts b/common/lib/pauseGates.ts @@ -93,6 +93,34 @@ 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. // diff --git a/editor/app/jobs/active/buildActiveJobs.ts b/editor/app/jobs/active/buildActiveJobs.ts @@ -507,6 +507,12 @@ function autoIdleNote( return "disk gate closed"; case "downloads-paused": return "downloads paused globally"; + case "lane-held": + return "lane held"; + case "lane-blocked": + return "its sweep is armed"; + case "engine-unreachable": + return "engine unreachable"; case "snoozed": return "snoozed"; case "disabled": diff --git a/editor/app/operations/components/dispatch.ts b/editor/app/operations/components/dispatch.ts @@ -108,7 +108,9 @@ export function idleReasonText( reason: AutoRunnerIdleReason | null, kind: AutoQueueKind, ): string | null { - const what = kind === "transcription" ? "transcribe" : "download"; + // "auto-transcribe" / "auto-download" for the two bucket lanes; the operation + // lanes' runners are named for the lane itself ("auto-digest"). + const what = kind === "transcription" ? "transcribe" : kind; switch (reason) { case "no-pending": return "nothing pending — every rule is empty"; @@ -128,6 +130,14 @@ export function idleReasonText( return "snoozed"; case "disabled": return `auto-${what} was switched off`; + case "lane-held": + return "the lane is held — resume it to start dispatching again"; + case "lane-blocked": + // NAMES THE SWEEP in the runner's own log; here it is the short form, + // because the operator is looking at the sweep panel beside this one. + return "this lane's sweep is armed — stop it and the runner resumes"; + case "engine-unreachable": + return "the lane's engine did not answer its probe — check the job log"; case "stopped": case null: return null; diff --git a/plans/tools/phase1-poll-timing.ts b/plans/tools/phase1-poll-timing.ts @@ -0,0 +1,46 @@ +#!/usr/bin/env tsx +// What the three-second status poll costs, over the live corpus. +// +// `/api/auto-queue/status` folds one `buildKind()` per lane, and the whole of +// `buildKind`'s cost is `computeLeafPending` — which reads and parses every +// channel snapshot the lane's projection touches. Slice 1.2 makes the digest +// and backfill lanes draw their operation union by default, so two lanes that +// read nothing start reading all 68 snapshots twice a tick. This is the meter +// that says whether the parse memo in `readChannelSnapshot` paid for itself. +// +// STRICTLY READ-ONLY, and offline: no editor, no runner, no snapshot regen. +// Usage, from anywhere in the repo: +// pnpm --filter yt-dlp-transcript-common exec tsx ../plans/tools/phase1-poll-timing.ts +// +// ROUNDS=n to change the sample count (default 3). Round 1 is cold — an empty +// memo and a cold page cache — and the rounds after it are what a poll running +// every three seconds actually pays. + +import { getPaths } from "../../common/lib/paths"; +import { computeLeafPending } from "../../common/controller/autoRunner"; +import { resetChannelSnapshotMemo } from "../../common/controller/channels"; +import { LANES, type AutoQueueKind } from "../../common/lib/autoQueueTypes"; + +const paths = getPaths(); +const ROUNDS = Number(process.env.ROUNDS ?? 3); +// NOMEMO=1 drops the parse memo before every round — what the poll would cost +// without it, which is the comparison the memo has to justify. +const NOMEMO = process.env.NOMEMO === "1"; + +for (let round = 0; round < ROUNDS; round++) { + if (NOMEMO) resetChannelSnapshotMemo(); + const started = Date.now(); + const per: string[] = []; + // Promise.all over the lanes, exactly as buildAutoQueueStatusPayload folds them. + await Promise.all( + LANES.map(async (lane: AutoQueueKind) => { + const at = Date.now(); + const pending = await computeLeafPending(lane, paths); + const total = Object.values(pending.counts).reduce((a, b) => a + b, 0); + per.push(`${lane}=${Date.now() - at}ms/${total}`); + }), + ); + console.log( + `round ${round + 1}: total ${Date.now() - started}ms [${per.sort().join(" ")}]`, + ); +}