import type { LiveInputs } from "./inputs"; import { type AutoRunnerStatus, type LeafPending, type NextUpView, type RecencyKeyView, laneDispatchRoot, } from "../controller/autoRunner"; import type { AutoQueueKind, AutoQueuePick, AutoQueueState, } from "../jobs/autoQueueState"; import { LANES } from "../lib/autoQueueTypes"; import { effectivePaceSeconds } from "../jobs/platformBackoff"; import type { AutoQueuePolicy } from "../jobs/autoQueuePolicy"; import { isGateHeld } from "../lib/pauseGates"; import type { FocusSummary } from "../lib/channelPriority"; import type { AutoQueueLanesPayload } from "./autoQueueLanes"; import { type PriorityView, laneFocusSummary } from "./channelPriority"; // Read-only payload for the Auto-Queue panel: per-kind runner status (running?, // what's in flight, in-flight counts per tree node), the effective policy, the // recent pick log, and snapshot-derived pending counts per leaf — so the UI can // show "cornbreadman: 12 pending" and which leaf is currently being serviced. // Backs a passive UI poll, like /api/scheduler/status. // A platform currently in a rate-limit (429) / network cooldown. `untilMs` is // the epoch ms the cooldown lapses; the client renders remaining time itself so // it counts down live. Only the download kind ever has these. export type PlatformCooldownView = { platform: string; untilMs: number; fails: number; // The platform's HOLD (release 17, slice RL): its failures outlasted the // cooldown cap, and `untilMs` is the next probe. Null for a plain cooldown. // A held platform is listed WHATEVER its probe time (review H2): once the // probe is overdue (`untilMs` passed) it stays held until a clean probe, a // clean manual run or Clear hold. `rateLimited` is false for a hold reached // through network failures alone. hold: { sinceMs: number; rateLimited: boolean } | null; }; // A platform whose request pace is above its static value (release 17, slice // RL): what every yt-dlp spawn against it now waits between requests. export type PlatformPaceView = { platform: string; sleepRequestsSeconds: number; baseSeconds: number; }; // A video whose subtitle fetch answered 429 while its media came down. // `held` once the count reached three (left alone for 7 days). export type SubtitleDeferralView = { videoId: string; channelSlug: string; count: number; lastAtMs: number; untilMs: number; held: boolean; }; // A video the auto-download pick skips until `untilMs` because it was // rate-limited recently (see jobs/unitOutcome.ts). Download kind only. export type VideoDeferralView = { videoId: string; channelSlug: string; untilMs: number; }; export type AutoQueueKindStatus = { kind: AutoQueueKind; // THE POLICY AS THE LANE DISPATCHES IT, not as it is stored. Every field is // the stored one except `root`, which is the COMPILED tree whenever // `settings.channelPriority` says anything (see ./channelPriority.ts). The // claim ladder is drawn from this, so a rung is a rule the runner actually // has — and with a priority model set the stored tree is not one. policy: AutoQueuePolicy; runner: AutoRunnerStatus; pendingByLeaf: Record; // The first few pending ids of each rule, in the order the runner would hand // them out — the drill-down behind each rung's count, and the only way to // check an ordering change by eye. pendingHeadByLeaf: Record; // videoId -> owning channel slug, for the drill-down's links. Covers exactly // the ids in pendingHeadByLeaf (plus nextUp's). ownerByVideo: Record; // videoId -> recency sort key, for the same ids, when the policy orders by // recency. Absent entries just render without a date. recencyByVideo: Record; // What the policy would dispatch next, and which rules it passed over. nextUp: NextUpView | null; picks: AutoQueuePick[]; // Platforms paused by a rate-limit/network backoff (download kind only). cooldowns: PlatformCooldownView[]; // Videos deferred after a rate limit, soonest to return first (download // kind only; empty elsewhere and whenever none is live). deferred: VideoDeferralView[]; // Platforms whose pace a rate limit raised, and the videos whose subtitles // are deferred (download kind only; release 17, slice RL). The deferrals // are listed while their window is open, soonest first. pace: PlatformPaceView[]; subtitleDeferred: SubtitleDeferralView[]; // This runner's lane gate — what the page's pause button draws, and what lets // the operations rail show a runner HOLDING for the first time. // // THE TWO KINDS ANSWER IT FROM DIFFERENT PLACES, on purpose. Transcription's // hold is LIVE, on the worker pool; the stored `autoQueue.transcription.held` // is only what the boot hook re-applies after a restart, and the e2e harness // rewrites settings wholesale between tests while the pool keeps its // pausedSnapshot. // Download has no live counterpart: its flag IS the gate, read at dispatch. held: boolean; // THE TREE ABOVE IS GENERATED, so the editor for it is read-only and the // channel priorities on /channels are where it is edited. False for a corpus // with no priorities set, which is every corpus until one is. policyCompiled: boolean; // THE FOCUS BANNER'S NUMBERS, FOR THIS LANE. Always present and inert when // `active` is false, so the banner is one component with one early return // rather than a conditional on the payload. The counts are this lane's own — // `focusPending`/`otherPending` differ per lane by construction. focus: FocusSummary; // The focus's display name, resolved on the server (the model stores a // siteId, and the banner does no I/O). Null when nothing is focused. focusName: string | null; }; export type AutoQueueStatusPayload = Record< AutoQueueKind, AutoQueueKindStatus > & { // The sweep half of the digest and backfill pipelines, plus the comparison // rail's bands. On the SAME payload as the runners rather than a second // endpoint, because the rail's whole purpose is that four lanes are read // together — two polls would let the rail and the focused lane disagree about // the same moment. // // Every lane now has a `status[lane]` entry too. The two are not duplicates: // this is what the SWEEPS are doing (armed, scoped, planned), and the entry // is what the lane's TREE says. Slice 1.3 deletes this half. lanes: AutoQueueLanesPayload; }; export type AutoQueueStatusInputs = Pick< LiveInputs, "settings" | "pool" | "now" > & { // Resolved ONCE for all four lanes by the shell: a focus set costs a channel // listing and, for a site focus, a sites read, and the four lanes compile from // the same one. priority: PriorityView; // Per-lane runner status. Injected because asking for it CONSTRUCTS the // runner, which a status poll must never do. runner: Record; // The auto-queue state document, read once for all four lanes (it used to be // read once per lane, four parses of the same file per poll). state: AutoQueueState; // Per-lane snapshot-derived pending work. pending: Record; lanes: AutoQueueLanesPayload; }; function buildKind( kind: AutoQueueKind, inputs: AutoQueueStatusInputs, ): AutoQueueKindStatus { const { settings, pool, priority } = inputs; const stored = settings.autoQueue[kind]; // The stored policy, with the DISPATCHED root in place of the stored one. // Spread rather than rebuilt, so `held`, `snoozeUntil`, `enabled`, `order` // and `maxWorkers` come through untouched — the same rule every writer of a // policy in this repo follows. const policy: AutoQueuePolicy = { ...stored, root: laneDispatchRoot(kind, stored, priority, priority.slugs), }; const runner = inputs.runner[kind]; const pending = inputs.pending[kind]; const now = inputs.now(); // Surface platforms still inside their cooldown window so the operator can see // why an otherwise-pending platform isn't being serviced (and a manual sync on // it would be refused). A manual sync/download 429 writes into this same map. const cooldowns: PlatformCooldownView[] = Object.entries( inputs.state[kind].platformBackoff, ) .filter(([platform, e]) => e.until > now || Boolean(inputs.state[kind].platformHolds?.[platform])) .map(([platform, e]) => { const hold = inputs.state[kind].platformHolds?.[platform]; return { platform, untilMs: e.until, fails: e.fails, hold: hold ? { sinceMs: hold.since, rateLimited: hold.rateLimited } : null, }; }) .sort((a, b) => b.untilMs - a.untilMs); const pace: PlatformPaceView[] = Object.entries( inputs.state[kind].platformPace ?? {}, ) // At the pace the hourly easing leaves it (effectivePaceSeconds); an // entry already back at its base is not listed. .map(([platform, p]) => ({ platform, sleepRequestsSeconds: effectivePaceSeconds(p, now), baseSeconds: p.baseSeconds, })) .filter((p) => p.sleepRequestsSeconds > p.baseSeconds) .sort((a, b) => (a.platform < b.platform ? -1 : a.platform > b.platform ? 1 : 0)); const subtitleDeferred: SubtitleDeferralView[] = Object.entries( inputs.state[kind].subtitleDeferrals ?? {}, ) .filter(([, d]) => d.until > now) .map(([videoId, d]) => ({ videoId, channelSlug: d.channelSlug, count: d.count, lastAtMs: d.lastAt, untilMs: d.until, held: d.count >= 3, })) .sort( (a, b) => a.untilMs - b.untilMs || (a.videoId < b.videoId ? -1 : a.videoId > b.videoId ? 1 : 0), ); const deferred: VideoDeferralView[] = Object.entries( inputs.state[kind].videoDeferrals ?? {}, ) .filter(([, d]) => d.until > now) .map(([videoId, d]) => ({ videoId, channelSlug: d.channelSlug, untilMs: d.until, })) // Ties by id in CODE-POINT order, not localeCompare: YouTube ids are // mixed-case and the strip's order must not depend on the server's locale. .sort( (a, b) => a.untilMs - b.untilMs || (a.videoId < b.videoId ? -1 : a.videoId > b.videoId ? 1 : 0), ); return { kind, policy, runner, pendingByLeaf: pending.counts, pendingHeadByLeaf: pending.head, ownerByVideo: pending.owner, recencyByVideo: pending.recency, nextUp: pending.nextUp, picks: inputs.state[kind].picks, cooldowns, deferred, pace, subtitleDeferred, // THE ASYMMETRY IS TRANSCRIPTION'S ALONE, and it is not a special case for // "the first lane": its hold is LIVE on the worker pool, while every other // lane's flag IS its gate. See lib/pauseGates.ts. held: kind === "transcription" ? pool.isPaused() : isGateHeld(settings, kind), policyCompiled: priority.compiled, // Keyed by COMPILED leaf id (`prio-focus-`), which is why this is // computed here and not on the client: it is only meaningful against the // counts of the tree the lane dispatches from. focus: laneFocusSummary(priority, pending.counts), focusName: priority.name, }; } // ONE ENTRY PER LANE, off LANES rather than off two hand-written names — so a // lane added to the model appears on this payload with no edit here, which is // the whole point of the widened type. export function buildAutoQueueStatusPayload( inputs: AutoQueueStatusInputs, ): AutoQueueStatusPayload { return { ...(Object.fromEntries( LANES.map((lane) => [lane, buildKind(lane, inputs)]), ) as Record), lanes: inputs.lanes, }; } // THE POLL'S EXPENSIVE HALF, SHARED: single-flight plus a short memo. // // The payload above is cheap; what it is built FROM is not. Each lane's // pending work is a fold over every channel's snapshot (the shell's // `computeLeafPending`, four lanes), and every surface that shows the lanes // polls it every 3 s (`useOperationsStatus`). Each poll used to compute its // own: on 2026-10-01, with the main thread busy regenerating a report, one // computation took 96 s, and every tab's poll started another behind it — a // queue that only grew. Now concurrent callers share the computation in flight, // and a result is reused for AUTO_QUEUE_STATUS_MEMO_MS after it lands, so N // pollers cost one fold per window. // // KEYED BY THE SETTINGS THE FOLD READS, AND OTHERWISE BY TIME. The fold reads // settings itself (`computeLeafPending`: each lane's policy and rule tree, // and `channelPriority` — paused channels, the focus, the compiled leaf ids), // so the shell passes a key built from those (`autoQueue` + `channelPriority`): // a rule added, a focus set or a channel paused changes the key and misses the // memo, so a count is never keyed by a tree that is no longer the one drawn. // What nothing tells the memo about — a snapshot rewritten, the runner picking // a video (its in-flight set is subtracted inside the fold) — is up to // AUTO_QUEUE_STATUS_MEMO_MS behind; the next poll is the correction. A lane's // hold, the runner's status and its picks are not behind it at all: the shell // reads them on every call (editor/app/operations/status.ts). export const AUTO_QUEUE_STATUS_MEMO_MS = 3_000; export type SingleFlightMemo = { // The memoized value while it is fresh AND was computed under `key`; else the // computation in flight under `key`; else `compute()`, started now and shared // with every caller asking with `key` until it settles. A rejection is not // memoized: the next caller computes again. get(key: string, compute: () => Promise): Promise; // Forget the value AND detach the computation in flight, whose result is // then dropped rather than stored (a test reset must not be answered with // the previous fixture's numbers). clear(): void; }; export function singleFlightMemo( opts: { ttlMs?: number; now?: () => number } = {}, ): SingleFlightMemo { const ttlMs = opts.ttlMs ?? AUTO_QUEUE_STATUS_MEMO_MS; const now = opts.now ?? Date.now; let value: { key: string; v: T; at: number } | null = null; let inFlight: { key: string; p: Promise; seq: number } | null = null; // Every computation started gets the next number; only the latest started // may store its value or clear the in-flight slot, so a computation under an // old key that lands late never overwrites a newer one. clear() bumps it too. let seq = 0; return { get(key, compute) { if (value && value.key === key && now() - value.at < ttlMs) { return Promise.resolve(value.v); } if (inFlight && inFlight.key === key) return inFlight.p; const mine = ++seq; const p = compute().then( (v) => { if (seq === mine) { value = { key, v, at: now() }; inFlight = null; } return v; }, (err: unknown) => { if (seq === mine) inFlight = null; throw err; }, ); inFlight = { key, p, seq: mine }; return p; }, clear() { seq++; value = null; inFlight = null; }, }; } // The editor's one instance, on globalThis like the registry: the polled route // and the operations pages are separate bundles, and the e2e reset route // (editor/app/api/test/invalidate-cache) drops it through the same global. declare global { // eslint-disable-next-line no-var var __yttAutoQueueStatusMemo__: SingleFlightMemo | undefined; } export function autoQueueStatusMemo(): SingleFlightMemo { if (!globalThis.__yttAutoQueueStatusMemo__) { globalThis.__yttAutoQueueStatusMemo__ = singleFlightMemo(); } return globalThis.__yttAutoQueueStatusMemo__ as SingleFlightMemo; }