import { readFile } from "node:fs/promises"; import { writeJsonAtomic } from "../lib/jsonFile-server"; import type { Paths } from "../lib/paths"; import { type AutoQueueRuntime, emptyAutoQueueRuntime, } from "./autoQueuePolicy"; import { type PlatformBackoffState, type PlatformHoldState, type PlatformPaceState, type SubtitleDeferralState, type VideoDeferralState, coercePlatformBackoff, coercePlatformHolds, coercePlatformPace, coerceSubtitleDeferrals, coerceVideoDeferrals, effectivePaceSeconds, } from "./platformBackoff"; // Persistent fairness state for the auto-queue runners. Unlike in-flight worker // counts (which are zero by definition after a restart and live only on the // runner singleton), the SWRR current-weights are "fairness memory": persisting // them avoids always restarting round-robin from the first child after a reboot. // // Best-effort and small: one runtime per kind plus a bounded pick log for the // observability panel. Written atomically (tmp + rename). This module is pure // I/O + types — the selection logic lives in autoQueuePolicy.ts and the // orchestration in autoRunner.ts. Mirrors syncSchedulerState.ts. // Defined in the model layer (lib/ may not import jobs/); re-exported here so // every existing `from "./autoQueueState"` import still resolves. export type { AutoQueueKind } from "../lib/autoQueueTypes"; import { LANES, type AutoQueueKind } from "../lib/autoQueueTypes"; // One recorded grant, newest-first, for the status panel. export type AutoQueuePick = { at: number; leafId: string; videoId: string; channelSlug: string; }; export type AutoQueueKindState = { runtime: AutoQueueRuntime; // Newest-first, bounded to AUTO_QUEUE_PICK_LOG_LIMIT entries. picks: AutoQueuePick[]; // Per-platform download cooldowns (rate-limit/network backoff). Only used by // the "download" kind; empty for transcription. Persisted so an Odysee 429 // cooldown survives a server restart instead of re-storming on boot. platformBackoff: PlatformBackoffState; // Per-video deferrals (a rate-limited video is skipped by the auto-download // pick for VIDEO_RATE_LIMIT_DEFER_MS). Download kind only; `{}` elsewhere. // Persisted so a restart does not re-hit the same video at fails+1. A file // written before the key existed coerces to `{}`; an older build ignores it. videoDeferrals: VideoDeferralState; // THE PACING MEMORY (release 17, slice RL) — download kind only, `{}` // elsewhere; a file written before a key existed coerces it to `{}` and an // older build drops it on its next write. See platformBackoff.ts. // Per-platform `--sleep-requests` above the static value, while it is above. platformPace: PlatformPaceState; // Platforms whose rate limit outlasted the cooldown cap: one probe per // `pacing.holdProbeMinutes` until a clean one. platformHolds: PlatformHoldState; // Videos whose SUBTITLE fetch answered 429 (their media did not): the count, // and until when the batch subtitle fetch leaves them alone. subtitleDeferrals: SubtitleDeferralState; }; // KEYED BY EVERY LANE, including the two with no executor yet. A lane whose // block is missing from the file coerces to empty state rather than to // `undefined`: `state[kind].picks` is read on the status poll for every lane the // console draws, and a lane added to LANES must not be able to make that throw // against a state file written before it existed. export type AutoQueueState = Record; export const AUTO_QUEUE_PICK_LOG_LIMIT = 50; export function emptyAutoQueueKindState(): AutoQueueKindState { return { runtime: emptyAutoQueueRuntime(), picks: [], platformBackoff: {}, videoDeferrals: {}, platformPace: {}, platformHolds: {}, subtitleDeferrals: {}, }; } export function emptyAutoQueueState(): AutoQueueState { return Object.fromEntries( LANES.map((lane) => [lane, emptyAutoQueueKindState()]), ) as AutoQueueState; } function coerceRuntime(value: unknown): AutoQueueRuntime { const out = emptyAutoQueueRuntime(); if (!value || typeof value !== "object") return out; const r = value as Record; if (r.currentWeights && typeof r.currentWeights === "object") { for (const [id, w] of Object.entries( r.currentWeights as Record, )) { if (typeof w === "number" && Number.isFinite(w)) { out.currentWeights[id] = w; } } } return out; } function coercePick(value: unknown): AutoQueuePick | null { if (!value || typeof value !== "object") return null; const r = value as Record; if ( typeof r.leafId !== "string" || typeof r.videoId !== "string" || typeof r.channelSlug !== "string" ) { return null; } return { at: typeof r.at === "number" && Number.isFinite(r.at) ? r.at : 0, leafId: r.leafId, videoId: r.videoId, channelSlug: r.channelSlug, }; } function coerceKindState(value: unknown): AutoQueueKindState { if (!value || typeof value !== "object") return emptyAutoQueueKindState(); const r = value as Record; const picks = Array.isArray(r.picks) ? r.picks .map(coercePick) .filter((p): p is AutoQueuePick => p !== null) .slice(0, AUTO_QUEUE_PICK_LOG_LIMIT) : []; return { runtime: coerceRuntime(r.runtime), picks, platformBackoff: coercePlatformBackoff(r.platformBackoff), videoDeferrals: coerceVideoDeferrals(r.videoDeferrals), platformPace: coercePlatformPace(r.platformPace), platformHolds: coercePlatformHolds(r.platformHolds), subtitleDeferrals: coerceSubtitleDeferrals(r.subtitleDeferrals), }; } // Read the state file, tolerating a missing/corrupt file by returning empty // state. Defensive: any stored shape that doesn't match is coerced so a // hand-edited file can't crash the runner. export async function readAutoQueueState(paths: Paths): Promise { let raw: unknown; try { raw = JSON.parse(await readFile(paths.autoQueueStateFile, "utf8")); } catch { return emptyAutoQueueState(); } if (!raw || typeof raw !== "object") return emptyAutoQueueState(); const r = raw as Record; const state = Object.fromEntries( LANES.map((lane) => [lane, coerceKindState(r[lane])]), ) as AutoQueueState; notePace(state); return state; } // Write the state atomically (tmp + rename), creating the .auto-queue dir on // first use. Each kind's pick log is trimmed to the limit on the way out. // UNIQUE PER WRITE, not per process. `persist()` is fire-and-forget and two // lanes' runners can persist within a few milliseconds of each other, so one // tmp name per process means the second write truncates the file the first is // still writing and then renames it over the real one — leaving a state.json // that reads as empty (tolerated) and loses a persisted platform cooldown // across a restart (not). `writeJsonAtomic` gives every write its own temp // name and chains writes to one path on globalThis — which this module's own // counter, one per module COPY, did not. export async function writeAutoQueueState( paths: Paths, state: AutoQueueState, ): Promise { const trim = (k: AutoQueueKindState): AutoQueueKindState => ({ runtime: k.runtime, picks: k.picks.slice(0, AUTO_QUEUE_PICK_LOG_LIMIT), platformBackoff: k.platformBackoff ?? {}, videoDeferrals: k.videoDeferrals ?? {}, platformPace: k.platformPace ?? {}, platformHolds: k.platformHolds ?? {}, subtitleDeferrals: k.subtitleDeferrals ?? {}, }); const out = Object.fromEntries( LANES.map((lane) => [lane, trim(state[lane] ?? emptyAutoQueueKindState())]), ) as AutoQueueState; notePace(out); await writeJsonAtomic(paths.autoQueueStateFile, out, { mkdir: true }); } // ONE SHARED OBJECT PER PROCESS — the holder every lane's runner and // `recordDownloadBackoff` read and write through. Its own global (not the // auto-runner singleton, which lives in controller/ where jobs/ may not reach), // cleared by the e2e harness beside `__yttAutoRunner__`. type AutoQueueStateHolder = { // The ONE persisted state object, shared by every lane's runner. state: AutoQueueState | null; stateFile: string | null; // The READ in flight, so every caller in the same window awaits ONE read and // receives ONE object. loading: Promise | null; loadingFile: string | null; // The download lane's pace as last read or written, by ANY reader of the // file — for `livePlatformPaceSeconds` when no runner holds the object. lastPace: PlatformPaceState | null; }; declare global { // eslint-disable-next-line no-var var __yttAutoQueueState__: AutoQueueStateHolder | undefined; } function getHolder(): AutoQueueStateHolder { if (!globalThis.__yttAutoQueueState__) { globalThis.__yttAutoQueueState__ = { state: null, stateFile: null, loading: null, loadingFile: null, lastPace: null, }; } return globalThis.__yttAutoQueueState__; } function notePace(state: AutoQueueState): void { getHolder().lastPace = state.download?.platformPace ?? {}; } // THE PACE A SPAWN USES, SYNCHRONOUSLY (release 17, slice RL). Every yt-dlp // argv in the repo is built by `channelExtraArgs`, which is synchronous and is // called from a dozen places; threading an awaited state read through each // would be a dozen chances to forget one. So this answers from what the process // already holds: the shared object when a runner holds it (the one every lane // and `recordDownloadBackoff` write through), else the pace the last read or // write of the file saw (the status poll reads it every few seconds), else // undefined — the caller then uses the static pace. Never reads the disk. export function livePlatformPaceSeconds(platform: string): number | undefined { const holder = getHolder(); const pace = holder.state?.download?.platformPace ?? holder.lastPace; const entry = pace?.[platform]; // With the hourly easing applied (platformBackoff.ts, PACE_TIME_DECAY_MS). return entry ? effectivePaceSeconds(entry, Date.now()) : undefined; } // THE PERSISTED STATE IS ONE OBJECT FOR EVERY LANE'S RUNNER, and it has to be. // // `writeAutoQueueState` serializes the WHOLE file — all four lanes — so two // runners each holding their own copy means every persist clobbers the other // lane's pick log and fairness memory with whatever that copy was read with. // With one slow lane that is nearly invisible (a transcription is minutes // apart), and it became obvious the moment two fast lanes ran together: the // digest lane's Recent picks panel emptied itself on every backfill dispatch, // and `/api/auto-queue/status` — which reads the FILE — showed zero picks for a // lane that was demonstrably working. `lane-runner.spec.ts`'s concurrency test // is what found it. // // One object mutated by both is what makes a write of it true for both. Its // lifetime is the holder's, which the e2e harness clears between specs // (api/test/invalidate-cache) — so a reset corpus does not inherit a previous // spec's picks or backoff. // // THE CACHE IS THE PROMISE, NOT THE VALUE — and until release 8 it was the // value, which made the paragraph above false at exactly the moment it matters. // Boot (`startAutoRunnersIfEnabled`) starts all four lanes' `runLoop`s without // awaiting them, so all four reached the `await readAutoQueueState` here before // any of them had set the cached state; each read its own copy and each // runLoop kept that private copy for its lifetime. Live, 2026-09-25: the // download lane persisted `fails: 13`, a later `until` and a 6 h video deferral // at 13:20; at 13:33:37 another lane persisted its boot-time copy of the whole // file (`fails: 12`, deferrals `{}`) and silently undid the pacing on disk. // Caching the in-flight read keyed by the state file means every caller in the // same window awaits one read and receives the same object. A reset (the e2e // harness dropping `__yttAutoQueueState__`) replaces the holder wholesale, so a // read still in flight resolves into the OLD holder and cannot leak across. // // It lives in jobs/ (moved from controller/autoRunner.ts in release 8) because // `recordDownloadBackoff` must write THROUGH it: a manual Sync's 429 cooldown // written straight to disk was erased by the next persist of any lane, and never // merged at all while the download lane was paused, stopped or at capacity. export async function sharedAutoQueueState( paths: Paths, ): Promise { const holder = getHolder(); const file = paths.autoQueueStateFile; if (holder.state && holder.stateFile === file) { return holder.state; } if (holder.loading && holder.loadingFile === file) { return holder.loading; } const loading = readAutoQueueState(paths).then( (state) => { // A later call for a DIFFERENT file superseded this read; it owns the // holder now, and this read's object goes only to its own awaiters. if (holder.loading === loading) { holder.state = state; holder.stateFile = file; holder.loading = null; holder.loadingFile = null; } return state; }, (err: unknown) => { // Never cache a failure: the next caller retries the read. if (holder.loading === loading) { holder.loading = null; holder.loadingFile = null; } throw err; }, ); holder.loading = loading; holder.loadingFile = file; return loading; } // The shared object if something in this process already holds (or is // loading) it for this file, else null — WITHOUT starting a read. For writers // outside the runner (`recordDownloadBackoff`): with no runner live there is no // in-memory copy to keep in step, and the plain disk read-modify-write is right. export async function liveAutoQueueState( paths: Paths, ): Promise { const holder = getHolder(); const file = paths.autoQueueStateFile; if (holder.state && holder.stateFile === file) return holder.state; if (holder.loading && holder.loadingFile === file) { try { return await holder.loading; } catch { return null; } } return null; } export function recordPick(state: AutoQueueKindState, pick: AutoQueuePick): void { state.picks.unshift(pick); if (state.picks.length > AUTO_QUEUE_PICK_LOG_LIMIT) { state.picks.length = AUTO_QUEUE_PICK_LOG_LIMIT; } }