import type { Platform } from "../lib/platform"; // The persisted tree/settings SHAPE lives in the model layer — lib/ may not // import jobs/, and settings.ts has to name an AutoQueueSettings. Re-exported // here so every existing `from "../jobs/autoQueuePolicy"` import still resolves. import type { AutoQueueGroup, AutoQueueKind, AutoQueueLeaf, AutoQueueMatch, AutoQueueMatchType, AutoQueueMode, AutoQueueNode, AutoQueueOrder, AutoQueuePolicy, AutoQueueSettings, } from "../lib/autoQueueTypes"; import { LANES, isGroup } from "../lib/autoQueueTypes"; export type { AutoQueueGroup, AutoQueueKind, AutoQueueLeaf, AutoQueueMatch, AutoQueueMatchType, AutoQueueMode, AutoQueueNode, AutoQueueOrder, AutoQueuePolicy, AutoQueueSettings, }; export { LANES, isGroup }; // THE SANITIZER MOVED DOWN A LAYER (one-core phase 3 slice 4a). Everything that // turns a raw settings.json value into a legal `AutoQueueSettings` now lives in // `lib/autoQueueSchema.ts`, with the rest of the settings schema — which is what // let `lib/settings.ts` stop importing this file. Re-exported here so every // existing `from "../jobs/autoQueuePolicy"` import still resolves, exactly as // the TYPES above are. export { AUTO_QUEUE_MODES, AUTO_QUEUE_ORDERS, AUTO_QUEUE_MAX_WORKERS_MAX, defaultAutoQueue, defaultAutoQueuePolicy, sanitizeAutoQueue, sanitizeAutoQueueOrder, } from "../lib/autoQueueSchema"; // Pure, side-effect-free policy engine for the automatic priority queue. It // decides WHICH pending video to process next, across all channels, from a // configurable tree. The runner (common/controller/autoRunner.ts) does the I/O // (reads snapshots to build the pending sets, acquires worker leases, persists // fairness state). Keeping this pure makes the selection rules unit-testable // without a server — exactly how syncScheduler.ts splits selectDueChannels from // the editor tick route. // // The tree is a direct analog of Linux HTB / cgroup hierarchies: internal nodes // (groups) carry a `mode` governing how their children compete; leaves carry a // `match` describing which videos they own. A capped subtree (its in-flight // worker count at its `maxWorkers` ceiling) reads as "no work" and the algorithm // falls through to the next-priority sibling — exactly like HTB's class ceil. // Buckets each runner kind draws from, in priority order. A leaf with no // explicit bucket draws from the whole list (union, deduped); the list order is // its internal priority. Single source of truth for the runner, the pending- // count helper, and the editor's bucket picker. export const TRANSCRIBE_BUCKETS = ["downloadedNoTranscript", "failedListed"] as const; // `chatOnlyPending` is LAST on purpose: it is the smallest and cheapest work on // the lane (one --skip-download pass per video, no media), and putting it ahead // of real downloads would let a chat backlog delay the corpus. Empty for every // channel that has not set `downloadFilter.rejectedLivestreams` — which is // every channel that predates the field — so the fold, the lane's work list and // the pick order are byte-identical for them. export const DOWNLOAD_BUCKETS = [ "partialDownloads", "undownloadedIds", "chatOnlyPending", ] as const; // Buckets a runner will NOT draw from unless asked. Replacing YouTube's // auto-captions with our own transcript costs an audio download plus a // transcription per video, on a corpus where ASR-only videos outnumber // manually-captioned ones ~100:1 — so it is never part of the default union. // Two ways in, both explicit: a leaf can target the bucket by name (per-channel // opt-in, see selectableBucketsForKind), or the runner's `replaceAutoSubs` // switch appends it to the TAIL of the default union (see // defaultBucketsForPolicy) — strictly lowest priority, since buildPendingByLeaf // walks buckets in list order and claiming is first-match-wins. export const TRANSCRIBE_OPT_IN_BUCKETS = ["downloadedAutoSubsOnly"] as const; export const DOWNLOAD_OPT_IN_BUCKETS = ["autoSubsOnly"] as const; // NO BUCKETS IS A REAL ANSWER, not a gap. The digest and backfill lanes draw // from operations, never from snapshot buckets — `snapshot.backfill[op].ids` is // their work list — so an empty array here is what says "this lane has no // bucket picker and no bucket projection", and every caller already handles an // empty list the way it handles a leaf that matches nothing. // // (The type is AutoQueueKind now that it lives in lib/autoQueueTypes.ts, which // this module already imports. The old note here warned against importing it // from autoQueueState.ts, which imports FROM this module — that cycle is gone // with the type.) export function bucketsForKind(kind: AutoQueueKind): readonly string[] { if (kind === "transcription") return TRANSCRIBE_BUCKETS; if (kind === "download") return DOWNLOAD_BUCKETS; return []; } export function optInBucketsForKind(kind: AutoQueueKind): readonly string[] { if (kind === "transcription") return TRANSCRIBE_OPT_IN_BUCKETS; if (kind === "download") return DOWNLOAD_OPT_IN_BUCKETS; return []; } // Everything a leaf may be pointed at for this kind: the default union plus the // opt-in buckets. Backs the editor's bucket dropdown AND the runner's snapshot // projection — a leaf can only find ids in a bucket the runner projected. export function selectableBucketsForKind( kind: AutoQueueKind, ): readonly string[] { return [...bucketsForKind(kind), ...optInBucketsForKind(kind)]; } // The bucket list a leaf with NO explicit bucket draws from. Identical to // bucketsForKind unless the policy opted into auto-caption replacement, which // appends the opt-in buckets at the tail so real work always drains first. export function defaultBucketsForPolicy( kind: AutoQueueKind, policy: Pick, ): readonly string[] { return policy.replaceAutoSubs ? [...bucketsForKind(kind), ...optInBucketsForKind(kind)] : bucketsForKind(kind); } // WHERE A BUCKET'S IDS LIVE IN A SNAPSHOT. One place, because the answer is not // uniform: `undownloadedIds` is a TOP-LEVEL field and every other bucket is // under `snapshot.buckets`, and that asymmetry was spelled out at each of the // two places that walked the list. export type BucketSource = { buckets?: Record | null; undownloadedIds?: string[] | null; }; // Returns the STORED array, not a copy. Every caller reads it, and this runs // 68 times per bucket per scheduling tick and per three-second status poll — // copying up to 11,000 strings there is the cost readChannelSnapshotShared's // memo exists to avoid. export function bucketIdsFrom(source: BucketSource, name: string): string[] { if (name === "undownloadedIds") return source.undownloadedIds ?? []; return source.buckets?.[name] ?? []; } // A BUCKET LANE'S WHOLE WORK LIST, folded from its default buckets in priority // order with duplicates dropped — the exact list a bucket-less leaf on this lane // would draw before the opt-in buckets are appended. // // ONE FOLD, TWO CALLERS, and that is the point of slice 1.5. The snapshot // GENERATOR writes it as `backfill.download.ids` / `backfill.transcription.ids` // so every lane's work list is one array in the snapshot, and the runner reads // it back — falling back to this same fold for a snapshot written before those // entries existed. Because both sides are this function, the migration cannot // change what a leaf draws: the entry a fresh snapshot carries is byte-identical // to the fold the old one gets. // // UNSORTED, DELIBERATELY. `undownloadedIds` is playlist order — newest first, // which is the auto-download runner's queue — and sorting it here would silently // reorder the download lane. The dedup preserves first-seen order for the same // reason. // // The OPT-IN buckets are NOT here. They enter only through // defaultDrawsForPolicy, at the tail, exactly as defaultBucketsForPolicy has // always appended them — a snapshot entry that folded them in would make // `replaceAutoSubs` a property of the corpus instead of of the policy. export function bucketLaneWorkIds( kind: AutoQueueKind, source: BucketSource, ): string[] { const out: string[] = []; const seen = new Set(); for (const name of bucketsForKind(kind)) { for (const id of bucketIdsFrom(source, name)) { if (seen.has(id)) continue; seen.add(id); out.push(id); } } return out; } // THE NAMES A BUCKET-LESS LEAF DRAWS, in priority order — what the runner // actually projects, as opposed to what `defaultBucketsForPolicy` describes. // // The two differ on a bucket lane and only there. Since slice 1.5 the runner // draws that lane's whole work list as ONE list, named for the operation the // lane dispatches (`laneWorkList`, from operations.bucketLaneOperationId) and // projected from `snapshot.backfill[op].ids` with the fold above as its // fallback. The result is identical to walking the default buckets one by one — // bucketLaneWorkIds IS that walk — so nothing a leaf draws moves; what changes // is that the list has a name, a home in the snapshot, and one definition. // // It stays in the BUCKET claim space (buildPendingByLeaf's `\0id`, not // `op\0id`), because on a bucket lane the operation and the bucket union are the // same work: a video claimed by an explicit `failedListed` leaf must not be // claimed AGAIN by a catch-all leaf drawing the operation. That is why the id is // projected into `ChannelWork.buckets` rather than into `.operations` — the // operation claim space exists to keep diarization and attribution apart on ONE // video, and there is no second operation here to keep apart from. // // `defaultBucketsForPolicy` is unchanged and still answers the question // policyDrawsBucket asks ("would the runner draw this bucket"), which is about // the POPULATION and not about the projection. export function defaultDrawsForPolicy( kind: AutoQueueKind, policy: Pick, laneWorkList: string | null, ): readonly string[] { if (!laneWorkList) return defaultBucketsForPolicy(kind, policy); return policy.replaceAutoSubs ? [laneWorkList, ...optInBucketsForKind(kind)] : [laneWorkList]; } // Would this runner kind, under this policy, draw `bucket` for this channel? // // The question a lane outside the runner has to ask before it hands work over: // "if I leave this video in that bucket, will the runner actually pick it up?" // It is NOT `policy.replaceAutoSubs` — that flag is only one of the two ways in. // A leaf naming an operation draws no bucket at all; a leaf naming a bucket // draws only that one (per-channel opt-in, no flag needed); a bucket-less leaf // draws defaultBucketsForPolicy — which is where replaceAutoSubs enters, and the // only place it does. // // Pure, like everything else here: the caller supplies the channel's slug and // platform, never a snapshot. export function policyDrawsBucket( kind: AutoQueueKind, policy: Pick, channel: Pick, bucket: string, ): boolean { const defaults = defaultBucketsForPolicy(kind, policy); return flattenLeaves(policy.root).some( (l) => !l.match.operation && matchesChannel(l.match, channel) && (l.match.bucket ? l.match.bucket === bucket : defaults.includes(bucket)), ); } // --- Runtime fairness state (persisted best-effort by autoQueueState.ts) ---- // Smooth Weighted Round-Robin (SWRR, nginx's algorithm) current-weight per node // id. Round-robin mode is SWRR with every weight forced to 1; weighted-fair uses // each child's `weight`. Mutated in place by selectNextWork when a group makes a // choice, so fairness advances only on an actual grant. export type AutoQueueRuntime = { currentWeights: Record; }; export function emptyAutoQueueRuntime(): AutoQueueRuntime { return { currentWeights: {} }; } // In-flight worker counts keyed by node id (a leaf and every ancestor on the // path to it). The runner increments each id on grant and decrements on // release; the resolver reads them to enforce maxWorkers ceilings. export type ActiveCounts = Record; // --- Pending-set construction (pure) --------------------------------------- // One channel's available work, grouped by snapshot bucket. The runner builds // these from each channel's snapshot.json (already filtered to genuinely // actionable ids), then hands them here so leaf matching stays pure/testable. export type ChannelWork = { slug: string; platform: Platform | null; // bucket name -> available video ids in that bucket buckets: Record; // operation id -> reachable video ids for that operation, from // snapshot.backfill[op].ids. A SEPARATE map for the OPERATION lanes, because // it is a separate CLAIM SPACE: buildPendingByLeaf claims a video as // `${operation}\0${id}`, which is what lets one video be pending for // diarization AND for attribution-diarized at the same time. Folding those // into `buckets` would collapse them into one claim and silently drop the // second operation's work. // // THE BUCKET LANES DELIBERATELY GO THE OTHER WAY (slice 1.5). Their work list // is also `snapshot.backfill[op].ids`, but it is projected into `buckets` // under the operation's id, because on those lanes the operation and the // bucket union are the SAME work — there is no second operation to keep // apart, and sharing the bucket claim space is what stops a catch-all leaf // re-claiming a video an explicit retry-bucket leaf already took. See // defaultDrawsForPolicy and controller/autoRunner.ts buildChannelWork. // // Optional: every projection written before operations existed omits it, and // a leaf naming an operation simply finds nothing. operations?: Record; // THE CHANNEL-SCOPED WORK A LEAF CANNOT NAME. The metadata scan is a fact // about the CHANNEL, not about a video — there is no id to put in a bucket, // because the whole point of the scan is that these videos have no directory // — so the download lane's pre-pick reads it off the projection here rather // than through `pick()`. Both optional: a projection that omits them offers // no scan, which is what every caller but the download runner wants. // // `scanUnscanned` is the snapshot's own `metadataScan.unscanned`, which is // defined to mean exactly what metadataScanTargets() will fetch. // `filtered` is "this channel has a download filter that compiles" — without // one a scan settles nothing and is pure cost against the source. scanUnscanned?: number; filtered?: boolean; }; // `retainLeaves(pending, root, "buckets" | "operations")` USED TO LIVE HERE, and // 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 (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 // operationsForLane gives it (see controller/autoRunner.ts buildChannelWork). // A leaf asking for anything else finds no list to draw from and comes back // empty on its own — the same zero, reached by construction instead of by a // second pass that had to be remembered at three call sites. // Pre-order (priority-order) flatten of every leaf in the tree. export function flattenLeaves(node: AutoQueueNode): AutoQueueLeaf[] { if (!isGroup(node)) return [node]; const out: AutoQueueLeaf[] = []; for (const child of node.children) out.push(...flattenLeaves(child)); return out; } // Exported for policyDrawsBucket's callers outside this module (the backfill // hand-off asks it about one channel it holds only a slug and a platform for), // which is why the parameter is the narrowest shape that answers the question // rather than a whole ChannelWork. export function matchesChannel( match: AutoQueueMatch, ch: Pick, ): boolean { if (match.type === "all") return true; if (match.type === "channel") return !!match.value && ch.slug === match.value; if (match.type === "platform") { return !!match.value && ch.platform === match.value; } return false; } // Assign every available video to exactly ONE leaf: the first (highest-priority, // pre-order) leaf whose match covers the video's channel and whose bucket // contains it. First-match-wins prevents the same video being claimed — and thus // double-processed — by two overlapping leaves (e.g. a channel leaf and an `all` // catch-all). A leaf with an explicit `match.bucket` draws from only that bucket; // a leaf with none draws from the union of `defaultBuckets` (the kind's full // list), in list order. The `claimed` Set dedups across AND within leaves, so a // video present in two buckets is taken once, under the earlier (higher-priority) // bucket. Returns leafId -> available video ids (in priority order). export function buildPendingByLeaf( root: AutoQueueNode, channels: ReadonlyArray, defaultBuckets: ReadonlyArray, opts?: { // Optional total order applied to each leaf's FINISHED list, after claiming. // Claiming itself is untouched — a video in two buckets is still attributed // once, to the higher-priority bucket — so this only decides which of a // leaf's own videos goes first. Supplied by the runner from the recency // 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; }, ): Record { const leaves = flattenLeaves(root); const pending: Record = {}; for (const leaf of leaves) pending[leaf.id] = []; // Claims are keyed `${operation}\0${id}`, NOT by id. // // The same video is legitimately pending for digest AND for diarization — // they are different work on the same input — so an id-keyed set would let // whichever leaf happened to run first steal the other operation's work and // silently drop it from the queue. Every tree written before operations // existed has no operation anywhere, so every key is "\0" + id and the // dedup behaves exactly as it always has. const claimed = new Set(); const defaultOperations = opts?.defaultOperations ?? []; for (const leaf of leaves) { // An operation leaf draws from its own id space and ignores buckets // 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 ? [{ 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 = draw.operation ? ch.operations?.[draw.list] : ch.buckets[draw.list]; if (!ids) continue; for (const id of ids) { const claim = `${draw.operation ? draw.list : ""}\u0000${id}`; if (claimed.has(claim)) continue; claimed.add(claim); pending[leaf.id].push(id); } } } } if (opts?.compare) { const compare = opts.compare; for (const ids of Object.values(pending)) ids.sort(compare); } return pending; } // --- Core selection (pure) -------------------------------------------------- export type WorkPick = { leafId: string; videoId: string; // Node ids from root..leaf. The runner increments active counts for every id // on grant and decrements them on release, so ancestor caps are enforced. path: string[]; }; function capSaturated(node: AutoQueueNode, active: ActiveCounts): boolean { const cap = node.maxWorkers; if (cap == null) return false; return (active[node.id] ?? 0) >= cap; } // Whether a subtree has at least one available, non-cap-saturated video. function hasWork( node: AutoQueueNode, pending: Record, active: ActiveCounts, ): boolean { if (capSaturated(node, active)) return false; if (!isGroup(node)) return (pending[node.id]?.length ?? 0) > 0; return node.children.some((c) => hasWork(c, pending, active)); } function weightOf(node: AutoQueueNode, mode: AutoQueueMode): number { if (mode !== "weighted-fair") return 1; const w = node.weight; return typeof w === "number" && w > 0 ? w : 1; } // Pick one leaf within a group via Smooth Weighted Round-Robin over only its // children that currently have work. Mutates runtime.currentWeights. The chosen // child is guaranteed to have work, so the recursive descent never dead-ends. function pickSWRR( group: AutoQueueGroup, eligible: AutoQueueNode[], runtime: AutoQueueRuntime, ): AutoQueueNode { const cw = runtime.currentWeights; let total = 0; let best: AutoQueueNode | null = null; for (const child of eligible) { const w = weightOf(child, group.mode); total += w; cw[child.id] = (cw[child.id] ?? 0) + w; if (!best || cw[child.id] > cw[best.id]) best = child; } best = best as AutoQueueNode; cw[best.id] -= total; return best; } function pick( node: AutoQueueNode, pending: Record, runtime: AutoQueueRuntime, active: ActiveCounts, parents: string[], ): WorkPick | null { if (capSaturated(node, active)) return null; const path = [...parents, node.id]; if (!isGroup(node)) { const ids = pending[node.id]; if (!ids || ids.length === 0) return null; return { leafId: node.id, videoId: ids[0], path }; } const eligible = node.children.filter((c) => hasWork(c, pending, active)); if (eligible.length === 0) return null; if (node.mode === "strict") { // Highest-priority child with work = first in declared order. return pick(eligible[0], pending, runtime, active, path); } // round-robin (all weights 1) and weighted-fair share one SWRR path. const chosen = pickSWRR(node, eligible, runtime); return pick(chosen, pending, runtime, active, path); } // Select the next video to process, or null when nothing is available (every // leaf empty or every path cap-saturated). Mutates `runtime` (advancing fairness // only on an actual grant). The caller increments `active` for each id in the // returned path on start and decrements on release. export function selectNextWork( root: AutoQueueNode, pending: Record, runtime: AutoQueueRuntime, active: ActiveCounts = {}, ): WorkPick | null { return pick(root, pending, runtime, active, []); }