import path from "node:path"; import type { Paths } from "../lib/paths"; import type { ChannelConfig } from "../lib/channelConfig"; import { mapConcurrent } from "../lib/concurrency"; import { getPaths } from "../lib/paths"; import { getSettings, type SiteSettings } from "../lib/settings"; import { diskGate } from "../lib/diskSpace"; import { formatBytes } from "../lib/format"; import { detectPlatform } from "../lib/platform"; import { compileDownloadFilter } from "../lib/downloadFilters"; import { runMetadataScanJob } from "./metadataScanJob"; import { getWorkerPool } from "../jobs/workerPool"; import { getRegistry } from "../jobs/registry"; import { runManagedFunction } from "../jobs/streamCommand"; 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, type WorkPick, buildPendingByLeaf, flattenLeaves, selectNextWork, bucketIdsFrom, bucketLaneWorkIds, defaultDrawsForPolicy, selectableBucketsForKind, } from "../jobs/autoQueuePolicy"; import { bucketLaneOperationId, operationsForLane, type Operation, } from "../lib/operations"; import { cheapestComparator, classifyOperationUnit, laneLimit, openOperationRun, operationLaneLive, preflightOperationRun, runOperationUnit, videoDurationSeconds, type OperationLane, type OperationRun, } from "./operationBatch"; import { type RecencyKey, buildRecencyKeys, makeRecencyComparator, } from "./recencyIndex"; // 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, type AutoQueueState, liveAutoQueueState, readAutoQueueState, sharedAutoQueueState, recordPick, writeAutoQueueState, } from "../jobs/autoQueueState"; import { currentPaceSeconds, downloadGapMs, isCoolingDown, isPlatformHeld, isVideoDeferred, mergeBackoffEntry, pruneDeferred, prunePlatformPacing, pruneSubtitleDeferrals, } from "../jobs/platformBackoff"; import { platformMinGapSeconds, staticSleepRequestsSeconds, } from "../ytdlp/platformArgs.mjs"; import { applyUnitOutcome } from "../jobs/unitOutcome"; import { type DownloadFailureClass } from "../lib/availability"; import { resolveCookiePolicy } from "../lib/cookiePolicy"; import { isAutoSubsOnly } from "../lib/subtitleProvenance"; import { readVideoFiles } from "../lib/videoStatus"; import { type DownloadOutcomeStatus } from "../lib/downloadOutcome"; import { downloadQueueKey } from "../lib/queueKeys"; import { isGateHeld } from "../lib/pauseGates"; import { inspectChannelMedia, markerHoldsText, readRelocationMarker, type ChannelMediaLocation, type ChannelMediaStatus, } from "../lib/channelMedia"; import { isMediaHeld, isTextHeld, mediaHoldText } from "../lib/channelMediaHold"; import { type ChannelPriority, type FocusSummary, type SiteChannelIndex, compileLaneRoot, focusSummary, isChannelPaused, isDefaultChannelPriority, resolveFocusSlugs, } from "../lib/channelPriority"; import { siteChannelIndex } from "../lib/site"; import { listChannelConfigs, readChannelConfig, readChannelSnapshotShared, } from "./channels"; import { transcribeOneFromQueue } from "./transcribeOneFromQueue"; import { findVideoSourceUrl } from "./undownloadedVideos"; import { downloadOneManaged } from "../ytdlp/downloadOneManaged"; import { resolveDownloadFormatPreset } from "../ytdlp/downloadFormat"; // The automatic priority-queue runners. Each kind (auto-transcribe / // auto-download) is ONE long-lived registry job (queueKey "" so it runs in // parallel to everything, with working cancel + drain). The job's body is a // fill-to-capacity loop: it reads pending work from channel snapshots, asks the // pure policy (autoQueuePolicy.selectNextWork) which video to do next across all // channels, and runs it. // // Priority is preserved by capping the number of OUTSTANDING units to the count // of eligible worker slots — never dumping a backlog into the worker pool's FIFO // park queue (which would decide order by hardware, not policy). Each time a slot // frees, the loop re-picks from fresh policy, so a higher-priority channel's // video claims the next freed slot (non-destructive "preempt at next slot"). export const AUTO_TRANSCRIBE_KIND = "auto-transcribe"; export const AUTO_DOWNLOAD_KIND = "auto-download"; // One per-video download launched by the auto-download runner. Distinct from // the long-lived runner job (AUTO_DOWNLOAD_KIND): each unit is a real managed // job on the channel's platform queue (downloadQueueKey), so it serializes // 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 = { transcription: AUTO_TRANSCRIBE_KIND, download: AUTO_DOWNLOAD_KIND, digest: AUTO_DIGEST_KIND, backfill: AUTO_BACKFILL_KIND, }; export function autoRunnerJobKind(kind: AutoQueueKind): string { 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 // re-reading snapshots (new downloads may have appeared, or a worker re-enabled). const IDLE_POLL_MS = 3000; // Re-read channel list + platforms at most this often (configs change rarely). const CHANNEL_LIST_TTL_MS = 30_000; // Cap on the session "recently completed" set, so a long-running server doesn't // grow it unbounded. FIFO eviction; the snapshot catches up within ~1s anyway. const COMPLETED_CAP = 5000; // THE DOWNLOAD LANE'S CHANNEL-SCOPED UNIT. A metadata scan has no video id, so // it gets a synthetic one under a prefix nothing else can produce (a video id // never contains a space) and a leaf id that exists in no tree — the pick is // not made by `selectNextWork` and never reaches one. Both are exported so the // console and the tests name the same strings the runner writes. export const METADATA_SCAN_UNIT_PREFIX = "metadata-scan "; export const METADATA_SCAN_LEAF_ID = "metadata-scan"; // --- Live status (for /api/auto-queue/status), per kind -------------------- export type AutoRunnerInFlight = { videoId: string; leafId: string; channelSlug: string; startedAt: number; // A SENTENCE FOR A UNIT THAT IS NOT A VIDEO. `videoId` is the map key and // therefore always set, but the download lane's metadata-scan unit is // channel-scoped: its key is a synthetic `metadata-scan ` that names no // video directory. A console must not link it as one, so a unit carrying a // note renders the note instead of a video link. Absent for every ordinary // unit, which is how the surfaces stay unchanged. note?: string; // The registry job running this unit, when it is one (a download unit is a // child job on its platform queue); absent for a task on the runner's job. jobId?: string; }; // Why the runner is up but dispatching nothing. Every one of these was already // computed inside next() (or limit()) and thrown away, so the operator's most // common question — "it's running, why isn't it doing anything?" — had no answer // on the page. Recorded on the live record at each early return and cleared on a // successful pick. export type AutoRunnerIdleReason = // The job record is gone or no longer running. | "stopped" // policy.enabled went false; the runner is shutting itself down. | "disabled" // policy.snoozeUntil is in the future — idle on purpose, not stopped. | "snoozed" // The download lane's own gate (`autoQueue.download.held`). Download only. | "downloads-paused" // diskGate refused to start more work. Download only. | "disk-gate" // Every pending video belongs to a platform inside a rate-limit/network // cooldown window. Download only. | "cooldown" // Every pending video left after the platform gates was rate-limited // recently and is deferred (videoDeferrals). Download only. | "deferred" // Every pending platform is HELD: its rate limit outlasted the cooldown cap, // and the lane runs one probe per `pacing.holdProbeMinutes` until one comes // back clean (release 17, slice RL). Download only. | "held" // Every pending platform is inside the gap between two units: // `sleepBetweenDownloadsSeconds` plus the adaptive pace above its base // (release 17, slice RL). Download only. | "paced" // No enabled, non-degraded worker slot exists. Transcription only. | "no-workers" // The worker pool is pause-all'd. Transcription only, and DISTINCT from // 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" // The lane's own gate is shut (`autoQueue[lane].held`). // 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 // the per-platform one-download-at-a-time gate). | "capped"; type RunnerLive = { jobId: string; inFlight: Map; active: ActiveCounts; startedAt: number; idleReason: AutoRunnerIdleReason | null; }; type AutoRunnerSingleton = { runners: Map; }; declare global { // eslint-disable-next-line no-var var __yttAutoRunner__: AutoRunnerSingleton | undefined; } function getSingleton(): AutoRunnerSingleton { if (!globalThis.__yttAutoRunner__) { globalThis.__yttAutoRunner__ = { runners: new Map(), }; } return globalThis.__yttAutoRunner__; } // THE PERSISTED STATE IS ONE OBJECT FOR EVERY LANE'S RUNNER — and for // `recordDownloadBackoff` too. The holder and `sharedAutoQueueState` live in // `jobs/autoQueueState.ts` (release 8) because `jobs/downloadBackoff.ts` has to // write through the same object and `jobs/` may not import `controller/`. // Re-exported here so every existing import keeps working. export { sharedAutoQueueState }; export type AutoRunnerStatus = { kind: AutoQueueKind; running: boolean; jobId: string | null; startedAt: number | null; inFlight: AutoRunnerInFlight[]; activeByNode: ActiveCounts; // Why nothing is being dispatched right now, or null when the last scheduling // decision was a grant. A stopped runner always reads "stopped". idleReason: AutoRunnerIdleReason | null; }; export function getAutoRunnerStatus(kind: AutoQueueKind): AutoRunnerStatus { const live = getSingleton().runners.get(kind); const running = !!live && getRegistry().get(live.jobId)?.status === "running"; return { kind, running, jobId: live?.jobId ?? null, startedAt: live?.startedAt ?? null, inFlight: live ? [...live.inFlight.values()] : [], activeByNode: live ? { ...live.active } : {}, idleReason: running ? (live?.idleReason ?? null) : "stopped", }; } // --- Pending-work construction from snapshots ------------------------------ type ChannelMeta = { slug: string; platform: ReturnType; // The parsed config, carried so the media guard below does not re-read all 68 // config.json files on every tick of four lanes and on every three-second // status poll. listChannelConfigs has already parsed them. config: ChannelConfig; }; const SNAPSHOT_READ_CONCURRENCY = 64; // Re-derived from the shared listChannelConfigs rather than repeating the // readdir-then-serial-read here. // // PAUSED IS A FILTER ON THE CHANNEL LIST, not a shape in the tree, and this is // the one predicate that makes it so (plans/channel-priority.md, decision 3). A // tree cannot express exclusion — an `{type:"all"}` catch-all matches // everything, and first-match-wins would let a catch-all placed above the Low // group swallow Low's work — so a paused channel is removed from the LIST every // leaf draws from instead. This function is the single source of that list for // BOTH the runner loop and `computeLeafPending`, so a paused channel is absent // from the lane's draw and from the status panel's pending counts in one edit. // // PER LANE, through `isChannelPaused(model, slug, lane)` — the EFFECTIVE tier // for the lane being listed, so the per-operation override map decides: a // channel with `{tier:"normal", overrides:{sync:"paused"}}` (what all 15 live // `excludeFromSync` channels migrate to) is still drawn by the download lane, // and one with `{tier:"paused"}` or `overrides:{download:"paused"}` is not. // Re-evaluated on CHANNEL_LIST_TTL_MS, which is the clock for this decision. // // It returns the PRE-FILTER slug list beside the filtered meta, off the one // read. `known` is an EXISTENCE filter for the focus resolver (a site whose // membership has outrun the corpus, or a hand-edited settings.json), and the // paused-filtered list is the wrong answer for it twice over: it would drop a // focused channel that this lane happens to have paused, and it would make the // runner's "M channels held" line disagree with the banner, which resolves // against every channel (editor/app/operations/channelPriorityView.ts). async function listChannelMeta( paths: Paths, kind: AutoQueueKind, priority: ChannelPriority, // A listing the CALLER already has. The status poll reads every channel's // config once and then asked four lanes for their pending work, each of which // re-read the whole directory — see computeLeafPending's `shared`. sharedConfigs?: readonly { slug: string; config: ChannelConfig }[], ): Promise<{ meta: ChannelMeta[]; slugs: string[] }> { const configs = sharedConfigs ?? (await listChannelConfigs(paths)); return { meta: configs .filter(({ slug }) => !isChannelPaused(priority, slug, kind)) .map(({ slug, config }) => ({ slug, platform: detectPlatform(config.url), config, })), slugs: configs.map(({ slug }) => slug), }; } // --- The compiled priority trees ------------------------------------------- // // THE MODEL COMPILES, IT IS NOT CONSULTED. `settings.channelPriority` holds one // tier per channel plus one focus selector; `compileLaneRoot` turns that into // the lane's `AutoQueueGroup` (focus > normal > low > catch-all, all strict), // and dispatch runs the ordinary engine over it. Nothing in // `buildPendingByLeaf`, `selectNextWork`, `operationBatch`, `laneLimit` or // `pauseGates` learns a second priority mechanism, and "a zero limit is a hold, // never a stop" (controller/operationBatch.ts:22-25) is preserved trivially // because nothing here ever returns a limit. // // A FOCUS HOLDS THE REST BECAUSE STRICT DESCENT ALREADY DOES, and no second // mechanism is added: `pick()` (jobs/autoQueuePolicy.ts) filters a strict // group's children to those WITH WORK and descends into the first of them, and // it is re-asked on every grant. So focus work present => nothing below it is // picked; focus work exhausted => the next group runs; new focus work arriving // => the very next pick retakes the lane. // // WHERE THE HAND-EDITED TREES GO. The compiled root REPLACES // `settings.autoQueue[lane].root` at `laneDispatchRoot` below — the stored tree // is not consulted at all while a model exists. That is what makes the compiler // the ONE WRITER of channel priority: a `PolicyTreeEditor` save can still put a // channel leaf in the stored tree, but it cannot change what this lane // dispatches, so the two cannot fight — the compiler simply wins. (S3's // `saveChannelPriorityAction` then also PERSISTS the compiled roots, so the // stored tree and this one agree on disk; S4 makes the editor read-only for // compiled groups, which is the UI catching up with this fact.) // // AN ABSENT MODEL CHANGES NOTHING, BYTE FOR BYTE. `isDefaultChannelPriority` // below is the gate: no focus and no channel entries means the compiler never // runs and the stored trees stand exactly as they are today. const PRIORITY_CONTEXT_TTL_MS = 60_000; // The resolved half of the model — the half that costs I/O. Rebuilt when the // stored document changes or the TTL lapses, the same two triggers the // operation lanes' run context uses (settings key + 60 s), and for the same // reason: `resolveFocusSlugs` reads `transcripts/sites/*/site.json` through // `siteChannelIndex`, and a `{kind:"site"}` focus tracks that file's membership rather // than freezing a list. The TTL is what makes a channel added to the focused // site join the focus without a settings write. // The half of it `laneDispatchRoot` needs, and the only half a caller OUTSIDE // this file can supply: the editor's status payload resolves the same two // fields its own way (editor/app/operations/channelPriorityView.ts) and then // asks THIS function for the tree, so the console can never draw a tree the // runner does not dispatch from. export type PriorityDispatchContext = { model: ChannelPriority; focusSlugs: readonly string[]; }; type PriorityContext = PriorityDispatchContext & { key: string; at: number; focusSlugs: string[]; }; // The key carries `known` because the focus set is resolved against it: a // channel created or deleted since the last resolution changes the answer, and // a cache keyed on the document alone would keep the old one for a minute. function contextKey( model: ChannelPriority, paths: Paths, known: readonly string[], ): string { return JSON.stringify([model, paths.sitesDir, known]); } let priorityContext: PriorityContext | null = null; // `settings` is a PARAMETER, not a read: both call sites already hold the // settings object for this tick (the runner's `next()` reads it to see // `enabled` and `snoozeUntil`; `computeLeafPending` reads it for the policy), // and a second `getSettings()` here made every tick and every three-second // status poll parse settings.json twice. // Exported for its own test: the cache key, the TTL and the "only a site focus // reads the sites directory" rule are three claims a caller cannot observe // through `laneDispatchRoot`, which is pure. export function resetPriorityContextForTest(): void { priorityContext = null; } export function priorityContextFor( paths: Paths, settings: SiteSettings, // Every channel slug that exists — the existence filter for the focus, the // same one the editor's status payload passes. Omitted only by a caller that // genuinely has no list. known: readonly string[] = [], ): PriorityContext { const model = settings.channelPriority; // The paths go in the key so two worktrees' runners in one process cannot // share a focus resolved against the other's sites directory. const key = contextKey(model, paths, known); const now = Date.now(); if ( priorityContext && priorityContext.key === key && now - priorityContext.at < PRIORITY_CONTEXT_TTL_MS ) { return priorityContext; } // ONLY A SITE FOCUS READS THE SITES DIRECTORY. `{kind:"channels"}` and // `{kind:"none"}` resolve from the document alone, so the common case pays // nothing. `siteChannelIndex` (lib/site.ts) is the ONE spelling of that read // — the sync tick, the /channels writer and the status payload ask it too. const index: SiteChannelIndex = model.focus.kind === "site" ? siteChannelIndex(paths) : {}; priorityContext = { key, at: now, model, focusSlugs: resolveFocusSlugs( model, index, known.length > 0 ? known : undefined, ), }; return priorityContext; } // THE ROOT THIS LANE ACTUALLY DISPATCHES FROM. One function, called by the // runner loop AND by computeLeafPending, so the status panel can never name a // leaf the runner does not have. // // `slugs` is the lane's own (already paused-filtered) channel list; // `compileLaneRoot` re-applies the same per-lane predicate, so handing it the // filtered list and handing it every slug produce the identical tree for THIS // lane. Compiling is pure and O(channels) — ~69 string pushes against the // ~6.5 MB of snapshot JSON the same tick folds — so it happens per tick and // only the focus resolution above is cached. export function laneDispatchRoot( kind: AutoQueueKind, policy: AutoQueuePolicy, ctx: PriorityDispatchContext, slugs: readonly string[], ): AutoQueueGroup { if (isDefaultChannelPriority(ctx.model)) return policy.root; return compileLaneRoot(kind, ctx.model, slugs, ctx.focusSlugs); } // The once-per-state-change line a lane writes while a focus is holding it. // // NOT AN IDLE REASON, and deliberately not: a lane whose focus group holds the // rest is not idle, it is dispatching focus work — `AutoRunnerIdleReason` stays // exactly as it is, and `no-pending` remains the true answer when the whole // tree is empty. This is the runner's LOG saying which of three states it is // in, so "why is only jeralyzer moving?" is answerable from the job log alone. // // The state is the three-valued thing, NOT the counts: the counts are in the // message but never in the key, or every completed focus unit would re-fire the // line. Same shape as the snooze line and the disk-gate line above. export function focusHoldState(summary: FocusSummary | null): string { if (!summary || !summary.active) return "none"; return summary.holding ? "hold" : "free"; } export function focusHoldLine( kind: AutoQueueKind, summary: FocusSummary | null, ): string | null { if (!summary || !summary.active) return null; const channels = `${summary.channelCount} channel${summary.channelCount === 1 ? "" : "s"}`; if (summary.holding) { return ( `Auto-${kind}: focus (${channels}) is holding this lane — ` + `${summary.focusPending} focus unit(s) pending, ` + `${summary.heldChannels} channel(s) held.` ); } return ( `Auto-${kind}: focus (${channels}) has no work left in this lane — ` + `${summary.otherPending} unit(s) released to the rest of the corpus.` ); } // The gate itself, as a closure so the runner keeps one line per transition and // the test can drive the transitions without a corpus. Logs on entering "hold" // and on entering "free"; says nothing while there is no active focus. export function makeFocusHoldReporter( kind: AutoQueueKind, onLog: (line: string) => void, ): (summary: FocusSummary | null) => void { let state = "none"; return (summary) => { const next = focusHoldState(summary); if (next === state) return; state = next; const line = focusHoldLine(kind, summary); if (line) onLog(line); }; } // ONCE PER STATE CHANGE, not once per tick. buildChannelWork runs on every // scheduling tick of all four lanes and on a three-second status poll; logging a // skip each time would write the same line thousands of times an hour and bury // the one that matters. Keyed by slug, so mounting the drive logs the recovery // too — an operator watching the log sees the channel leave and come back. const mediaSkipLogged = new Map(); // WHETHER A LANE SKIPS A CHANNEL, from its inspect answer (release 17). The // digest lane writes only text, so only the text hold — a `legacy` channel, a // `data/` it cannot read, a tier migration in flight — keeps it off a channel; // it runs while the media is moving, stalled or unmounted. The transcription, // download and backfill lanes open the big files and keep the media hold. export function isChannelHeldForLane( kind: AutoQueueKind, location: Pick, ): boolean { if (kind === "digest") { return isTextHeld(location.status) || !location.text.readable; } return isMediaHeld(location.status); } function noteSkippedForMedia( slug: string, status: ChannelMediaStatus, detail?: string, ): void { if (mediaSkipLogged.get(slug) === status) return; mediaSkipLogged.set(slug, status); // The hold's own words first (lib/channelMediaHold.ts — "held: its media is // moving …" for a channel with a relocation marker), then the status and the // inspector's detail as they always were. console.log( `[auto] skipping ${slug}, ${mediaHoldText(status) ?? "held"}: ` + `media ${status}${detail ? ` — ${detail}` : ""}`, ); } function noteMediaReachable(slug: string): void { if (!mediaSkipLogged.has(slug)) return; mediaSkipLogged.delete(slug); console.log(`[auto] ${slug}: media reachable again — the hold is lifted`); } // For tests and for /api/test/invalidate-cache: the log-once memory is // process-local state, not a decision, and nothing downstream reads it. export function resetChannelMediaSkipLog(): void { mediaSkipLogged.clear(); } // Read each channel's snapshot and project the buckets this runner kind cares // about into ChannelWork, plus a videoId -> owning channel map (a platform/all // leaf spans channels, so the pick needs the owner to locate the video dir). async function buildChannelWork( paths: Paths, kind: AutoQueueKind, meta: ReadonlyArray, // The operation id spaces to project alongside the buckets: exactly the ones // THIS LANE DISPATCHES, never the whole catalog. See laneOperationIds. operations: ReadonlyArray = [], ): Promise<{ channels: ChannelWork[]; owner: Map }> { // 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() }; } const channels: ChannelWork[] = []; const owner = new Map(); // Read every snapshot concurrently, then fold sequentially — `owner` is // first-writer-wins, so the fold must stay in `meta` order to stay // 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) => readChannelSnapshotShared(paths, m.slug), ); // GUARD 2 OF FOUR (see plans/relocate-channel-media.md). This is the runners' // ONLY channel-level chokepoint, and it is the one that matters most: the // download lane reading an unmounted channel's snapshot as "everything // undownloaded" is an instruction to re-fetch the whole channel onto the disk // that was too full to hold it. Three syscalls per channel per tick, run // concurrently alongside the snapshot reads. const media = await mapConcurrent(meta, SNAPSHOT_READ_CONCURRENCY, (m) => // The config is passed, not re-read: without it this is a 69th, 70th … // config.json read per tick for a field listChannelConfigs already parsed. inspectChannelMedia(paths, m.slug, m.config), ); for (const [i, { slug, platform }] of meta.entries()) { const snap = snaps[i]; if (!snap) continue; // A SKIP IS A SKIP. The lane keeps running every other channel: this is // never a lane stop, and it is not the LANE's hold — "a zero limit is a // hold, never a stop" (pauseGates.ts) is a different mechanism and is // untouched by it. It is the CHANNEL's media hold (lib/channelMediaHold.ts, // the words both builds use): a relocation marker, an unmounted or stalled // drive, a link and a config that disagree. It lifts by itself — the next // tick after the marker is removed (a move completed or abandoned; the // movers forget the memo) or the drive is back reads the channel again. // // THE DIGEST LANE READS TEXT (release 17): it is held only by the text // hold — a `legacy` channel, or a `data/` it cannot read — so a digest // runs on a channel whose media is moving, stalled or unmounted. The // transcription, download and backfill lanes open the big files and keep // the media hold. const location = media[i]; if (location && isChannelHeldForLane(kind, location)) { noteSkippedForMedia(slug, location.status, location.detail); continue; } noteMediaReachable(slug); const buckets: Record = {}; // Project every bucket a leaf could be pointed at — including the opt-in // auto-caption ones. A leaf that names a bucket the runner never projected // would silently find no work; projecting them here costs nothing when no // leaf (and no policy switch) asks for them, because buildPendingByLeaf // only walks the buckets its `defaultBuckets` / `match.bucket` name. for (const name of selectableBucketsForKind(kind)) { const ids = bucketIdsFrom(snap, name); buckets[name] = ids; for (const id of ids) if (!owner.has(id)) owner.set(id, slug); } // THE LANE'S OWN WORK LIST, under the id of the operation it dispatches — // slice 1.5. A bucket-less leaf draws THIS rather than walking the default // buckets itself (see defaultDrawsForPolicy), so all four lanes now answer // "what does this lane have to do" out of `snapshot.backfill[op].ids`. // // THE FALLBACK IS THE MIGRATION. No live snapshot is regenerated by the // slice that added the entry, so a channel reads its work list through // bucketLaneWorkIds until its next regen — the same fold // generateChannelSnapshot writes, so the two cannot disagree. // // Projected into `buckets`, NOT into `operations`, and that is the whole // reason nothing moves: the bucket claim space is `\0id`, so a video an // explicit `failedListed` leaf already claimed is not claimed a second time // by a catch-all leaf drawing this list. The operation claim space exists to // keep diarization and attribution apart on ONE video, and a bucket lane has // no second operation to keep apart from. const laneWorkList = bucketLaneOperationId(kind); if (laneWorkList) { const ids = snap.backfill?.[laneWorkList]?.ids ?? bucketLaneWorkIds(kind, snap); buckets[laneWorkList] = ids; for (const id of ids) if (!owner.has(id)) owner.set(id, slug); } let ops: Record | undefined; if (operations.length > 0) { ops = {}; for (const op of operations) { // `ids` IS the reachable set — missing + stale + partial, never // missing-input, deferred or blocked (see OperationSnapshotEntry). A // leaf pointed at an operation therefore claims only work the lane can // actually do, which is the same contract a bucket carries. const ids = snap.backfill?.[op]?.ids ?? []; ops[op] = ids; for (const id of ids) if (!owner.has(id)) owner.set(id, slug); } } // THE CHANNEL-SCOPED HALF, off the snapshot already read. The download // lane's pre-pick uses these two to decide whether to run a metadata scan // for this channel before it downloads anything from it; every other lane // ignores them. Both come out of `snap`/`config` with no extra I/O — the // scan backlog is a number the snapshot already publishes, and the filter // is the config listChannelConfigs already parsed. const config = meta[i].config; channels.push({ slug, platform, buckets, operations: ops, scanUnscanned: snap.metadataScan?.unscanned ?? 0, filtered: compileDownloadFilter(config.downloadFilter) !== null, }); } return { channels, owner }; } // WHAT THIS LANE MAY PROJECT: the operations the lane actually dispatches, and // nothing else. // // 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. // // 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[] { 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 // existed. One place, so the runner and the status page can never disagree about // what an absent field means. export function orderOf(policy: Pick): AutoQueueOrder { return policy.order ?? "listed"; } // Build the comparator buildPendingByLeaf sorts each leaf with, plus the keys it // was built from (the UI renders them next to the drill-down ids, which is how an // ordering change gets verified by eye). Returns a null comparator for // order:"listed" so the historical path is "don't sort", not "sort by identity". // // Never throws: a recency lookup failing is not a reason to stop dispatching // work, so a broken index degrades to today's ordering. async function recencyOrdering( kind: AutoQueueKind, paths: Paths, policy: AutoQueuePolicy, meta: ReadonlyArray, channels: ReadonlyArray, ): Promise<{ compare: ((a: string, b: string) => number) | null; keys: Map; }> { 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. const candidateIds = new Set(); const owner = new Map(); for (const ch of channels) { for (const ids of [ ...Object.values(ch.buckets), ...Object.values(ch.operations ?? {}), ]) { for (const id of ids) { candidateIds.add(id); if (!owner.has(id)) owner.set(id, ch.slug); } } } try { const keys = await buildRecencyKeys({ paths, meta, candidateIds, owner, // Only auto-download faces videos with NOTHING on disk to date them by. // Every transcription candidate is already downloaded, so its // metadata.info.json is there for the tail read and the playlist reads // would be pure cost. interpolate: kind === "download", }); 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 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, "newest"), }), 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, candidateIds: ReadonlySet, ): Promise> { const ids = [...candidateIds]; const out = new Map(); 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. const PENDING_HEAD = 20; // The video the policy would hand out next, and enough context to say WHY that // one. `skippedLeafIds` are the rules ahead of the chosen one in priority order // that had nothing pending — which is the whole answer to "why is it working on // rule 2?". export type NextUpView = { videoId: string; channelSlug: string | null; leafId: string; path: string[]; recency: RecencyKey | null; skippedLeafIds: string[]; }; export type LeafPending = { counts: Record; // The first PENDING_HEAD ids of each leaf, in the exact order the runner would // hand them out. head: Record; // videoId -> owning channel slug, for linking a drill-down id to its page. owner: Record; // videoId -> recency sort key, only for the ids in `head`, and only when the // policy actually orders by recency. recency: Record; // What the policy would pick right now, or null when it would pick nothing. nextUp: NextUpView | null; }; // Per-leaf pending (snapshot-derived) work for the status panel — what the // operator sees as "cornbreadman: 12 pending". Uses the same matching AND the // same ordering as the runner, so the numbers and the drill-down line up with // what would actually be picked. // // `shared` IS THE FOUR-LANE POLL'S WAY OUT OF READING EVERYTHING FOUR TIMES. // /operations asks all four lanes on a ~3 s poll, and each call listed every // channel's config (a readdir plus a read per channel — 30 on Jeralyzer) and // parsed the auto-queue state document again, having already read both itself // for the rest of the payload. Nothing about either read is per-lane, so a // caller that holds them passes them in and the poll pays once. Omitted, each // is read here exactly as before, so every other caller is unchanged. export async function computeLeafPending( kind: AutoQueueKind, paths: Paths = getPaths(), shared?: { configs?: readonly { slug: string; config: ChannelConfig }[]; state?: AutoQueueState; }, ): Promise { const settings = getSettings(); const policy = settings.autoQueue[kind]; // THE SAME TWO PRIORITY DECISIONS THE RUNNER MAKES, in the same order: the // paused filter on the channel list, then the compiled root. This function's // whole contract is that its numbers are the runner's numbers, so both sides // of channel priority have to be here too — otherwise the panel would count // pending work for a paused channel, or attribute it to a stored leaf the // runner is not dispatching from. const { meta, slugs } = await listChannelMeta( paths, kind, settings.channelPriority, shared?.configs, ); const ctx = priorityContextFor(paths, settings, slugs); const root = laneDispatchRoot(kind, policy, ctx, meta.map((m) => m.slug)); const laneOperations = laneOperationIds(kind); const { channels, owner } = await buildChannelWork( paths, kind, meta, laneOperations, ); const { compare, keys } = await recencyOrdering( kind, paths, policy, meta, channels, ); // THE LANE'S OWN PROJECTION IS THE FILTER. `channels` above carries only the // buckets and operations THIS lane dispatches (see laneOperationIds), so a // leaf naming anything else finds no list and comes back empty — the zero // retainLeaves used to apply afterwards, reached by construction. const pending = buildPendingByLeaf( root, channels, defaultDrawsForPolicy(kind, policy, bucketLaneOperationId(kind)), { ...(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. const live = getSingleton().runners.get(kind); if (live && live.inFlight.size > 0) { removeIds(pending, new Set(live.inFlight.keys())); } const counts: Record = {}; const head: Record = {}; const ownerOut: Record = {}; const recency: Record = {}; for (const [leafId, ids] of Object.entries(pending)) { counts[leafId] = ids.length; const slice = ids.slice(0, PENDING_HEAD); head[leafId] = slice; for (const id of slice) { const slug = owner.get(id); if (slug) ownerOut[id] = slug; const key = keys.get(id); if (key) recency[id] = key; } } // THE TRAP: selectNextWork advances SWRR fairness by mutating // runtime.currentWeights in place. This runs on a 3-second status poll, so // asking the live runtime would let merely HAVING the page open skew a // round-robin group's rotation. Deep-clone first; the clone is discarded. // The live shared object when a runner holds one (only cloned from, below), // else the file. const state = shared?.state ?? (await liveAutoQueueState(paths)) ?? (await readAutoQueueState(paths)); const runtime = { currentWeights: { ...state[kind].runtime.currentWeights }, }; const pick = selectNextWork( root, pending, runtime, live ? { ...live.active } : {}, ); let nextUp: NextUpView | null = null; if (pick) { const order = flattenLeaves(root); const at = order.findIndex((l) => l.id === pick.leafId); nextUp = { videoId: pick.videoId, channelSlug: owner.get(pick.videoId) ?? null, leafId: pick.leafId, path: pick.path, recency: keys.get(pick.videoId) ?? null, skippedLeafIds: at <= 0 ? [] : order .slice(0, at) .filter((l) => (counts[l.id] ?? 0) === 0) .map((l) => l.id), }; } return { counts, head, owner: ownerOut, recency, nextUp }; } // Back-compat shim for callers that only ever wanted the counts. export async function computeLeafPendingCounts( kind: AutoQueueKind, paths: Paths = getPaths(), ): Promise> { return (await computeLeafPending(kind, paths)).counts; } function countPending(pending: Record): number { let total = 0; for (const ids of Object.values(pending)) total += ids.length; return total; } function removeIds( pending: Record, exclude: ReadonlySet, ): void { if (exclude.size === 0) return; for (const leafId of Object.keys(pending)) { pending[leafId] = pending[leafId].filter((id) => !exclude.has(id)); } } // Count of worker slots that could run a transcription right now (enabled and // not degraded — busy ones still count as capacity, just occupied). This is the // ceiling on outstanding units so the loop never parks a backlog. function eligibleSlots(): number { return getWorkerPool() .summary() .filter((s) => s.state === "enabled" && !s.degraded).length; } // --- The fill-to-capacity loop --------------------------------------------- // One unit of work selected by the policy, carried from the runPool `next` // (selection + reservation) into `run` (execution + release). // WHICH CHANNEL GETS A METADATA SCAN THIS TICK, or none — the download lane's // pre-pick, as a pure function over the projection the runner already built. // // Pure because the rule is the whole feature and everything else about it is // plumbing. Three conditions, and each one is a bug if it goes: // // - A FILTER. Without one a scan settles nothing: it would fetch a title per // listed video, change no verdict, and spend the source's patience to do // it. `filtered` is "the channel's downloadFilter compiles", which is the // same predicate `settledByTitleFilterIds` opens the store for. // - A BACKLOG. `scanUnscanned` is the snapshot's own `metadataScan.unscanned`, // defined to mean exactly what `metadataScanTargets()` will fetch, so a // channel offered here is a channel the scan has work for. // - ONCE PER RUNNER. A channel whose scan left the backlog above zero — a // soft block, a cooldown, a batch that stopped early — is not retried by // this runner. Re-offering it on the next three-second tick would hammer // the source that just refused us; the operator's Run and the next restart // both still work. // // The platform gate is the same one the downloads honour: a busy or cooling // platform offers no scan either, because a scan is one more thing asking that // source for titles. // // A FOCUS HOLDS THE SCAN TOO. The compiled priority tree makes a focus hold // every download pick by strict descent (see laneDispatchRoot), and the pre-pick // does not go through that tree — so without this, "only jeralyzer is moving" // would be false the moment another channel had titles to read, on the same // network the focus is trying to have to itself. The rule mirrors what strict // descent does rather than re-implementing it: while the focus still has // pending work, only a focused channel is offered a scan; once it is exhausted // the lane is free and so is the scan. // // FIRST IN `channels` ORDER, which is `metaCache` order, which is the priority // order every other pick on this lane already uses. export function pickMetadataScanChannel( channels: ReadonlyArray, opts: { scanned: ReadonlySet; platformSkip: ReadonlySet; platformOf: (slug: string) => string; // The resolved focus, when one is holding the lane. Absent (or with // `holding: false`) means every channel is eligible, which is the default // and the byte-identical path for a corpus with no focus configured. focus?: { slugs: ReadonlySet; holding: boolean }; }, ): { slug: string; targets: number } | null { const held = opts.focus?.holding === true; for (const c of channels) { const targets = c.scanUnscanned ?? 0; if (!c.filtered || targets <= 0) continue; if (held && !opts.focus!.slugs.has(c.slug)) continue; if (opts.scanned.has(c.slug)) continue; if (opts.platformSkip.has(opts.platformOf(c.slug))) continue; return { slug: c.slug, targets }; } return null; } type Picked = { pick: WorkPick; channelSlug: string; unitPlatform: string | null; // THE CHANNEL-SCOPED UNIT. Set only by the download lane's metadata-scan // pre-pick: `pick.videoId` is then a synthetic key naming no video and // `pick.path` is empty, so no node's active count moves for it. `run()` // branches on this and nothing else. scan?: { slug: string; targets: number }; }; async function runLoop( kind: AutoQueueKind, paths: Paths, live: RunnerLive, onLog: (line: string) => void, setProgress: (snap: JobProgress) => void, signal: AbortSignal, ctx: JobRunContext, ): Promise { const tracker = makeTaskTracker(ctx, onLog); // SHARED with every other lane's runner in this process — see // sharedAutoQueueState. A per-runner copy makes each persist erase the other // lanes' pick logs. const state = await sharedAutoQueueState(paths); const kindState = state[kind]; const runtime = kindState.runtime; // 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. // A HELD platform's entry is kept however old: a hold ends with a clean // probe, not with time (release 17, slice RL). prunePlatformPacing(kindState, Date.now()); pruneSubtitleDeferrals(kindState.subtitleDeferrals, Date.now()); // And lapsed per-video deferrals (live ones survive a restart by design — a // restart must not re-hit the rate-limited video at fails+1). pruneDeferred(kindState.videoDeferrals, 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(); const completedOrder: string[] = []; const counters = { transcribed: 0, failed: 0, skipped: 0 }; 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, root: AutoQueueGroup, laneOperations: ReadonlyArray, ): 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 // download for a platform that's already busy. Keyed by platform name (or // "unknown" for unrecognized hosts). Unused for transcription. const platformInFlight = new Map(); const PER_PLATFORM_CAP = 1; // THE GAP BETWEEN TWO UNITS ON ONE PLATFORM (release 17, slice RL): the epoch // ms before which the platform takes no new unit. Set when a unit settles to // `sleepBetweenDownloadsSeconds` plus the adaptive pace above its base // (downloadGapMs) — the lane used to start the next unit the moment one // settled, ignoring the setting every batch download honours. In memory: a // restart is itself a gap. const platformNextStartAt = new Map(); // Channels this runner has already scanned (or tried to). See the pre-pick in // next() for why a scan is never retried inside one runner's lifetime. const scannedThisRun = new Set(); // Platforms this tick may not touch — busy or cooling down. Hoisted out of the // download branch because the metadata-scan pre-pick honours it too: a // platform in a 429 cooldown must not be asked for 1,800 titles either. const platformSkip = new Set(); const platformKey = (slug: string, slugToPlatform: Map) => slugToPlatform.get(slug) ?? "unknown"; // videoId -> child download job id, for the units the runner currently has on // the registry's platform queues. Used to hard-cancel them (kill yt-dlp or // drop a still-queued unit) when the runner itself is cancelled. Download-only. const childJobIds = new Map(); let metaCache: ChannelMeta[] = []; // Every channel slug, pre-paused-filter, on the same TTL as metaCache. let knownSlugs: string[] = []; let metaAt = 0; // The root the last pick was made from. next() sets it every tick; // runOperationPick reads it to resolve the leaf its own pick named. let dispatchRoot: AutoQueueGroup = getSettings().autoQueue[kind].root; // Whether the last next() saw the disk gate closed. next() runs on every // scheduling tick, so without this the log fills with one identical line per // tick for as long as the disk is full — which is precisely the situation in // which the log needs to stay readable. Logged on each transition instead. let diskIdle = false; const reportFocusHold = makeFocusHoldReporter(kind, onLog); // A graceful stop (the job record is gone after an e2e reset, or the policy was // disabled) is modeled as a soft drain: stop picking, let in-flight finish. // Combined with the job's real drain signal so either one ends the run. const stopController = new AbortController(); const drainSignal = AbortSignal.any([ctx.drainSignal, stopController.signal]); // NOTE: do NOT call getWorkerPool().reconfigure() here. At server boot the // runner is the first thing to touch the pool, and reconfigure() sets the // pool's `initialized` flag WITHOUT applying the saved "Set as default" // arrangement (.worker-defaults.json) — so a pre-empted ensureInit() would // never disable the workers the operator turned off, and the runner would use // them anyway. Leaving it to the pool's lazy ensureInit() (triggered by the // first eligibleSlots()/acquire() below) seeds settings AND applies defaults // exactly once. Live settings changes are pushed by saveSettingsAction's own // reconfigure(applyEnabled:true), so nothing is lost by not refreshing here. // On a HARD cancel, actively cancel each in-flight/queued child download job: a // running one's yt-dlp is SIGTERM'd (via downloadOneManaged's signal), a // still-queued one is dropped — either way launchUnit's `await res.done` // resolves and the unit's run() settles so runPool's in-flight clears. (Drain // is soft: it leaves them to finish. Transcription units settle via the hard // signal threaded into the unit, so they need no explicit cancel here.) const onHardAbort = () => { for (const jid of childJobIds.values()) getRegistry().cancel(jid); }; 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 }; // THE DUPLICATE PLAN IS BUILT ONCE PER RUNNER, not once per refresh. It reads // the whole corpus-wide duplicates report, and this context is rebuilt on a // 60-second TTL — so re-deriving it would put that read on a timer forever to // learn about a cluster detection has not run since. A per-channel job builds // one per job for the same reason: it is per-RUN state, and a runner's run is // its whole lifetime. let clusterPlan: NonNullable["clusterPlan"] = null; const LANE_RUN_TTL_MS = 60_000; const laneRunKey = (settings: ReturnType): string => JSON.stringify([ settings.digest, settings.backfill, settings.diarization, settings.attribution, ]); const refreshLaneRun = async (): Promise => { if (!isOperationLane(kind) || laneRun.refreshing) return; const settings = getSettings(); const key = laneRunKey(settings); // The TTL covers a FAILED probe too, or a dead ollama would be re-asked on // every three-second poll. The log line already says the lane is idle and // why; the cost of the TTL is up to a minute of idling after the engine // comes back. const fresh = laneRun.key === key && Date.now() - laneRun.at < LANE_RUN_TTL_MS; if (fresh && (laneRun.run !== null || laneRun.error !== null)) return; // NEVER SWAP THE CONTEXT OUT FROM UNDER AN IN-FLIGHT UNIT. Its llm and // unit-executor leases decrement counters on the object it was dispatched // with, so replacing that object mid-flight would leak a slot from the // limit's point of view — permanently, since nothing ever decrements it // again. if (laneRun.run && live.inFlight.size > 0) return; laneRun.refreshing = true; try { const opened = await openOperationRun({ lane: kind, paths, onLog, ...(clusterPlan ? { clusterPlan } : {}), }); clusterPlan = opened.digest?.clusterPlan ?? clusterPlan; // THE WHOLE LEDGER SURVIVES THE REFRESH — the OBJECT, not a copy of two // of its fields. // // Every field on it is per-RUN, and a runner's run is its whole lifetime: // the metered spend cap is a per-run ceiling (resetting it every sixty // seconds would turn a $5 cap into $5 a minute) and the disk-floor latch // is the same shape. The two COUNTERS are why this has to be the object. // The idle check above cannot close the window: it runs before the await // below, so `limit()` keeps reading the old run while this one is being // opened and `next()` may dispatch one more unit against it. That unit's // `llmActive--` / `unitActive--` and its own `costUsd` land on the run it // was dispatched with — and if that run's ledger is not this one's, the // lane under-counts its leases forever and drops that unit's spend from // the cap. if (laneRun.run) opened.live = laneRun.run.live; // 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 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 ceiling; } const slots = eligibleSlots(); // A zero limit means runPool never calls next(), so this is the only place // "every worker is off or degraded" can be observed — without it the page // would show a running runner with no explanation at all. // // THE PAUSE IS ASKED FIRST, and it has to be. pauseAll() DISABLES every // worker (remembering each one's state for the resume), so a paused pool is // indistinguishable from an empty one by slot count alone — and // /operations/transcription said "no enabled worker to run it" next to its // own Resume Transcriptions button. Same shape as the download lane's // downloads-paused: an idle reason, not a stop. if (slots === 0 && live.inFlight.size === 0) { live.idleReason = getWorkerPool().isPaused() ? "workers-paused" : "no-workers"; } return Math.min(slots, ceiling); }; // Select the next unit and RESERVE its slot (so the next selection sees it), // or return null when there's nothing to pick right now / the runner should // stop. runPool calls this only while below the limit and not draining. const next = async (): Promise => { // Stop if the job record is gone (e2e registry reset) or no longer running // (a hard cancel flips it to "cancelled" and also fires `signal`). const rec = getRegistry().get(live.jobId); if (!rec || rec.status !== "running") { live.idleReason = "stopped"; stopController.abort(); return null; } const settings = getSettings(); const policy: AutoQueuePolicy = settings.autoQueue[kind]; if (!policy.enabled) { live.idleReason = "disabled"; onLog(`Auto-${kind} disabled — stopping runner.`); stopController.abort(); return null; } // Snooze: idle, do NOT stop. Same shape as the downloads-pause branch below // — settings are re-read every iteration, so the runner wakes by itself the // moment the deadline passes, and because the deadline lives in settings.json // it survives a server restart. getSettings() already normalizes a lapsed // snooze to null; the clock check is belt-and-braces for a clock that moved. const snoozeUntil = policy.snoozeUntil ?? null; if (snoozeUntil !== null && snoozeUntil > Date.now()) { if (live.idleReason !== "snoozed") { onLog( `Auto-${kind} snoozed until ${new Date(snoozeUntil).toLocaleString()}.`, ); } live.idleReason = "snoozed"; return null; } // Global downloads pause: like the `enabled` flag, this is re-read each // iteration. Rather than stop the runner, idle it (return null) so it // resumes dispatching on the next tick once unpaused — no restart needed. // The gate itself is defined once, in lib/pauseGates.ts. if (kind === "download" && isGateHeld(settings, "download")) { live.idleReason = "downloads-paused"; 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 // preflight, the thing that runs by itself did not. // // Idle (return null), don't stop: the gate is a refusal to START more work, // and it self-heals, so the runner must still be here to notice. Download // only — transcription writes a transcript.json next to audio it already // has, so stopping it frees nothing. if (kind === "download") { // The CORPUS volume, deliberately, and not a channel's. This check runs // before the pick, so there is no channel yet — and the lane spans 68 of // them, on however many volumes. The per-channel answer is taken further // down, in runYtdlp, where the slug is known and the bytes are about to // be written; a full SSD idling the lane here is the conservative half of // the pair, not the whole of it. const gate = await diskGate(paths, settings); if (!gate.ok) { if (!diskIdle) { diskIdle = true; onLog(`Auto-download idle: ${gate.message}.`); } live.idleReason = "disk-gate"; return null; } if (diskIdle) { diskIdle = false; onLog( `Auto-download resuming: ${formatBytes(gate.freeBytes)} free, ` + `above the ${formatBytes(gate.resumeBytes)} resume mark.`, ); } } // Refresh the (rarely-changing) channel list/platforms on a TTL. // // THE PAUSED FILTER RIDES THIS CLOCK. `listChannelMeta` drops every channel // whose effective tier for THIS lane is `paused`, so the 30 s TTL is also // how long a pause takes to reach dispatch. `metaAt === 0` rather than // `metaCache.length === 0` is the cache-miss test now: a lane on which // every channel is paused has a legitimately empty list, and the old // sentinel would re-read 68 configs on every three-second tick for it. // // THE META REFRESH COMES FIRST, because the focus resolution needs the // channel list it produces: `known` is what drops a focus slug the corpus // no longer has, and resolving without it makes this lane's "M channels // held" line disagree with the banner by one. const now = Date.now(); if (metaAt === 0 || now - metaAt > CHANNEL_LIST_TTL_MS) { const listed = await listChannelMeta( paths, kind, settings.channelPriority, ); metaCache = listed.meta; knownSlugs = listed.slugs; metaAt = now; } const ctx = priorityContextFor(paths, settings, knownSlugs); // THE COMPILED ROOT ENTERS HERE, and this is the only place it does for the // dispatch path: everything below — the projection, the completed filter, // the pick and the leaf lookup in runOperationPick — reads `dispatchRoot`, // never `policy.root`. See laneDispatchRoot. const root = laneDispatchRoot(kind, policy, ctx, metaCache.map((m) => m.slug)); dispatchRoot = root; const slugToPlatform = new Map( metaCache.map((m) => [m.slug, m.platform ?? "unknown"]), ); const laneOperations = laneOperationIds(kind); const { channels, owner } = await buildChannelWork( paths, kind, metaCache, laneOperations, ); // Re-derived each iteration from the freshly-read policy, like `enabled`, so // toggling the replace-auto-captions lane (or the ordering) takes effect // without a restart. The recency index behind the comparator has its own // 30s TTL, so this is a map lookup per candidate on all but one tick in ten. const { compare } = await recencyOrdering( kind, paths, policy, metaCache, channels, ); // THE LANE'S OWN PROJECTION IS THE FILTER — see computeLeafPending. This is // the dispatch path, so it is the one where getting it wrong runs the wrong // engine on the wrong video. const pending = buildPendingByLeaf( root, channels, defaultDrawsForPolicy(kind, policy, bucketLaneOperationId(kind)), { ...(compare ? { compare } : {}), defaultOperations: laneOperations }, ); // 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, root, laneOperations); const pendingBeforeGates = countPending(pending); // Say — ONCE per transition — whether a focus is holding this lane. Read off // the `prio-*` leaf ids in the map just built, so it costs one pass over // keys and no new read, and it is skipped entirely while no focus resolves. // Computed ONCE and reused by the scan pre-pick below, which has to honour // the same hold: `holding` is `focusPending > 0`, i.e. exactly "strict // descent is still inside the focus group". const focus = ctx.focusSlugs.length > 0 ? focusSummary(ctx.model, ctx.focusSlugs, pending) : null; reportFocusHold(focus); // 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. let anyCooling = false; // Download only: some pending video was dropped because it is deferred // after a recent rate limit (see unitOutcome.ts). let anyDeferred = false; // Download only: a platform skipped for its hold (one probe at a time) or // for the gap between units — each its own answer to "why idle?". let anyHeld = false; let anyPaced = false; platformSkip.clear(); if (kind === "download") { // Merge in any cooldown a manual sync/import wrote to the shared state // (read-modify-write from outside the runner) since our last persist. // `until` and `fails` each merge to their max, so a sync-set 429 cooldown // is honored even though our own next persist would otherwise clobber it, // and neither side's escalation count is lost. try { const persisted = (await readAutoQueueState(paths)).download .platformBackoff; for (const [pf, e] of Object.entries(persisted)) { kindState.platformBackoff[pf] = mergeBackoffEntry( kindState.platformBackoff[pf], e, ); } } catch { // Best-effort: a transient read failure just skips this iteration's merge. } for (const [pf, n] of platformInFlight) { if (n >= PER_PLATFORM_CAP) platformSkip.add(pf); } for (const pf of Object.keys(kindState.platformBackoff)) { if (isCoolingDown(kindState.platformBackoff, pf, now)) { platformSkip.add(pf); // A platform skipped for a cooldown is a different answer to "why is // it idle?" than one skipped for being busy, so the two are tracked // apart rather than both reading as "capped". A held platform's // cooldown is its next probe. if (isPlatformHeld(kindState.platformHolds, pf)) anyHeld = true; else anyCooling = true; } } for (const [pf, at] of platformNextStartAt) { if (at > now && !platformSkip.has(pf)) { platformSkip.add(pf); anyPaced = true; } } if (platformSkip.size > 0) { for (const leafId of Object.keys(pending)) { pending[leafId] = pending[leafId].filter((id) => { const slug = owner.get(id); return !(slug && platformSkip.has(platformKey(slug, slugToPlatform))); }); } } // A video rate-limited recently is skipped until its deferral lapses, so // after a cooldown the runner moves on rather than re-picking it. if (Object.keys(kindState.videoDeferrals).length > 0) { for (const leafId of Object.keys(pending)) { pending[leafId] = pending[leafId].filter((id) => { if (!isVideoDeferred(kindState.videoDeferrals, id, now)) return true; anyDeferred = true; return false; }); } } } // ── THE SCAN COMES BEFORE THE DOWNLOADS ──────────────────────────────── // // A filtered channel with unscanned listed videos is a channel whose // download queue is WRONG, not merely incomplete: every non-matching video // in it will be prefetched, rejected and discarded one at a time, at one // yt-dlp invocation each, against a source whose patience is the scarce // resource. The scan answers the same question for the whole channel in one // batch. So it is not another kind of work competing with downloads — it is // the thing that decides what the downloads ARE, and it goes first. // // ONE PER CHANNEL PER RUNNER, and never a second while one is in flight. // `scannedThisRun` is the session's memory: a channel whose scan left the // backlog above zero (a soft block, a cooldown, a batch that stopped early) // is NOT retried by this runner — the operator's Run button and the next // restart both still work, and re-offering it on the next three-second tick // would hammer the very source that just refused us. // // It takes a platform slot exactly as a download unit does, so a scan and a // download never hit one source at once, and it inherits the platform // cooldown for free: a platform in `platformSkip` offers no scan either. if (kind === "download") { const scanCandidate = pickMetadataScanChannel(channels, { scanned: scannedThisRun, platformSkip, platformOf: (slug) => platformKey(slug, slugToPlatform), ...(focus ? { focus: { slugs: new Set(ctx.focusSlugs), holding: focus.holding, }, } : {}), }); if (scanCandidate) { const { slug, targets } = scanCandidate; const unitPlatform = platformKey(slug, slugToPlatform); scannedThisRun.add(slug); const videoId = `${METADATA_SCAN_UNIT_PREFIX}${slug}`; const note = `scanning ${targets} listed video${ targets === 1 ? "" : "s" } for ${slug}`; live.idleReason = null; // No `path`, so no node's active count moves: the scan is not a leaf's // work and must not consume a leaf's or a group's maxWorkers. live.inFlight.set(videoId, { videoId, leafId: METADATA_SCAN_LEAF_ID, channelSlug: slug, startedAt: Date.now(), note, }); platformInFlight.set( unitPlatform, (platformInFlight.get(unitPlatform) ?? 0) + 1, ); // IN THE PICK LOG LIKE ANY OTHER UNIT. The log is "what did this lane // do", and a tick that ran a scan and no download would otherwise read // as a tick that did nothing at all. recordPick(kindState, { at: Date.now(), leafId: METADATA_SCAN_LEAF_ID, videoId, channelSlug: slug, }); persist(); onLog(`Auto-download: ${note}.`); return { pick: { leafId: METADATA_SCAN_LEAF_ID, videoId, path: [] }, channelSlug: slug, unitPlatform, scan: { slug, targets }, }; } } const pick = selectNextWork(root, pending, runtime, live.active); if (!pick) { // Attribute the idleness. Nothing pending at all is a different situation // from work that exists but is unreachable, and both differ from work the // platform gates just took away — the page says which. if (pendingBeforeGates === 0) live.idleReason = "no-pending"; else if (countPending(pending) === 0) { live.idleReason = anyCooling ? "cooldown" : anyHeld ? "held" : anyPaced ? "paced" : anyDeferred ? "deferred" : "capped"; } else live.idleReason = "capped"; return null; } const channelSlug = owner.get(pick.videoId); if (!channelSlug) { // Shouldn't happen (every pending id came from a channel) — guard anyway. markCompleted(pick.videoId); return null; } live.idleReason = null; const unitPlatform = kind === "download" ? platformKey(channelSlug, slugToPlatform) : null; // Reserve the slot and account active counts up the path BEFORE returning, // so the next selection sees the reservation (and ancestor caps reflect it). live.inFlight.set(pick.videoId, { videoId: pick.videoId, leafId: pick.leafId, channelSlug, startedAt: Date.now(), }); for (const id of pick.path) live.active[id] = (live.active[id] ?? 0) + 1; if (unitPlatform) { platformInFlight.set( unitPlatform, (platformInFlight.get(unitPlatform) ?? 0) + 1, ); } recordPick(kindState, { at: Date.now(), leafId: pick.leafId, videoId: pick.videoId, channelSlug, }); persist(); onLog(`Auto-${kind}: ${channelSlug}/${pick.videoId} (leaf ${pick.leafId})`); 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 => { const laneOperations = laneOperationIds(kind); const opened = laneRun.run; if (!opened) return { outcome: "skipped" }; // THE ROOT THE PICK CAME FROM, not the stored one: a compiled leaf id // (`prio--`) does not exist in the stored tree, and looking it // up there would silently fall back to the lane's whole operation union. const leaf = flattenLeaves(dispatchRoot).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 => { const { pick, channelSlug, unitPlatform } = picked; let result: UnitResult = { outcome: "failed" }; try { // A RELOCATION THAT STARTED BETWEEN THE PICK AND THE RUN. // // The pending lists are built from snapshots and the reservation above is // taken before this runs; a relocate job can begin in that window. Its // rsync is copying the very `data/` this unit is about to write into — // the omnimirror shape (2026-09-13), where a sidecar written mid-copy // left the verify refusing a 131 GB transfer. A marker is the signal, the // same one every other guard already honours. // // THIS IS THE BACKSTOP, NOT THE DEFENCE. GUARD 2 in buildChannelWork // already drops the channel from the NEXT tick's pending list — any // marker makes inspectChannelMedia report `in-transition`, which is // neither `ok` nor `in-place` — so the only window this closes is // pick→run, and it is measured in milliseconds. // // WHAT THE `finally` DOES WITH IT DIFFERS BY LANE, and neither answer // costs anything. On a BUCKET lane it calls markCompleted, so this video // is retired for the session and not offered again until the runner // restarts or the snapshot regenerates: one unit per channel, accepted, // and far cheaper than a copy refused at its last step. On an OPERATION // lane it retires nothing — runOperationPick owns the (operation, video) // keys and never ran — which is equally fine, because guard 2 has taken // the channel off the list before the next tick could offer it again. // // NOT THE DIGEST LANE (release 17): a media move leaves the text where it // is, and a digest writes only text. A tier migration rebuilds `data/` // itself, and holds it too. const marker = await readRelocationMarker(paths, channelSlug); if (marker && (kind !== "digest" || markerHoldsText(marker))) { onLog( `Auto-${kind}: skipping ${channelSlug}/${pick.videoId}, ` + `${mediaHoldText("in-transition")} — a relocation ` + `(${marker.direction}) to ${marker.target} is in flight.`, ); // `return` still runs the `finally`, which is what releases the // reservation and folds this outcome — a skip, not a failure, so it // never touches the platform backoff. result = { outcome: "skipped" }; return; } if (picked.scan) { result = await launchMetadataScan({ paths, slug: picked.scan.slug, targets: picked.scan.targets, tracker, onLog, onChildJob: (jid) => { childJobIds.set(pick.videoId, jid); const u = live.inFlight.get(pick.videoId); if (u) u.jobId = jid; }, }); return; } 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); const u = live.inFlight.get(pick.videoId); if (u) u.jobId = jid; }, }); } catch (err) { result = { outcome: "failed" }; onLog(`Auto-${kind} unit error on ${pick.videoId}: ${String(err)}`); } finally { childJobIds.delete(pick.videoId); live.inFlight.delete(pick.videoId); for (const id of pick.path) { live.active[id] = Math.max(0, (live.active[id] ?? 0) - 1); } if (unitPlatform) { platformInFlight.set( unitPlatform, Math.max(0, (platformInFlight.get(unitPlatform) ?? 0) - 1), ); } 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)); } // Pacing: a rate-limit / network failure backs the whole platform off, // and a rate-limited video is also deferred so the next pick after the // cooldown is a DIFFERENT video — see unitOutcome.ts. Any other outcome // retires the video for the session; a success clears the cooldown. const settledAt = Date.now(); const unitSettings = getSettings(); const effect = applyUnitOutcome( kindState, { platform: unitPlatform, videoId: pick.videoId, channelSlug, outcome: result.outcome, ...(result.failureClass ? { failureClass: result.failureClass } : {}), }, settledAt, Math.random, unitSettings.pacing, ); if (effect.line) onLog(effect.line); // The gap before this platform's next unit, at the pace this outcome // left it (release 17, slice RL). if (unitPlatform) { const base = staticSleepRequestsSeconds(unitPlatform); const gap = downloadGapMs( unitSettings.sleepBetweenDownloadsSeconds, currentPaceSeconds(kindState.platformPace, unitPlatform, base), base, { minSeconds: platformMinGapSeconds(unitPlatform) }, ); if (gap > 0) platformNextStartAt.set(unitPlatform, settledAt + gap); else platformNextStartAt.delete(unitPlatform); } // 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 (effect.markCompleted && !isOperationLane(kind)) { markCompleted(pick.videoId); } // Persist so the cooldown (and reset) survive a restart. persist(); // The runner job has no channelSlug, so the registry's own // snapshot-on-finish hook doesn't fire — request a regen here so the // channel's pending counts refresh and a restart sees fresh buckets. requestChannelSnapshot(paths, channelSlug); } }; onLog(`Auto-${kind} runner started.`); // One audited fill-to-capacity primitive replaces the old hand-written // wake/waitNext loop AND its terminal drain-wait. runPool stops pulling on a // drain (or stop), lets in-flight units finish, and never spins the event // loop on an already-aborted signal — the "Drain all hangs the app" bug class. await runPool({ next, run, limit, signal, drainSignal, idlePollMs: IDLE_POLL_MS, }); persist(); onLog( `Auto-${kind} runner finished: ${counters.transcribed} done, ${counters.failed} failed, ${counters.skipped} skipped.`, ); } // ONE METADATA SCAN, as a child job on the channel's platform download queue. // // The SAME job the operator's Run button dispatches — kind, queue key, replay // spec and the shared per-platform cooldown all come from // controller/metadataScanJob.ts — so the registry serializes it against a // manual sync on that platform exactly as it serializes a download unit, and // there is one definition of what a scan IS. // // `background: true` for the same reason a download unit is: a clicked action // preempts queued units without interrupting a running one. // // A SCAN IS NEVER A FAILURE FOR BACKOFF PURPOSES. runMetadataScan does not // throw on a rate limit — it stops, records the shared cooldown through its own // onPlatformBackoff, and keeps what it flushed. So this returns "skipped" or // "transcribed" and never a `failureClass`: the cooldown is already recorded by // the time we get here, and a second backoff entry from the runner would // double-count one refusal. async function launchMetadataScan(args: { paths: Paths; slug: string; targets: number; tracker: ReturnType; onLog: (line: string) => void; onChildJob?: (jobId: string) => void; }): Promise { const config = await readChannelConfig(args.paths, args.slug); if (!config) return { outcome: "skipped" }; const res = await runMetadataScanJob({ paths: args.paths, slug: args.slug, channelConfig: config, background: true, }); if (!res.ok) { args.onLog( `Auto-download: could not enqueue the metadata scan for ${args.slug}: ${res.error}`, ); return { outcome: "skipped" }; } args.onChildJob?.(res.jobId); const term = await res.done; if (term.status === "cancelled") return { outcome: "skipped" }; if (term.status === "failed") { args.onLog(`Auto-download: the metadata scan for ${args.slug} failed.`); return { outcome: "skipped" }; } return { outcome: "transcribed" }; } type LaunchArgs = { kind: AutoQueueKind; paths: Paths; pick: WorkPick; channelSlug: string; tracker: ReturnType; onLog: (line: string) => void; signal: AbortSignal; // Soft-cancel (drain): forwarded into the transcription unit's worker-pool // acquire so a parked (not-yet-started) unit unblocks and skips on drain — // otherwise the runner's drain loop waits forever for it to clear. A running // transcription is unaffected (drainSignal never reaches the engine). Unused by // the download branch (downloads serialize on the registry queue, not the pool, // and a drain lets the in-flight download finish by design). drainSignal?: AbortSignal; // Download only: called with the child download job's id once enqueued, so the // runner can hard-cancel it when stopped. No-op for transcription units. onChildJob?: (jobId: string) => void; }; // The outcome of one unit. `failureClass` is set only for download failures // (the runner uses rate_limit/network to drive per-platform backoff); it's // derived from the returned outcome record via classifyDownloadFailure. type UnitResult = { outcome: "transcribed" | "skipped" | "failed"; failureClass?: DownloadFailureClass; }; // Whether this unit is an auto-captions replacement: the video's only // transcript is a YouTube ASR VTT. Decided from disk state (one readdir + a 4 KB // VTT head read) rather than plumbed from the pick, so it stays correct no // matter which leaf/bucket selected the video — including the default union // under `replaceAutoSubs`. Both branches below need it: the transcription gate // would otherwise skip the video as "already transcribed", and the download // branch would fetch subtitles instead of audio. async function isAutoSubsUnit( paths: Paths, channelSlug: string, videoId: string, ): Promise { const dir = path.join(paths.channelsDir, channelSlug, "data", videoId); const files = await readVideoFiles(dir, { checkUntranscribable: true }); return isAutoSubsOnly(dir, files); } // Run a single unit of work. Returns the outcome for counters/backoff. Hard // cancel (signal) aborts an in-flight unit; drain is handled by the loop (it // stops launching new units), so the drain signal is intentionally NOT // forwarded here. async function launchUnit(args: LaunchArgs): Promise { const replaceAutoSubs = await isAutoSubsUnit( args.paths, args.channelSlug, args.pick.videoId, ); if (args.kind === "transcription") { if (replaceAutoSubs) { args.onLog( `Auto-transcribe: ${args.pick.videoId} has only YouTube auto-captions — transcribing over them.`, ); } const res = await transcribeOneFromQueue({ paths: args.paths, channelSlug: args.channelSlug, videoId: args.pick.videoId, audioFilename: "audio.mp3", strict: false, tracker: args.tracker, // Show which channel each in-flight transcription belongs to (the runner // is cross-channel, so the bare videoId wouldn't say). taskLabel: `${args.channelSlug}/${args.pick.videoId}`, onLog: args.onLog, signal: args.signal, // Drain (soft cancel): unblocks this unit's worker-pool acquire if it's // parked waiting for a slot, so the runner's drain can complete instead of // hanging. A unit already transcribing finishes (the engine only sees the // hard signal). drainSignal: args.drainSignal, // Acquire the worker pool slot at background priority: a manual (foreground) // transcribe preempts queued auto units for the next freed slot. The running // unit is never interrupted — it drains, then the slot goes to the manual // waiter. Mirrors how auto-download units yield to a manual sync. background: true, // Relax the "already transcribed" gate for an ASR-only video: only OUR // transcript.json counts, so the auto-captions get transcribed over. replaceAutoSubs, }); return { outcome: res.outcome }; } // Download: resolve the video's source URL and run the same managed // single-video download the import/download-missing paths use (so archive // appends, auth/no-subs fallbacks and audio-integrity checks all apply) — but // as its OWN managed job on the channel's platform queue (downloadQueueKey), // the SAME queue a manual sync uses. The registry then serializes this unit // against any sync on that platform, so the two never spawn yt-dlp at once. // Marked `background` so a clicked sync preempts queued units (without // interrupting a running one). The runner's per-platform in-flight gate still // keeps this to one outstanding unit per platform so the queue isn't flooded. const rawConfig = await readChannelConfig(args.paths, args.channelSlug); if (!rawConfig) return { outcome: "skipped" }; // A `handling: "youtube"` channel's normal download passes --write-auto-subs // --write-subs --skip-download: it would re-fetch the very auto-captions we're // replacing and never touch the audio. Force transcribe-handling for this one // video so it downloads audio instead; the channel's stored config is // untouched. (This is the same override the manual bucket button passes as // handlingOverride.) Nothing to do when the channel already downloads audio. const config = replaceAutoSubs && rawConfig.handling !== "transcribe" ? { ...rawConfig, handling: "transcribe" as const } : rawConfig; if (config !== rawConfig) { args.onLog( `Auto-download: ${args.pick.videoId} has only YouTube auto-captions — downloading audio (handling override: transcribe).`, ); } const url = await findVideoSourceUrl( args.paths, args.channelSlug, args.pick.videoId, config, ); if (!url) { args.onLog(`Auto-download: no source URL for ${args.pick.videoId} — skipped`); return { outcome: "skipped" }; } const settings = getSettings(); // The child job's fn sets these synchronously before it returns; the runner // reads them after `await res.done` resolves (single-threaded → no race). let unitStatus: DownloadOutcomeStatus | undefined; let unitFailureClass: DownloadFailureClass | undefined; const res = await runManagedFunction({ kind: AUTO_DOWNLOAD_UNIT_KIND, queueKey: downloadQueueKey(config), paths: args.paths, channelSlug: args.channelSlug, videoId: args.pick.videoId, background: true, fn: async (onLog, signal, _setProgress, ctx) => { // Per-video progress lives on THIS child job's row (its own ctx), so the // unit shows up in Active Jobs like a manual single-video download. const task = makeTaskTracker(ctx, onLog).start({ id: args.pick.videoId, label: `${args.channelSlug}/${args.pick.videoId}`, kind: "download", }); try { // downloadOneManaged does NOT throw on a 429/network failure — it // returns a record whose last attempt carries the stderr tail + // availability class, already classified into failureClass. const record = await downloadOneManaged({ channelSlug: args.channelSlug, channelConfig: config, paths: args.paths, videoUrl: url, onLog: task.onLog, // The child job's own abort signal: registry.cancel(childJobId) aborts // it (→ kills yt-dlp) when the runner is hard-cancelled. signal, cookiePolicy: resolveCookiePolicy(settings, config), inlineTranscribeOnFallback: settings.inlineTranscribeOnFallback, globalSkipLiveDownloads: settings.skipLiveDownloads, downloadFormatPreset: resolveDownloadFormatPreset({ channel: config.downloadFormat, global: settings.downloadFormat, }), appendArchive: true, }); unitStatus = record.status; unitFailureClass = record.failureClass; // Land a failed download as a failed JOB (red row) rather than a silent // "done"; the runner reads the captured status/class regardless. A // short-audio download kept its file (so it won't be re-queued) but // produced no usable audio, so surface it as failed too. if ( record.status === "failed" || record.status === "failed-corrupt-source" || record.status === "failed-short-audio" ) { throw new Error( record.status === "failed-short-audio" ? "download produced truncated audio (short-audio)" : `download failed (${record.failureClass ?? "unknown"})`, ); } } finally { task.end(); } }, }); if (!res.ok) { args.onLog( `Auto-download: could not enqueue ${args.pick.videoId}: ${res.error}`, ); return { outcome: "failed" }; } args.onChildJob?.(res.jobId); const term = await res.done; // Cancelled (runner stopped, or dropped while still queued) → not a failure. if (term.status === "cancelled") return { outcome: "skipped" }; if (unitStatus === "skipped-filtered") return { outcome: "skipped" }; // The filter declined the media and the live chat was fetched instead. A // SUCCESS for the lane — the unit did the work it was picked for, the video // leaves chatOnlyPending, and a platform whose chat pass came back clean has // its cooldown cleared exactly as a download would clear it. if (unitStatus === "chat-only") return { outcome: "transcribed" }; // Complete-but-malformed source: terminal and kept on disk. Treat as skipped // (not failed) so it doesn't drive backoff and isn't re-picked for download. if (unitStatus === "corrupt-full-source") return { outcome: "skipped" }; // A success whose SUBTITLE fetch alone answered 429 says so: the runner // defers the video's subtitles and leaves the platform alone (release 17). if (unitStatus && unitStatus.startsWith("ok")) { return unitFailureClass === "subs_rate_limit" ? { outcome: "transcribed", failureClass: "subs_rate_limit" } : { outcome: "transcribed" }; } // failureClass drives the runner's per-platform backoff (rate_limit/network). return { outcome: "failed", failureClass: unitFailureClass }; } // --- Lifecycle ------------------------------------------------------------- // Launch a runner job for `kind` if its policy is enabled and one isn't already // running. Idempotent across HMR reloads / repeated calls (it checks the live // registry). Returns the job id, or null when not started. export async function startAutoRunner( kind: AutoQueueKind, paths: Paths = getPaths(), ): Promise { const settings = getSettings(); const policy = settings.autoQueue[kind]; if (!policy.enabled) return null; const singleton = getSingleton(); const existing = singleton.runners.get(kind); if (existing && getRegistry().get(existing.jobId)?.status === "running") { return existing.jobId; // already running } const live: RunnerLive = { jobId: "", inFlight: new Map(), active: {}, startedAt: Date.now(), idleReason: null, }; const result = await runManagedFunction({ kind: autoRunnerJobKind(kind), queueKey: "", // run in parallel to everything; not serialized paths, fn: async (onLog, signal, setProgress, ctx) => { live.jobId = ctx.jobId; try { await runLoop(kind, paths, live, onLog, setProgress, signal, ctx); } finally { if (singleton.runners.get(kind) === live) { singleton.runners.delete(kind); } } }, }); if (!result.ok) return null; live.jobId = result.jobId; singleton.runners.set(kind, live); 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 — 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 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 // server startup. // // 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 { for (const lane of LANES) await startAutoRunner(lane, paths); } // Stop a runner (hard cancel its job). In-flight units are aborted. Returns true // if one was running. export function stopAutoRunner(kind: AutoQueueKind): boolean { const live = getSingleton().runners.get(kind); if (!live) return false; return getRegistry().cancel(live.jobId); } // Drain a runner (graceful stop): stop picking new work, let in-flight units // finish, then the loop exits and the job ends "done". Returns true if one was // running. Same soft-cancel the Active Jobs "Drain" button uses, just addressed // by kind instead of job id. export function drainAutoRunner(kind: AutoQueueKind): boolean { const live = getSingleton().runners.get(kind); if (!live) return false; return getRegistry().requestDrain(live.jobId); }