import { existsSync } from "node:fs"; import { open as openFile, readFile } from "node:fs/promises"; import path from "node:path"; import { open } from "lmdb"; import type { Paths } from "../lib/paths"; import { mapConcurrent } from "../lib/concurrency"; import { extractVideoId } from "../lib/videoId"; import { uploadKeyFor } from "./keptVideos"; import type { AutoQueueOrder } from "../jobs/autoQueuePolicy"; import type { ChannelConfig } from "../lib/channelConfig"; import { isDriveNotAnswering, onDrive, stalledLocationForPath, } from "../lib/storageHealth"; // Upload-date lookup for the auto-queue's "newest first" ordering. // // The problem: the runners pick pending[leaf][0], and that order comes from // snapshot.json's bucket arrays, which channelSnapshot sorts lexicographically // by video id. For YouTube ids that is arbitrary; for YYYYMMDD_-prefixed dir // names it is oldest-first. Neither lets a today's-upload jump a 9,000-video // backlog. To sort by recency we need an upload date per candidate video, for // tens of thousands of candidates, on every scheduling tick — so it has to be // nearly free. // // It nearly is, because the transcript index already stores it. buildIndex's // `byChannel` sub-DB is keyed [channelSlug, uploadDate, id] with a constant // value, so a KEY-ONLY range scan reads every video's date without decoding a // single value. What that scan cannot supply is covered by a bounded tail read. // // Four layers, in order: // 1. the LMDB `byChannel` scan — ~76 ms warm for all 78,583 indexed videos, // 99% of the corpus. BUT it is the TRANSCRIPT index, so it holds only // videos that already have a transcript. Measured against the live corpus, // that covers just 122 of the 870 videos auto-transcribe actually has // pending — the bucket it orders is `downloadedNoTranscript`, which is by // definition the set this layer cannot see. It is still layer 1 because it // dates the whole corpus for one cheap scan, and it supplies the anchors // layer 3 interpolates between; // 2. an 8 KB TAIL read of each remaining video's metadata.info.json, regexed // for upload_date. Measured on the real backlog: 868/870 hits (99.8%) at // 0.19 ms/file, 164 ms for the lot. This is the layer that actually orders // auto-transcribe. Memoized process-wide — an upload date never changes — // so the cost is paid once per video, not once per refresh; // 3. playlist-neighbour interpolation, for auto-download — an UNdownloaded // video has no metadata.info.json at all, so it has no date anywhere on // disk. The stored playlist is newest-first, so the nearest PRECEDING // dated entry is an upper bound on its age; // 4. uploadKeyFor's YYYYMMDD_ dir-name prefix, then "" (sorts oldest). // // Reading each metadata.info.json in FULL was measured and rejected: ~6.5 min // and ~41 GiB of I/O. So was roster.json's firstSeenAt (99.7% of entries // collapse onto one seed timestamp) and the served stats pages (per-site // filtered, capped at 30,990 rows). // // Everything here degrades rather than throws: a missing, locked or corrupt // index just means fewer known dates, never a runner that fails to dispatch. // A video's recency sort key. `key` is a YYYYMMDD string (or one of the two // sentinels below); `estimated` marks a key that was interpolated rather than // read, so the UI can render it as "≈2026-08-12" instead of claiming precision. export type RecencyKey = { key: string; estimated: boolean }; // Sorts above every real date. Given to an undownloaded video that sits ahead of // every dated entry in its channel's newest-first playlist — i.e. it was // uploaded after everything we have. That is precisely the video this feature // exists to promote, so it goes to the very top of the archive, not merely to // the top of its own channel. Guarded by requiring at least one dated anchor in // the playlist (see interpolateFromPlaylist), so a brand-new channel with // nothing downloaded cannot flood the head of the queue. const NEWER_THAN_ANYTHING = "￿"; // Sorts below every real date — what uploadKeyFor already returns for an // undatable id. Named only for readability at the call sites. const UNKNOWN = ""; // --- Layer 1: the LMDB byChannel key scan ----------------------------------- // [channelSlug, uploadDate, id] — buildIndex.ts's ChannelKey. type ChannelKey = [string, string, string]; export type UploadDateIndex = { // slug -> (videoId -> YYYYMMDD). Empty when the index is unavailable. datesBySlug: Map>; close(): Promise; }; const EMPTY_INDEX: UploadDateIndex = { datesBySlug: new Map(), close: async () => {}, }; // Read-only view of the index's upload dates, scanned per channel. Modeled on // openChannelSigner (common/lib/channelSignature.ts): guard on existsSync, wrap // the open in try/catch, and degrade to an empty stub rather than fail — the // runner must keep dispatching even when the index is mid-rebuild or absent. export function openUploadDateIndex( paths: Paths, slugs: ReadonlyArray, ): UploadDateIndex { if (!existsSync(paths.lmdbPath)) return EMPTY_INDEX; let root: ReturnType; try { root = open({ path: paths.lmdbPath, readOnly: true, maxDbs: 14 }); } catch { return EMPTY_INDEX; } try { const byChannel = root.openDB({ name: "byChannel", encoding: "msgpack", }); const datesBySlug = new Map>(); for (const slug of slugs) { const dates = new Map(); // Key-only iteration: getRange yields {key, value} but we never touch // .value, and lmdb-js decodes values lazily, so no msgpack decode happens. for (const { key } of byChannel.getRange({ start: [slug], end: [slug, "￿"], })) { const k = key as ChannelKey; if (k[0] !== slug) break; // First writer wins, matching the archive's own dedup: an id can only // appear once per channel, but a mid-rebuild index may briefly hold a // stale row alongside a fresh one. if (!dates.has(k[2])) dates.set(k[2], k[1]); } datesBySlug.set(slug, dates); } return { datesBySlug, close: async () => { await root.close(); }, }; } catch { void root.close().catch(() => {}); return EMPTY_INDEX; } } // --- Layer 2: the metadata.info.json tail read ------------------------------- // yt-dlp writes upload_date near the END of metadata.info.json, so an 8 KB tail // finds it without paying for a multi-hundred-KB file (some carry every // subtitle track and comment). Measured 868/870 on the live corpus. const TAIL_BYTES = 8192; const UPLOAD_DATE_RE = /"upload_date":\s*"(\d{8})"/; // Upload dates are immutable, so a date read once is a date forever. Memoizing // process-wide turns this layer from a per-refresh cost into a one-off: after // the first fill the whole layer is map lookups. `null` memoizes a genuine miss // (no metadata, or upload_date outside the tail) so it is not retried every // refresh either. const tailMemo = new Map(); // Backstop against unbounded growth on a corpus far larger than this one. const TAIL_MEMO_CAP = 200_000; // How many tail reads one build may perform. At 0.19 ms each this is ~1.5 s // worst case, and only on the FIRST refresh after a huge backlog appears — // subsequent refreshes hit the memo. Anything past the cap falls through to // layer 4 for now and is picked up by a later refresh, which is the right // failure mode: an unrankable video sorts oldest, i.e. to the back of a // newest-first queue, rather than blocking a scheduling tick. const TAIL_READS_PER_BUILD = 8000; const TAIL_READ_CONCURRENCY = 32; // Hitting the cap is EXPECTED once, not a fault: the first pass that orders a // corpus-wide operation (digest's 77,508 pending, say) has tens of thousands of // undated ids and the memo is cold. It self-corrects — each refresh keys another // 8,000 and the rest sort oldest meanwhile — so the line says so explicitly, // and says it once, or a 3-second poll would print it for an hour. let tailCapLogged = false; async function readTailUploadDate(file: string): Promise { let fh; try { fh = await openFile(file, "r"); const { size } = await fh.stat(); if (size === 0) return null; const len = Math.min(TAIL_BYTES, size); const buf = Buffer.allocUnsafe(len); await fh.read(buf, 0, len, size - len); // latin1 never throws on a multi-byte sequence split by the tail boundary, // and the pattern we want is pure ASCII. return buf.toString("latin1").match(UPLOAD_DATE_RE)?.[1] ?? null; } catch { return null; } finally { await fh?.close().catch(() => {}); } } // Date the ids in `wanted` from their on-disk metadata. Removes each id it // keys from `wanted`, like interpolateFromPlaylist. // // `drives` maps a channel whose TEXT is on another drive — a `legacy` channel, // the retired whole-directory layout, keyed by its `dataDir` — to that target. // Since release 17 a relocated channel's text (`metadata.info.json` here) is on // the corpus disk and only its media is on the far drive, so its reads pass no // drive at all: onDrive is keyed by file kind. A legacy channel's reads go // through `onDrive`: none while that drive's location is // stalled (lib/storageHealth.ts) — each would hold an I/O thread until the drive // came back, 32 at a time — at most `inFlightPerLocation` (4) in flight on it, // and one that has not answered within the budget (3 s by default) marks it // stalled. An id not read for that reason is NOT // memoized as a miss: it falls through to layers 3 and 4 for now, and a later // refresh with the drive answering reads it. const NOT_READ = Symbol("not read: the drive is not answering"); async function datesFromMetadata( paths: Paths, owner: ReadonlyMap, wanted: Set, out: Map, drives: ReadonlyMap = new Map(), ): Promise { const todo: string[] = []; for (const id of wanted) { const memo = tailMemo.get(id); if (memo !== undefined) { if (memo) { out.set(id, { key: memo, estimated: false }); wanted.delete(id); } continue; } const slug = owner.get(id); if (slug === undefined) continue; const drive = drives.get(slug); if (drive && stalledLocationForPath(drive)) continue; todo.push(id); if (todo.length >= TAIL_READS_PER_BUILD) { if (!tailCapLogged) { tailCapLogged = true; console.log( `[recency] dating ${TAIL_READS_PER_BUILD} of ${wanted.size} undated candidates this pass; ` + `the remainder sort oldest until a later refresh keys them (expected on a first corpus-wide pass)`, ); } break; } } if (todo.length === 0) return; const dates = await mapConcurrent( todo, TAIL_READ_CONCURRENCY, async (id): Promise => { const slug = owner.get(id) as string; const file = path.join( paths.channelsDir, slug, "data", id, "metadata.info.json", ); const drive = drives.get(slug); if (!drive) return readTailUploadDate(file); try { return await onDrive(drive, () => readTailUploadDate(file)); } catch (err) { if (isDriveNotAnswering(err)) return NOT_READ; throw err; } }, ); for (const [i, id] of todo.entries()) { const date = dates[i]; if (date === NOT_READ) continue; if (tailMemo.size < TAIL_MEMO_CAP) tailMemo.set(id, date); if (!date) continue; out.set(id, { key: date, estimated: false }); wanted.delete(id); } } // --- Layer 3: playlist-neighbour interpolation ------------------------------ async function readPlaylistIds( paths: Paths, slug: string, ): Promise { try { const raw = await readFile( path.join(paths.channelsDir, slug, "playlist"), "utf8", ); const ids: string[] = []; for (const line of raw.split("\n")) { const trimmed = line.trim(); if (!trimmed) continue; const id = extractVideoId(trimmed); if (id) ids.push(id); } return ids; } catch { return []; } } // Estimate a date for each id in `wanted` from its position in the channel's // newest-first playlist: take the nearest PRECEDING entry whose date we know. // An id ahead of every known date is newer than everything we have and gets // NEWER_THAN_ANYTHING. // // Requires at least one dated anchor: with no anchor we know nothing about the // channel's timeline, and handing every id the top sentinel would let a freshly // added 9,000-video channel monopolize the head of the queue. Anchorless // channels fall through to layer 4 instead. // // Exported for the unit test: its two inputs (a newest-first playlist and a // partial date map) are exactly what a fixture can supply, whereas driving it // through buildRecencyKeys would require standing up an LMDB. export function interpolateFromPlaylist( playlistIds: ReadonlyArray, dates: ReadonlyMap, // Mutated: an id keyed here is removed, so the caller can stop early once // every missing id has been estimated. wanted: Set, out: Map, ): void { let anchored = false; for (const id of playlistIds) { if (dates.has(id)) { anchored = true; break; } } if (!anchored) return; let last: string | null = null; for (const id of playlistIds) { const known = dates.get(id); if (known) { last = known; continue; } if (!wanted.delete(id)) continue; out.set(id, { key: last ?? NEWER_THAN_ANYTHING, estimated: true, }); } } // --- The public builder ------------------------------------------------------ export type BuildRecencyKeysArgs = { paths: Paths; // The channels whose work is in play — the runner's own channel-meta list. // The config, when the caller holds it, is what tells layer 2 which channels // are on another drive (their reads go through the stall watchdog). meta: ReadonlyArray<{ slug: string; config?: Pick | null; }>; // Every video id the caller might sort. Ids outside this set are not keyed. candidateIds: ReadonlySet; // videoId -> owning channel slug. Required for the tail-read layer, which has // to know which channel dir a video lives in. Ids absent from this map skip // layer 2. The runner derives it from the same projection it builds the // candidate set from, so the two can never disagree. owner?: ReadonlyMap; // Read playlists for layer 2. Auto-transcribe only ever considers already- // downloaded videos, which the index already dates, so it passes false and // skips the reads entirely. Default true. interpolate?: boolean; // Bypass the TTL cache (tests, and the offline sanity script). fresh?: boolean; }; // What the expensive half of a build produces: everything derived from disk, // independent of which ids the caller happens to be asking about. Cached, so a // runner ticking every few seconds pays for the index scan and the playlist // reads at most once per TTL. type RecencySources = { at: number; // Every indexed video's date, merged across channels (an id belongs to one // channel, so the merge is lossless in practice; first writer wins). dates: Map; // Per channel, kept separately because interpolation needs a channel's own // anchors, not the corpus's. datesBySlug: Map>; // slug -> playlist ids in listing (newest-first) order. Only populated when // the caller asked to interpolate. playlists: Map; interpolated: boolean; }; // Matches the runner's own CHANNEL_LIST_TTL_MS: the two caches expire together, // so a channel added mid-run becomes visible to the meta list and to its dates // on the same tick rather than one lagging the other. const RECENCY_TTL_MS = 30_000; let cache: RecencySources | null = null; let cacheKey = ""; // Exported for tests, which need a clean slate between fixtures. export function clearRecencyCache(): void { cache = null; cacheKey = ""; tailMemo.clear(); tailCapLogged = false; } async function loadSources( paths: Paths, slugs: ReadonlyArray, interpolate: boolean, fresh: boolean, ): Promise { const key = `${paths.lmdbPath}\u0000${slugs.join(",")}`; const now = Date.now(); if ( !fresh && cache && cacheKey === key && now - cache.at < RECENCY_TTL_MS && // A cached scan that skipped playlists cannot serve a request that needs // them; the reverse is fine. (!interpolate || cache.interpolated) ) { return cache; } const index = openUploadDateIndex(paths, slugs); const dates = new Map(); const datesBySlug = index.datesBySlug; try { for (const perChannel of datesBySlug.values()) { for (const [id, date] of perChannel) if (!dates.has(id)) dates.set(id, date); } } finally { await index.close().catch(() => {}); } const playlists = new Map(); if (interpolate) { for (const slug of slugs) { // Require at least one indexed video. Layer 2 can add anchors later, but a // channel with no transcript at all is one whose whole playlist would fall // to layer 4 anyway — so skip the read rather than parse it and discard it. if (!datesBySlug.get(slug)?.size) continue; const ids = await readPlaylistIds(paths, slug); if (ids.length > 0) playlists.set(slug, ids); } } const sources: RecencySources = { at: now, dates, datesBySlug, playlists, interpolated: interpolate, }; if (!fresh) { cache = sources; cacheKey = key; } return sources; } export async function buildRecencyKeys({ paths, meta, candidateIds, owner, interpolate = true, fresh = false, }: BuildRecencyKeysArgs): Promise> { const out = new Map(); if (candidateIds.size === 0) return out; const slugs = meta.map((m) => m.slug); const sources = await loadSources(paths, slugs, interpolate, fresh); // Layer 1: the index scan. // // TWO KEY SPACES MEET HERE, and they are not the same space. The index is // keyed by the METADATA id (buildIndex writes `summary.id`, i.e. the id the // platform reports), while every candidate id reaching this function is a // DIRECTORY NAME — snapshot.backfill[op].ids, the snapshot buckets and a // readdir of data/ all name directories. They agree for ~85.5% of this // corpus and diverge for the rest: `the-quartering-rumble/data/v1007ay` has // metadata id `vxe1ae`, because Rumble's dir is the URL slug and its // metadata id is the embed id. // // A flat corpus-wide `dates` map turns that divergence into a WRONG DATE // rather than a missing one: channel A's directory name can collide with // channel B's metadata id and silently inherit B's upload date. So when the // caller supplies `owner` — every scheduling path does — the lookup is // scoped to the owning channel's own map, which contains that channel's ids // only. A dir name that is not a metadata id in its own channel then falls // through to layer 2, which reads the date out of the directory it names and // is therefore correct by construction. // // Without `owner` there is nothing to scope by, so the flat map stands; that // path is tests and the offline sanity script, not a runner. const missing = new Set(); for (const id of candidateIds) { const slug = owner?.get(id); const date = slug ? sources.datesBySlug.get(slug)?.get(id) : sources.dates.get(id); if (date) out.set(id, { key: date, estimated: false }); else missing.add(id); } // Layer 2: on-disk metadata. This is the one that actually orders // auto-transcribe, whose entire candidate set is by definition absent from the // transcript index. if (owner && missing.size > 0) { const drives = new Map(); for (const m of meta) { // The retired `dataDir` only: `mediaDir` holds no metadata. const dir = m.config?.dataDir?.trim(); if (dir) drives.set(m.slug, dir); } await datesFromMetadata(paths, owner, missing, out, drives); } // Layer 3: playlist interpolation, for videos with nothing on disk at all. // Anchors are the index dates PLUS anything layer 2 just read — a downloaded // -but-untranscribed video is invisible to the index yet makes a perfectly // good anchor, and folding it in tightens every estimate around it. if (interpolate && missing.size > 0) { for (const [slug, playlistIds] of sources.playlists) { if (missing.size === 0) break; const anchors = new Map(sources.datesBySlug.get(slug) ?? []); if (owner) { for (const [id, key] of out) { if (!key.estimated && key.key && owner.get(id) === slug) { anchors.set(id, key.key); } } } interpolateFromPlaylist(playlistIds, anchors, missing, out); } } // Layer 4: the YYYYMMDD_ dir-name prefix, else "" (oldest). Reuses the same // helper the on-disk keep-window uses, so an id keys identically here and // there rather than growing a second set of fallback rules. for (const id of candidateIds) { if (out.has(id)) continue; out.set(id, { key: uploadKeyFor(undefined, id), estimated: false }); } return out; } // --- Comparator --------------------------------------------------------------- // Total order over video ids for a given policy `order`. Returns null for // "listed", which the caller passes straight through to buildPendingByLeaf as an // absent comparator — so the historical order is reproduced by not sorting at // all, not by sorting with an identity comparator. // // AND NULL FOR "cheapest", which is NOT a recency order at all: it is keyed by // DURATION, and the comparator for it is supplied by the lane's runner. Falling // through to the date sort below would have made it silently mean oldest-first // (the `olderFirst = -1` arm), which is a real ordering an operator did not ask // for — so an order this module cannot answer returns "don't sort" and the lane // keeps today's behaviour until its runner supplies one. // // Equal keys compare 0 deliberately: Array#sort is stable, so videos sharing an // upload date keep their "listed" order (playlist order for undownloadedIds, // id order elsewhere) instead of being shuffled by an arbitrary tiebreak. export function makeRecencyComparator( keys: ReadonlyMap, order: AutoQueueOrder, ): ((a: string, b: string) => number) | null { if (order === "listed" || order === "cheapest") return null; // Keys are YYYYMMDD strings, so lexicographic order IS chronological order. // "newest" therefore sorts DESCENDING: a smaller (older) key must come later, // which is a positive comparator result. const olderFirst = order === "newest" ? 1 : -1; return (a, b) => { const ka = keys.get(a)?.key ?? UNKNOWN; const kb = keys.get(b)?.key ?? UNKNOWN; if (ka === kb) return 0; return ka < kb ? olderFirst : -olderFirst; }; }