import { tierChannelMedia } from "../lib/mediaTier-server"; import path from "node:path"; import { mkdir, readdir, readFile } from "node:fs/promises"; import { writeFileAtomic } from "../lib/jsonFile-server"; import { execa } from "execa"; import pLimit from "p-limit"; import { readArchive } from "../lib/archive"; import { type AudioFormat, type ChannelConfig, type ChannelHandling, } from "../lib/channelConfig"; import { getSettings } from "../lib/settings"; import { patchChannelConfig } from "../controller/channels"; import { diskGate } from "../lib/diskSpace"; import { detectPlatform } from "../lib/platform"; import { channelExtraArgs, channelPaceSeconds, channelPlatform, pacedPlatformArgs, platformMinGapSeconds, staticSleepRequestsSeconds, } from "./channelArgs"; import { isRealAudioFile } from "../lib/videoStatus"; import { readVttProvenance } from "../lib/subtitleProvenance"; import { extractVideoId } from "../lib/videoId"; import type { Paths } from "../lib/paths"; import { EXCLUDED_FROM_DOWNLOAD, classifyDownloadFailure, type Availability, } from "../lib/availability"; import type { DownloadOutcomeRecord } from "../lib/downloadOutcome"; import { alwaysCookies, resolveCookiePolicy, type ResolvedCookiePolicy, } from "../lib/cookiePolicy"; import { resolveEffectiveAvailability } from "../lib/availability-server"; import { withMetadataHistory } from "../lib/metadataHistory-server"; import { settledByTitleFilterIds, upsertMetadataScan, type MetadataScanEntry, } from "../controller/metadataScanStore"; import { backfillAvailabilityFromMetadata } from "../controller/backfillAvailability"; import { runAvailabilityCheck } from "../controller/checkAvailability"; import { writeMaybeMissing } from "../controller/maybeMissingStore"; import { mergeRoster, mergeRosterFile, observationsFromUrls, recordSweep, seedRosterIfAbsent, writeRoster, type Roster, } from "../controller/rosterStore"; import { acceptListing, type ListingDecision } from "../controller/acceptListing"; import { deriveChannelSets } from "../controller/channelSets"; import { resolveShardItems } from "../controller/shard"; import { isFullSweepDue, resolveFullSweepIntervalMinutes, } from "../jobs/deepSync"; import { clearSubtitleDeferral, platformCooldownRemainingMs, readSubtitleDeferrals, recordSubtitleDeferral, } from "../jobs/downloadBackoff"; import { SUBTITLE_HOLD_AFTER, downloadGapMs, isSubtitleDeferred, } from "../jobs/platformBackoff"; import { hasSubtitleRateLimit } from "../lib/availability"; import { computeKeepWindow, type KeepWindow } from "../controller/keptVideos"; import { downloadOneManaged } from "./downloadOneManaged"; import { STDERR_TAIL_BYTES } from "./runOneYtdlp"; import { resolveDownloadFormatPreset, resolveDownloadFormatSelector, type DownloadFormatPreset, } from "./downloadFormat"; import type { TaskTracker } from "../jobs/taskHooks"; import type { JobProgress } from "../jobs/registry"; export type YtdlpMode = | "store-playlist" | "download-from-playlist" | "download-missing" | "download-missing-subs" | "download-one-audio" | "retry-bucket" | "sync"; export type RunYtdlpOpts = { channelSlug: string; mode: YtdlpMode; channelConfig: ChannelConfig; paths: Paths; onLog: (s: string) => void; signal: AbortSignal; // Soft-cancel: when aborted, finish the current download but start no new // ones. Separate from the hard `signal`. drainSignal?: AbortSignal; // When provided, each per-video download is tracked as a per-operation task // with its own parsed progress bar on the Active Jobs screen. tracker?: TaskTracker; // download-missing only: drop --download-archive so videos already listed // in the archive still get re-fetched (recovery from a wrong/stale archive). ignoreArchive?: boolean; // download-from-playlist / download-missing: when false, omit // --abort-on-error so yt-dlp continues past per-video failures. Defaults // to true to preserve historical behavior. abortOnError?: boolean; // download-missing only: split the playlist N ways and only fetch the Mth // shard. Saved/resumed via channels//shard-download-missing.json. shardTotal?: number; shardIndex?: number; // download-missing only: when set, resolve + persist the shard slice and // return immediately without downloading. Backs the "Save shard" button. saveShardOnly?: boolean; // download-missing only: when sharding is active, the run refines the job's // progress target to its own slice (progressBaseline + shard items left to // fetch) so the Active Jobs bar measures the shard, not the whole channel. setProgress?: (snap: JobProgress) => void; progressBaseline?: number; // download-one-audio only: full webpage URL of the single video to fetch. singleVideoUrl?: string; // Override channelConfig.audioFormat for this run. Honored by download-one-audio // and by the managed download paths (passed through to downloadOneManaged). audioFormatOverride?: AudioFormat; // Per-run override of the yt-dlp `-f` download format (see // common/ytdlp/downloadFormat.ts). Honored by download-one-audio and the // managed download paths (passed through to downloadOneManaged). Beats the // per-channel and global defaults. downloadFormatOverride?: DownloadFormatPreset; // Per-run persistence overrides (Phase 2), threaded to downloadOneManaged for // every managed download in this run. keepSourceVideoOverride forces keep // (true) / discard (false); extractImmediately forces extract-now + discard // even for a video the keep-latest rule would persist. keepSourceVideoOverride?: boolean; extractImmediately?: boolean; // download-one-audio only: appended after configArgs, before the URL. extraYtdlpArgs?: string[]; // sync only: force this run to be a FULL SWEEP regardless of the configured // cadence (see common/jobs/deepSync.ts). Undefined leaves the decision to the // cadence gate; false forces the cheap paged walk. Backs the explicit // "Full sweep" action and its job replay. forceFullSweep?: boolean; // retry-bucket only: video IDs to limit the run to (matched against the // saved playlist by extracted ID). IDs not present in the playlist are // reported and skipped. bucketIds?: ReadonlyArray; // retry-bucket only (the "Needs cookies" bucket button): force cookie mode // "always" for this run, so every invocation carries the configured cookie // value and the defer-mode batch exclusion is bypassed. Warns (and behaves // like a plain run) when no cookie value is configured. forceCookies?: boolean; // retry-bucket only: override channelConfig.handling for this run without // mutating the channel config on disk. Useful for retrying old "youtube" // videos as "transcribe". handlingOverride?: ChannelHandling; // retry-bucket only (the "YouTube auto-captions only" bucket): these videos // ALREADY have an English VTT, which destinationExists() normally reads as // "already downloaded" — the run would prefilter every id away. With this set, // an English VTT that still sniffs as YouTube ASR no longer counts as a // destination, so the audio download proceeds. A manual (or unclassifiable) // caption track still suppresses the download, so a human transcript is never // the reason we re-fetch. Paired with handlingOverride: "transcribe" by the // caller, since a youtube-handling download would fetch subs, not audio. replaceAutoSubs?: boolean; // Called when a per-video download fails with a rate-limit (HTTP 429) or // network error, so the caller can record the SHARED per-platform cooldown // (see common/jobs/downloadBackoff.ts). This entangles manual sync/download // with the auto-download runner: a 429 hit by either path backs off the other. // Kept as a callback so this module stays decoupled from the auto-queue state. onPlatformBackoff?: ( failureClass: "rate_limit" | "network", ) => void | Promise; // Called once when the run finished — not thrown, not cancelled — with no // platform backoff recorded during it: the source answered cleanly, so a // manual run settles the platform the way a clean lane probe does (its // backoff and any hold clear; release 17 review H2). Returns a line for the // log, or null. onPlatformClean?: () => Promise; // Internal: how many requests this run made of the source (a listing, a // per-video download, a raw spawn). Set by runYtdlpMode and shared through // every `{...opts}` copy by reference; a run that asked nothing is not clean // (release 17 re-review R3). requestCounter?: { n: number }; }; // The batch modes whose child writes media into `data//` — yt-dlp's own // `-x` output, the batch download loops — and so end with the media tier's // hook. Not store-playlist (writes `playlist`) and not download-missing-subs // (writes subtitles, which are text). const MEDIA_WRITING_MODES: ReadonlySet = new Set([ "download-from-playlist", "download-missing", "download-one-audio", "retry-bucket", "sync", ]); // A video dir's mtime is compared with the run's start; a filesystem with // coarse timestamps rounds down, so the window opens this much earlier. const TIER_SINCE_SLACK_MS = 2_000; export async function runYtdlp(opts: RunYtdlpOpts): Promise { if (!opts.channelConfig.url) { throw new Error("Channel has no `url` configured"); } // A shard save computes a slice and writes it; it downloads nothing. if (!MEDIA_WRITING_MODES.has(opts.mode) || opts.saveShardOnly) { return runYtdlpMode(opts); } // THE MEDIA TIER'S HOOK FOR A BATCH (release 17): after the child returns — // done, failed or cancelled, whatever it finished — every video dir this run // touched has its media moved into channels//media when the channel // has one. The per-video managed downloads tier as they go // (downloadOneManaged); this catches what a raw yt-dlp run finalised itself. // Never throws, and costs one lstat on a classic channel. const startedAt = Date.now(); try { await runYtdlpMode(opts); } finally { await tierChannelMedia(opts.paths, opts.channelSlug, { since: startedAt - TIER_SINCE_SLACK_MS, onLog: opts.onLog, }); } } // A CLEAN RUN SETTLES THE PLATFORM (release 17 review H2): with no rate limit // or network failure recorded while it ran, the run is the probe a held // platform waits for — the lane may be off, or have nothing pending on it. async function runYtdlpMode(opts: RunYtdlpOpts): Promise { if (!opts.onPlatformClean) return runYtdlpModeInner(opts); let backedOff = false; const inner = opts.onPlatformBackoff; const requestCounter = { n: 0 }; await runYtdlpModeInner({ ...opts, requestCounter, onPlatformBackoff: async (failureClass) => { backedOff = true; await inner?.(failureClass); }, }); // A run that asked the source nothing (download-missing with nothing to // fetch, download-missing-subs with every video deferred) says nothing about // the platform, and must not lift a due hold. if (backedOff || opts.signal.aborted || requestCounter.n === 0) return; try { const line = await opts.onPlatformClean(); if (line) opts.onLog(line); } catch { /* shared-state write is best-effort */ } } async function runYtdlpModeInner(opts: RunYtdlpOpts): Promise { switch (opts.mode) { case "store-playlist": await storePlaylist(opts); return; case "download-from-playlist": await downloadFromPlaylist(opts); return; case "download-missing": await downloadMissing(opts); return; case "download-missing-subs": await downloadMissingSubs(opts); return; case "download-one-audio": await downloadOneAudio(opts); return; case "retry-bucket": await retryBucket(opts); return; case "sync": await sync(opts); return; } } async function retryBucket(opts: RunYtdlpOpts): Promise { if (!opts.bucketIds || opts.bucketIds.length === 0) { opts.onLog("retry-bucket: no video IDs supplied. Nothing to do.\n"); return; } await downloadPlaylistManaged(opts, { prefilter: "destination-exists", idAllowList: opts.bucketIds, handlingOverride: opts.handlingOverride, }); } function channelRoot(opts: RunYtdlpOpts): string { return path.join(opts.paths.channelsDir, opts.channelSlug); } function abortableSleep(ms: number, signal: AbortSignal): Promise { if (signal.aborted) return Promise.resolve(); return new Promise((resolve) => { const t = setTimeout(() => { signal.removeEventListener("abort", onAbort); resolve(); }, ms); const onAbort = () => { clearTimeout(t); resolve(); }; signal.addEventListener("abort", onAbort, { once: true }); }); } // Cookies, then the platform's fixed args, then the channel's own — the one // builder in channelArgs.ts. function configArgs(config: ChannelConfig, cookies?: string): string[] { return channelExtraArgs(config, cookies); } // The run-level cookie policy: settings + channel overrides, with the // retry-bucket forceCookies flag overriding the mode to "always" (that's the // "Download with cookies" button). Logs a warning when cookies are forced but // no value is configured anywhere — the run then proceeds cookie-less. function resolveRunCookiePolicy( opts: RunYtdlpOpts, channelConfig: ChannelConfig, ): ResolvedCookiePolicy { const policy = resolveCookiePolicy(getSettings(), channelConfig); if (!opts.forceCookies) return policy; if (!policy.cookies) { opts.onLog( "Force-cookies run requested, but no cookies-from-browser value is configured (global or channel). Proceeding without cookies.\n", ); } return { ...policy, mode: "always" }; } // Defer-mode batch exclusion: a video whose effective availability is // needs_auth is skipped by batch runs when the resolved cookie mode is // "defer" (it lands in the snapshot's Needs-cookies bucket instead of being // re-attempted on every sync/download-missing). members_only/private are // already covered by EXCLUDED_FROM_DOWNLOAD. forceCookies runs resolve to // mode "always" and therefore bypass this. export async function isDeferredAuthExcluded( videoDir: string, policy: ResolvedCookiePolicy, ): Promise { if (policy.mode !== "defer") return false; const cls = await resolveEffectiveAvailability(videoDir); return cls === "needs_auth"; } const OUTPUT_ARGS: string[] = [ "-o", "data/%(id)s/audio.%(ext)s", "-o", "subtitle:data/%(id)s/transcript", "-o", "infojson:data/%(id)s/metadata", // Suppress per-playlist info.json/description files that otherwise land at // the channel root and trip up the index build with malformed entries. "--no-write-playlist-metafiles", ]; // yt-dlp's `%(id)s` is the *extractor* id, which only matches our canonical, // URL-derived id on YouTube. On Twitch (`v` prefix), Rumble, and Odysee it // diverges, splitting a video's bytes from the app's `data//` // dir. For single-video spawns we know the URL up front, so we pin the output // path literally to the canonical id and stop depending on the extractor id. // Falls back to %(id)s for URLs whose canonical id is missing or not a safe // directory name (the post-download reconcile pass cleans those up). // `mediaName` overrides the base name of the main media output (default // "audio"). App-side extraction downloads the source container as // `source-media.` so it's never confused with an extractable audio file; // every other caller keeps the historical `audio.`. export function outputArgsForUrl( url: string, opts: { mediaName?: string } = {}, ): string[] { const mediaName = opts.mediaName ?? "audio"; const id = dataDirIdForUrl(url); if (id) { return [ "-o", `data/${id}/${mediaName}.%(ext)s`, "-o", `subtitle:data/${id}/transcript`, "-o", `infojson:data/${id}/metadata`, "--no-write-playlist-metafiles", ]; } if (mediaName !== "audio") { return [ "-o", `data/%(id)s/${mediaName}.%(ext)s`, "-o", "subtitle:data/%(id)s/transcript", "-o", "infojson:data/%(id)s/metadata", "--no-write-playlist-metafiles", ]; } return OUTPUT_ARGS; } // The data// directory outputArgsForUrl pins a URL's output to, or null // when it falls back to yt-dlp's %(id)s (no canonical id, or one that is not a // safe directory name) — i.e. the directory is only known afterwards. export function dataDirIdForUrl(url: string): string | null { const id = extractVideoId(url); return id && /^[\w.-]+$/.test(id) && id !== "." && id !== ".." ? id : null; } // Thrown by enumeratePlaylistUrls when a listing is rate-limited part-way. // // A full enumeration that the platform rate-limited part-way through. What it // printed is a PREFIX of the listing, not a listing: acting on it as one would // flag every entry past the cut-off as missing, and failing the sync on it // throws away a download walk that would have worked. Thrown only for a // `rate_limit`-classified exit (HTTP 429 and friends, classifyDownloadFailure); // every other non-zero exit still throws the plain error. // // `platform` is the cooldown key the Sync gate and the auto-download runner // share (`detectPlatform(url) ?? "unknown"`); `pagesReached` is the last // `Downloading page N` the extractor logged (null for an extractor that does // not log pages); `count` is how many entries were printed before the cut. export class EnumerationIncompleteError extends Error { readonly platform: string; readonly pagesReached: number | null; readonly count: number; constructor( exitCode: number | undefined, platform: string, pagesReached: number | null, count: number, ) { super( `yt-dlp exited with code ${exitCode} (rate-limited${ pagesReached !== null ? ` at page ${pagesReached}` : "" } of the listing, ${count} entries)`, ); this.name = "EnumerationIncompleteError"; this.platform = platform; this.pagesReached = pagesReached; this.count = count; } } // The last `Downloading page N` an extractor logged — yt-dlp's paged channel // extractors print `[RumbleChannel] : Downloading page 155` per page. export function lastListingPage(stderr: string): number | null { let last: number | null = null; for (const m of stderr.matchAll(/Downloading page (\d+)/g)) { last = Number(m[1]); } return last; } // The per-platform cooldown key: the one the Sync gate (pipelineActions) and // the auto-download runner both use, so a sweep's 429 pauses exactly what // theirs would. function cooldownPlatformKey(config: ChannelConfig): string { return detectPlatform(config.url) ?? "unknown"; } // Enumerate a channel's video URLs via `--flat-playlist --print url` (metadata // only, no downloads). Pass `range` to fetch a single newest-first page via // `-I start:end`; sync uses this to walk the channel incrementally. async function enumeratePlaylistUrls( opts: RunYtdlpOpts, root: string, range?: { start: number; end: number }, ): Promise { const args: string[] = [ "--flat-playlist", "--skip-download", "--print", "url", ]; if (range) { args.push("--lazy-playlist", "-I", `${range.start}:${range.end}`); } // Cookie mode "always" covers enumeration too (a members-only or otherwise // gated channel may not even list without cookies). args.push( ...configArgs( opts.channelConfig, alwaysCookies(resolveRunCookiePolicy(opts, opts.channelConfig)), ), opts.channelConfig.url!, ); opts.onLog(`$ ${opts.paths.ytdlpBin} ${args.join(" ")}\n`); if (opts.requestCounter) opts.requestCounter.n++; const child = execa(opts.paths.ytdlpBin, args, { cwd: root, cancelSignal: opts.signal, all: false, buffer: true, reject: false, }); child.stderr?.on("data", (c: Buffer) => opts.onLog(c.toString("utf8"))); const result = await child; const urls = String(result.stdout ?? "") .split("\n") .map((s) => s.trim()) .filter(Boolean); // yt-dlp exit code convention: 101 = "break-on-existing" / "max-downloads" // (clean stop, not an error). Treat it the same as 0. A non-zero/101 exit // while the run was cancelled is the cancel itself, not a failure. if ( result.exitCode !== 0 && result.exitCode !== 101 && !opts.signal.aborted ) { const stderr = String(result.stderr ?? ""); if (classifyDownloadFailure(stderr, undefined) === "rate_limit") { throw new EnumerationIncompleteError( result.exitCode, cooldownPlatformKey(opts.channelConfig), lastListingPage(stderr), urls.length, ); } throw new Error(`yt-dlp exited with code ${result.exitCode}`); } if (result.exitCode === 101) { opts.onLog(`yt-dlp stopped on existing entry (exit 101).\n`); } return urls; } // Thin exported wrapper over the module-private enumeratePlaylistUrls for // callers (e.g. the quick availability check) that only need a one-shot fresh // flat-playlist fetch and shouldn't have to fabricate a full RunYtdlpOpts. // enumeratePlaylistUrls only reads channelConfig, paths.ytdlpBin, onLog and // signal — never mode/tracker — so the cast is safe. export async function fetchFlatPlaylistUrls(opts: { channelConfig: ChannelConfig; paths: Paths; channelSlug: string; onLog: (s: string) => void; signal: AbortSignal; }): Promise { const root = path.join(opts.paths.channelsDir, opts.channelSlug); return enumeratePlaylistUrls( { ...opts, mode: "store-playlist" } as RunYtdlpOpts, root, ); } // Lightweight metadata probe for the URL-first new-channel flow. Fetches a // SINGLE flat-playlist entry (`--playlist-items 1`, no media) and prints the // best available name field, so the create form can pre-fill the display name // from a pasted URL. Crucially it also proves yt-dlp *can* handle the URL — a // host the app's detectPlatform doesn't recognize but yt-dlp does still probes // successfully, which is how an unknown-platform channel gets created (routed to // its own `platform:` serial queue by downloadQueueKey). Returns what it // could read; a null name just means yt-dlp printed nothing usable. export async function probeChannelMeta(opts: { url: string; paths: Paths; onLog?: (s: string) => void; signal?: AbortSignal; timeoutMs?: number; }): Promise<{ name: string | null; entryCount: number | null }> { const { url, paths } = opts; const log = opts.onLog ?? (() => {}); // yt-dlp's field-fallback template prints the first non-null of the listed // fields, so a channel URL yields "Veritasium" (channel), a bare playlist its // title, etc. `--flat-playlist --playlist-items 1` keeps it to one cheap // metadata read with no download. const args = [ "--flat-playlist", "--skip-download", "--playlist-items", "1", "--print", "%(channel,uploader,playlist_title,playlist,uploader_id,title)s", // No channel config exists yet; the platform's fixed args still apply // (a Rumble probe 403s without `--impersonate`), at its current pace. ...pacedPlatformArgs(detectPlatform(url)), url, ]; log(`$ ${paths.ytdlpBin} ${args.join(" ")}\n`); const child = execa(paths.ytdlpBin, args, { cancelSignal: opts.signal, all: false, buffer: true, reject: false, timeout: opts.timeoutMs ?? 60_000, }); child.stderr?.on("data", (c: Buffer) => log(c.toString("utf8"))); const result = await child; if ( result.exitCode !== 0 && result.exitCode !== 101 && !opts.signal?.aborted ) { throw new Error( `yt-dlp exited with code ${result.exitCode}${ result.stderr ? `: ${String(result.stderr).trim().slice(-400)}` : "" }`, ); } const lines = String(result.stdout ?? "") .split("\n") .map((s) => s.trim()) .filter((s) => s && s !== "NA"); return { name: lines[0] ?? null, entryCount: null }; } // Atomically replace the channel's stored `playlist` file (tmp + rename, so a // crashed write never leaves a truncated list behind). Shared by store-playlist // and the sync full sweep, which both enumerate the whole channel and are the // only writers of this file. async function writePlaylistFile( root: string, urls: ReadonlyArray, onLog: (s: string) => void, ): Promise { const playlistPath = path.join(root, "playlist"); await writeFileAtomic(playlistPath, urls.join("\n") + (urls.length ? "\n" : "")); onLog(`Wrote ${urls.length} URLs to ${playlistPath}\n`); } // Distinct canonical ids in the channel's currently-stored `playlist` — the // size of the last listing we accepted, and the reference the shrink guard // measures a fresh enumeration against. Missing file = no reference. async function storedListingCount(root: string): Promise { try { const raw = await readFile(path.join(root, "playlist"), "utf8"); const urls = raw.split("\n").map((s) => s.trim()).filter(Boolean); return new Set(observationsFromUrls(urls).map((o) => o.id)).size; } catch { return 0; } } // Put a full enumeration through the roster and the shrink guard. BOTH // full-listing writers go through here — store-playlist and the sync full // sweep — so neither can rewrite `playlist` from a listing the other would have // refused. store-playlist overwriting unconditionally was the most direct // instance of the reported bug. // // Note the order: the roster is seeded from the pre-existing state BEFORE // anything is overwritten, and the fresh observations are merged in // UNCONDITIONALLY — even for a rejected listing. Merging is additive, so a bad // fetch can only ever add; an entry we decline to record is an entry we can // lose, which is the failure this whole part exists to prevent. async function acceptEnumeration( opts: RunYtdlpOpts, root: string, urls: ReadonlyArray, ): Promise<{ decision: ListingDecision; roster: Roster; listedIds: Set; now: string; }> { const now = new Date().toISOString(); const observed = observationsFromUrls(urls); const listedIds = new Set(observed.map((o) => o.id)); const [previousListedCount, seeded] = await Promise.all([ storedListingCount(root), seedRosterIfAbsent(opts.paths, opts.channelSlug, { now, onLog: opts.onLog, }), ]); const decision = acceptListing({ listedCount: listedIds.size, previousListedCount, lastSweep: seeded.lastSweep, shrinkGuardPercent: getSettings().syncScheduler.fullSweepShrinkGuardPercent, }); const roster = recordSweep(mergeRoster(seeded, observed, now, "listing"), { at: now, listedCount: listedIds.size, verdict: decision.verdict, }); await writeRoster(opts.paths, opts.channelSlug, roster); return { decision, roster, listedIds, now }; } async function storePlaylist(opts: RunYtdlpOpts): Promise { const root = channelRoot(opts); await mkdir(root, { recursive: true }); const urls = await enumeratePlaylistUrls(opts, root); const { decision } = await acceptEnumeration(opts, root, urls); if (!decision.accept) { opts.onLog(`Store playlist: ${decision.reason}.\n`); return; } await writePlaylistFile(root, urls, opts.onLog); } async function downloadFromPlaylist(opts: RunYtdlpOpts): Promise { await downloadPlaylistManaged(opts, { prefilter: "archive" }); } async function downloadMissing(opts: RunYtdlpOpts): Promise { await downloadPlaylistManaged(opts, { prefilter: "destination-exists" }); } type PrefilterMode = "archive" | "destination-exists"; async function downloadPlaylistManaged( opts: RunYtdlpOpts, cfg: { prefilter: PrefilterMode; idAllowList?: ReadonlyArray; handlingOverride?: ChannelHandling; }, ): Promise { const { prefilter, idAllowList, handlingOverride } = cfg; const effectiveHandling: ChannelHandling = handlingOverride ?? opts.channelConfig.handling; const effectiveChannelConfig: ChannelConfig = handlingOverride && handlingOverride !== opts.channelConfig.handling ? { ...opts.channelConfig, handling: handlingOverride } : opts.channelConfig; const root = channelRoot(opts); const dataDir = path.join(root, "data"); const playlistPath = path.join(root, "playlist"); const archivePath = path.join(root, "archive"); let playlistText: string; try { playlistText = await readFile(playlistPath, "utf8"); } catch { throw new Error( `No saved playlist at ${playlistPath}. Run "Store playlist" first.`, ); } let urls = playlistText .split("\n") .map((s) => s.trim()) .filter(Boolean); // Post-reconcile, a video's on-disk dir is always its canonical id // (extractVideoId of the URL), so no slug→dir indirection is needed here. // The one exception is the *archive* file, which keeps yt-dlp's native // extractor ids — archiveIdForUrl bridges that. if (idAllowList) { const allow = new Set(idAllowList); const before = urls.length; const matched: string[] = []; const matchedIds = new Set(); for (const url of urls) { const dirId = extractVideoId(url); if (!dirId) continue; if (allow.has(dirId)) { matched.push(url); matchedIds.add(dirId); } } const missing = [...allow].filter((id) => !matchedIds.has(id)); opts.onLog( `Retry allowlist: kept ${matched.length} of ${before} playlist URLs (${missing.length} requested IDs not in playlist).\n`, ); if (missing.length > 0) { opts.onLog(` Missing: ${missing.join(", ")}\n`); } urls = matched; } if (handlingOverride) { opts.onLog( `Handling override: using ${handlingOverride} for this run (channel config unchanged).\n`, ); } const tofetch: string[] = []; if (prefilter === "archive") { // Skip URLs whose ID is already recorded in the channel's archive file. const archive = await readArchive(archivePath); let skipped = 0; if (opts.ignoreArchive) { tofetch.push(...urls); } else { for (const url of urls) { const archiveId = await archiveIdForUrl(url, dataDir); if (archiveId && archive.ids.has(archiveId)) { skipped++; continue; } tofetch.push(url); } } opts.onLog( `Prefilter: ${tofetch.length} new, ${skipped} already archived (of ${urls.length} total).\n`, ); } else { // Skip URLs whose destination file (transcript.en.vtt for youtube // handling, audio.* or transcript.json for transcribe handling) is // already on disk. Catches videos missing from a stale archive file. let alreadyComplete = 0; let unidentifiable = 0; for (const url of urls) { const dirId = extractVideoId(url); if (!dirId) { unidentifiable++; tofetch.push(url); continue; } if ( await destinationExists(dataDir, dirId, effectiveHandling, { replaceAutoSubs: opts.replaceAutoSubs, }) ) { alreadyComplete++; continue; } tofetch.push(url); } opts.onLog( `Prefilter: ${tofetch.length} missing destination files, ${alreadyComplete} already complete (of ${urls.length} total)${ unidentifiable ? `, ${unidentifiable} could not be identified` : "" }.\n`, ); } // Exclude videos whose effective availability marks them as permanently // unavailable (members_only, deleted, private). Applies to both prefilter // modes: failed download attempts don't write to the archive, so the // archive-prefilter path would also keep retrying them. Under cookie mode // "defer", needs_auth videos are excluded too — they wait in the snapshot's // Needs-cookies bucket for a manual cookie run instead of being re-attempted // (and re-failing) on every batch. A forceCookies run resolves to mode // "always" and so bypasses the defer exclusion. const runCookiePolicy = resolveRunCookiePolicy(opts, effectiveChannelConfig); // ONE read of the channel-level metadata-scan store for the whole batch. const settledByFilter = await settledByTitleFilterIds( opts.paths, opts.channelSlug, effectiveChannelConfig, ); const excludedCounts = { members_only: 0, deleted: 0, private: 0 }; let deferredAuthCount = 0; let settledByFilterCount = 0; const filteredTofetch: string[] = []; for (const url of tofetch) { const dirId = extractVideoId(url); if (dirId) { // Settled by the channel's download filter, from the set resolved once // above. Checked FIRST, so a filtered channel's batch doesn't run the // metadata prefetch for every non-match on every run just to reach the // same verdict the scan already reached. if (settledByFilter.has(dirId)) { settledByFilterCount++; continue; } const cls = await resolveEffectiveAvailability( path.join(dataDir, dirId), ); if ( cls && (EXCLUDED_FROM_DOWNLOAD as ReadonlyArray).includes(cls) ) { excludedCounts[cls as keyof typeof excludedCounts]++; continue; } if (runCookiePolicy.mode === "defer" && cls === "needs_auth") { deferredAuthCount++; continue; } } filteredTofetch.push(url); } const totalExcluded = excludedCounts.members_only + excludedCounts.deleted + excludedCounts.private; if (totalExcluded > 0) { opts.onLog( `Excluded ${totalExcluded} from this run (members_only=${excludedCounts.members_only}, deleted=${excludedCounts.deleted}, private=${excludedCounts.private}). Clear via Diagnostics > Recheck if a video became public again.\n`, ); } if (settledByFilterCount > 0) { opts.onLog( `Settled by this channel's download filter: ${settledByFilterCount} skipped (already decided). Change the include/exclude patterns in Configure to re-evaluate them.\n`, ); } if (deferredAuthCount > 0) { opts.onLog( `Cookie mode is "defer": needs_auth deferred=${deferredAuthCount} — run the "Needs cookies" bucket to download them with cookies.\n`, ); } tofetch.length = 0; tofetch.push(...filteredTofetch); // Sharding (download-missing): slice the *missing* set into a stable 1/N // partition and persist it here — after the prefilter has determined what's // missing, but BEFORE any download starts. Snapshotting the missing-set // slice to disk is the whole point: the missing set shrinks as videos land, // so on a resume run with the same (total, index) resolveShardItems returns // the saved slice instead of re-slicing the now-smaller set, and the run // picks up exactly the videos this shard still owes. Only the plain // download-missing path is sharded — never the archive-prefilter // (download-from-playlist) path, nor the id-allowlisted retry-bucket run. const shardResult = prefilter === "destination-exists" && !idAllowList ? await resolveShardItems({ paths: opts.paths, slug: opts.channelSlug, op: "download-missing", fullItems: tofetch, totalShards: opts.shardTotal, shardIndex: opts.shardIndex, onLog: (m) => opts.onLog(`${m}\n`), }) : null; const items = shardResult ? shardResult.items : tofetch; // "Save shard" probe: the slice is already persisted by resolveShardItems // above; stop here without downloading anything. if (opts.saveShardOnly) { if (shardResult && shardResult.config) { opts.onLog( `Shard download-missing: save-only — slice ${shardResult.config.shardIndex + 1}/${shardResult.config.totalShards} (${shardResult.config.items.length} item(s)) persisted; not downloading.\n`, ); } else { opts.onLog( `Shard download-missing: save-only requested but no shard inputs were provided; nothing persisted.\n`, ); } return; } // With a shard active, refine the job's progress total to this slice: // baseline (channel-wide downloadCount at run start) + the shard's items // still needing download. `current` is read channel-wide on the Active Jobs // screen, so a baseline-relative, subset-sized range still reaches 100% when // the shard finishes. if ( shardResult && shardResult.source !== "full" && opts.setProgress && opts.progressBaseline !== undefined ) { opts.setProgress({ metric: "downloads", initial: opts.progressBaseline, target: opts.progressBaseline + items.length, }); opts.onLog( `Shard download-missing: progress scoped to ${items.length} shard item(s) needing download.\n`, ); } if (items.length === 0) { opts.onLog("Nothing to fetch.\n"); await touchLastFullDownload(opts); return; } if (opts.ignoreArchive) { opts.onLog( "Ignoring archive: archive entries will NOT be appended for successful downloads in this run.\n", ); } const { okCount, failedCount, skippedCount, processedCount, firstFailure } = await runManagedDownloads( opts, items, effectiveChannelConfig, runCookiePolicy, ); if (firstFailure && !opts.signal.aborted) { throw firstFailure; } opts.onLog( `Managed download complete: ${okCount} succeeded, ${failedCount} failed` + (skippedCount ? `, ${skippedCount} skipped by filter` : "") + ` (${processedCount}/${items.length} processed).\n`, ); await touchLastFullDownload(opts); await safeBackfillAvailability(opts); } export type ManagedDownloadsDeps = { downloadOne?: typeof downloadOneManaged; sleep?: (ms: number, signal: AbortSignal) => Promise; // Where a video's subtitle 429 is recorded (release 17, slice RL). recordSubtitleDeferral?: (videoId: string, channelSlug: string) => Promise; // The platform's current pace (default: channelPaceSeconds). paceSeconds?: () => number; }; // True when the download filter declined the video before any media request: // every attempt it made was the n: 0 metadata prefetch (or its cookie retry) // and it ended `skipped-filtered`. Attempts n >= 1 are the real fetches — the // download attempts and the chat-only pass (a failed chat pass leaves the // status skipped-filtered, and still asked). // // A per_video FAILURE is deliberately NOT here, even one that never got past // the prefetch (release 9 review): a prefetch is still a request, and a run of // "unavailable" answers is how a throttled source looks (see SOFT_BLOCK_STREAK // in metadataScan.ts). Release 9's own reason — YouTube's explicit soft block // ("This content isn't available, try again later") reading as `deleted` → // per_video — is gone since release 10: it is `rate_limit` now (isSoftBlock), // backs the platform off and stops the batch. The pace stays anyway. export function declinedWithoutMediaFetch( outcome: DownloadOutcomeRecord, ): boolean { if (outcome.attempts.some((a) => a.n >= 1)) return false; return outcome.status === "skipped-filtered"; } type ManagedRunResult = { okCount: number; failedCount: number; // Videos declined by an app-level filter (e.g. skip-live). Neither success // nor failure: not archived, never aborts the batch. skippedCount: number; processedCount: number; // Non-null only when abortOnError is on and a non-per-video failure occurred. firstFailure: Error | null; }; // Serialize a list of video URLs through downloadOneManaged. Shared by the // playlist-prefilter modes and `sync`. Serialized (pLimit(1)) to keep logs // readable and avoid hammering the source with parallel requests (which is // what often triggers needs_auth in the first place); the per-video startup // cost is the price of being able to retry each independently. Honors the // hard `signal`, the soft `drainSignal`, the per-channel sleep, and the // per-operation tracker. export async function runManagedDownloads( opts: RunYtdlpOpts, urls: ReadonlyArray, effectiveChannelConfig: ChannelConfig, // Cookie policy for every managed download in this run. Callers that // already resolved it (for their own prefilter) pass it through so the // forceCookies-without-value warning isn't logged twice. cookiePolicyOverride?: ResolvedCookiePolicy, // Test seams: the per-video download and the between-downloads sleep. Every // production caller passes neither. deps: ManagedDownloadsDeps = {}, ): Promise { const downloadOne = deps.downloadOne ?? downloadOneManaged; const sleep = deps.sleep ?? abortableSleep; const settings = getSettings(); const cookiePolicy = cookiePolicyOverride ?? resolveRunCookiePolicy(opts, effectiveChannelConfig); const inlineTranscribeOnFallback = settings.inlineTranscribeOnFallback; const globalSkipLiveDownloads = settings.skipLiveDownloads; // Resolve the download format once per run (override > channel > global), // leaving the platform-aware "auto" expansion to downloadOneManaged where the // per-video extractor is known. const downloadFormatPreset = resolveDownloadFormatPreset({ override: opts.downloadFormatOverride, channel: effectiveChannelConfig.downloadFormat, global: settings.downloadFormat, }); const sleepSeconds = opts.channelConfig.sleepBetweenDownloadsSeconds ?? settings.sleepBetweenDownloadsSeconds; // THE PACE ABOVE ITS BASE joins the gap (release 17, slice RL): after a // rate limit doubled the platform's pace, a batch spaces its videos further // apart too. Read per gap, so a 429 in this batch slows the rest of it. const basePace = staticSleepRequestsSeconds(channelPlatform(effectiveChannelConfig)); // A platform's floor under the gap (BitChute's 60 s, jittered) — see // PLATFORM_MIN_GAP_SECONDS in platformArgs.mjs. const minGap = platformMinGapSeconds(channelPlatform(effectiveChannelConfig)); const paceNow = deps.paceSeconds ?? (() => channelPaceSeconds(effectiveChannelConfig)); const recordSubs = deps.recordSubtitleDeferral ?? ((id: string, slug: string) => recordSubtitleDeferral(id, slug, opts.paths)); // Title-filter rejections from this batch, written to the channel-level // metadata-scan store ONCE at the end. Per-video writes would rewrite the // whole file for every non-matching video — the exact shape of work a // filtered channel produces most of. const filterRejections: Record = {}; // The channel's keep-latest window as a cutoff, computed ONCE per run from the // current on-disk catalog. Each video downloaded below is classified against it // by its own upload date — so the newest videos (not yet on disk) still get // persisted. A keepLatest of 0/undefined yields an inert window. const keepWindow: KeepWindow = await computeKeepWindow({ paths: opts.paths, channelSlug: opts.channelSlug, keepLatest: effectiveChannelConfig.keepLatest ?? 0, }); const limit = pLimit(1); let failedCount = 0; let skippedCount = 0; const abortOnError = opts.abortOnError !== false; let firstFailure: Error | null = null; let processedCount = 0; // Set once free disk falls below the configured floor: like a drain, the // in-flight download finishes but no new ones start, so the batch ends // cleanly (status "done") with partial progress preserved on disk. let lowDiskStopped = false; await Promise.all( urls.map((url) => limit(async () => { if (opts.signal.aborted) return; // Drain (soft-cancel): finish the in-flight download, start no more. if (opts.drainSignal?.aborted) return; if (firstFailure && abortOnError) return; if (lowDiskStopped) return; // THIS channel's volume, not the corpus's: a relocated channel's // downloads land on the platter through the data/ symlink, so the SSD // being full is not a reason to stop them (and vice versa). const gate = await diskGate(opts.paths, settings, { dir: path.join(channelRoot(opts), "data"), }); if (!gate.ok) { lowDiskStopped = true; opts.onLog( `Stopping batch: ${gate.message}. Remaining videos skipped.\n`, ); return; } const task = opts.tracker?.start({ id: extractVideoId(url) ?? url, label: extractVideoId(url) ?? url, kind: "download", }); let outcome; try { if (opts.requestCounter) opts.requestCounter.n++; outcome = await downloadOne({ channelSlug: opts.channelSlug, channelConfig: effectiveChannelConfig, paths: opts.paths, videoUrl: url, onLog: task ? task.onLog : opts.onLog, signal: opts.signal, cookiePolicy, appendArchive: !opts.ignoreArchive, inlineTranscribeOnFallback, globalSkipLiveDownloads, keepWindow, keepSourceVideoOverride: opts.keepSourceVideoOverride, extractImmediately: opts.extractImmediately, audioFormatOverride: opts.audioFormatOverride, downloadFormatPreset, onFilterRejected: (id, entry) => { filterRejections[id] = entry; }, }); } finally { task?.end(); } // A subtitle 429 on a download that otherwise succeeded: the video's // subtitles are deferred, the platform is not backed off, the batch // goes on (release 17, slice RL). Best-effort, like the backoff. if (outcome.failureClass === "subs_rate_limit") { try { await recordSubs(outcome.videoId, opts.channelSlug); } catch { /* shared-state write is best-effort */ } } if (outcome.status === "skipped-filtered") { skippedCount++; } else if (outcome.status === "chat-only") { // The media was skipped and the live chat was fetched. Counted as a // SKIP, because the batch's `okCount` means "videos downloaded" and // this one deliberately was not — the chat is a different artifact // and the report's own chatOnly bucket is where it is counted. skippedCount++; } else if (outcome.status === "corrupt-full-source") { // Completed but stayed malformed after one re-download: terminal and // kept on disk, but not a usable download. Count as skipped (not ok, // not a retryable failure) so the batch summary stays honest. skippedCount++; } else if (outcome.status === "failed") { failedCount++; // downloadOneManaged classifies against the full stderr tail; fall // back to the truncated attempt error only if it's somehow absent. const lastAttempt = outcome.attempts[outcome.attempts.length - 1]; const failureClass = outcome.failureClass ?? classifyDownloadFailure( lastAttempt?.error ?? "", lastAttempt?.availabilityClass, ); // A rate-limit / network failure pauses the WHOLE platform via the // shared cooldown, so the auto-download runner (and a later manual // sync) back off too. Best-effort: never let bookkeeping break the // download loop. Fires regardless of abortOnError / firstFailure. if (failureClass === "rate_limit" || failureClass === "network") { try { await opts.onPlatformBackoff?.(failureClass); } catch { /* shared-state write is best-effort */ } } // Per-video failures (private, members-only, age-gated, deleted) are // properties of that video, not signals the batch should stop. Only // abort on rate-limit / network / unclassified-fatal. if (abortOnError && !firstFailure && failureClass !== "per_video") { firstFailure = new Error( `Managed download aborted (${failureClass}) for ${url}`, ); } } processedCount++; const isLast = processedCount >= urls.length; const willAbortLoop = firstFailure !== null && abortOnError; // A VIDEO THE FILTER DECLINED DOES NOT PAY FOR A FETCH IT NEVER MADE // (release 9, B2). On 2026-09-25 paramount-tactical's download-missing // slept 193 times for 14 archives: every video the download filter // declined slept 30 s after nothing but its metadata prefetch. Every // other outcome still sleeps — a real fetch, success or failure, and // every failure, per-video ones included (see declinedWithoutMediaFetch). const gapMs = downloadGapMs(sleepSeconds, paceNow(), basePace, { minSeconds: minGap, }); if ( gapMs > 0 && !isLast && !willAbortLoop && !opts.signal.aborted && !declinedWithoutMediaFetch(outcome) ) { opts.onLog(`Sleeping ${gapMs / 1000}s before next download...\n`); await sleep(gapMs, opts.signal); } }), ), ); // ONE write for the whole batch. Deliberately outside the abort checks: a // cancelled run still learned what those videos are called, and throwing that // away would make the next run re-fetch their metadata for nothing. if (Object.keys(filterRejections).length > 0) { try { await upsertMetadataScan( opts.paths, opts.channelSlug, { entries: filterRejections }, new Date().toISOString(), ); } catch (err) { opts.onLog( `Failed to record ${Object.keys(filterRejections).length} metadata scan entries: ${(err as Error).message}\n`, ); } } return { okCount: processedCount - failedCount - skippedCount, failedCount, skippedCount, processedCount, firstFailure, }; } // yt-dlp's --sub-langs supports a comma-separated list with shell-glob style // wildcards (e.g. `en.*`, `all`) and exclusions (`-fr`). Replicate just enough // of that here to decide whether a given track key from a video's metadata // (subtitles + automatic_captions) is one we'd expect on disk after running // yt-dlp with the channel's configured `subLangs`. function compileSubLangMatcher(spec: string): (track: string) => boolean { const include: string[] = []; const exclude: string[] = []; for (const raw of spec.split(",")) { const trimmed = raw.trim(); if (!trimmed) continue; if (trimmed.startsWith("-")) exclude.push(trimmed.slice(1)); else include.push(trimmed); } const matches = (track: string, pattern: string): boolean => { if (pattern === "all") return true; if (pattern.includes("*")) { const re = new RegExp( "^" + pattern .replace(/[.+?^${}()|[\]\\]/g, "\\$&") .replace(/\\\*/g, ".*") + "$", ); return re.test(track); } return pattern === track; }; return (track: string) => { if (exclude.some((p) => matches(track, p))) return false; return include.some((p) => matches(track, p)); }; } const SUB_EXT_PATTERN = "(?:vtt|json|json3|srv1|srv2|srv3)"; function trackOnDisk(entries: string[], track: string): boolean { const escaped = track.replace(/[.+?^${}()|[\]\\]/g, "\\$&"); const re = new RegExp(`^transcript\\.${escaped}\\.${SUB_EXT_PATTERN}$`); return entries.some((e) => re.test(e)); } // Fetch only the missing subtitle tracks for one already-downloaded video. // Uses the literal canonical output path (outputArgsForUrl) so the .vtt files // land in the same data// dir as the rest of the video. // THE BATCH SUBTITLE FETCH, per video. Returns the stderr it printed so the // caller can tell a subtitle 429 (deferred, release 17 slice RL) from any // other failure; throws on any non-zero exit as before. async function downloadSubsForUrl( opts: RunYtdlpOpts, root: string, url: string, subLangs: string, ): Promise { const args: string[] = [ "--ignore-config", "--restrict-filenames", ...outputArgsForUrl(url), "--write-auto-subs", "--write-subs", "--sub-langs", subLangs, "--skip-download", "--no-write-info-json", "--no-overwrites", ...configArgs( opts.channelConfig, alwaysCookies(resolveRunCookiePolicy(opts, opts.channelConfig)), ), "--", url, ]; return runChildAndStream(opts, root, args); } async function downloadMissingSubs(opts: RunYtdlpOpts): Promise { const root = channelRoot(opts); const dataDir = path.join(root, "data"); const playlistPath = path.join(root, "playlist"); let playlistText: string; try { playlistText = await readFile(playlistPath, "utf8"); } catch { throw new Error( `No saved playlist at ${playlistPath}. Run "Store playlist" first.`, ); } const urls = playlistText .split("\n") .map((s) => s.trim()) .filter(Boolean); const subLangs = opts.channelConfig.subLangs ?? "en.*,live_chat"; const matchesSubLang = compileSubLangMatcher(subLangs); const tofetch: string[] = []; let upToDate = 0; let unidentifiable = 0; let noMetadata = 0; let noExpectedTracks = 0; // Videos whose subtitles answered 429 recently (release 17, slice RL): left // alone until their deferral lapses — 6 h after a strike, 7 days from the // third. The video page's own download still fetches them. const subDeferrals = await readSubtitleDeferrals(opts.paths).catch( () => ({}) as Awaited>, ); const deferredNow: string[] = []; for (const url of urls) { // Post-reconcile a video lives in data//, so the URL's // canonical id is the dir name directly — no slug→dir indirection. const dirId = extractVideoId(url); if (!dirId) { unidentifiable++; continue; } if (isSubtitleDeferred(subDeferrals, dirId, Date.now())) { deferredNow.push(dirId); continue; } const videoDir = path.join(dataDir, dirId); let metaRaw: string; try { metaRaw = await readFile( path.join(videoDir, "metadata.info.json"), "utf8", ); } catch { // No corresponding data dir yet — skip; download-from-playlist / // download-missing should pull the video and its metadata first. noMetadata++; continue; } let parsedMeta: { subtitles?: Record; automatic_captions?: Record; }; try { parsedMeta = JSON.parse(metaRaw); } catch { noMetadata++; continue; } const advertised = new Set(); for (const k of Object.keys(parsedMeta.subtitles ?? {})) advertised.add(k); for (const k of Object.keys(parsedMeta.automatic_captions ?? {})) { advertised.add(k); } const expected = Array.from(advertised).filter(matchesSubLang); if (expected.length === 0) { noExpectedTracks++; continue; } const entries = await readdir(videoDir).catch(() => [] as string[]); const missing = expected.filter((t) => !trackOnDisk(entries, t)); if (missing.length === 0) { upToDate++; continue; } tofetch.push(url); } opts.onLog( `Prefilter: ${tofetch.length} videos missing subs, ${upToDate} already up to date` + (noExpectedTracks ? `, ${noExpectedTracks} with no tracks matching ${JSON.stringify(subLangs)}` : "") + (noMetadata ? `, ${noMetadata} without metadata (skipped)` : "") + (unidentifiable ? `, ${unidentifiable} unidentifiable URLs` : "") + (deferredNow.length ? `, ${deferredNow.length} deferred after a subtitle rate limit (${deferredNow .map((id) => { const d = (subDeferrals as Record)[id]; return `${id} ${d.count}×${d.count >= SUBTITLE_HOLD_AFTER ? " — left alone" : ""} until ${new Date(d.until).toISOString().slice(0, 16).replace("T", " ")} UTC`; }) .join(", ")})` : "") + `.\n`, ); if (tofetch.length === 0) { opts.onLog("Nothing to fetch.\n"); return; } // One yt-dlp spawn per video (serialized), mirroring the download paths. const abortOnError = opts.abortOnError !== false; const limit = pLimit(1); let processed = 0; let failed = 0; let firstFailure: Error | null = null; // THE GAP BETWEEN TWO VIDEOS (release 17 review L3): the same // sleepBetweenDownloadsSeconds plus the pace above its base the lane and the // batch downloads wait, so a channel whose subtitles are being refused is // not asked once per video back to back. const subsSleepSeconds = opts.channelConfig.sleepBetweenDownloadsSeconds ?? getSettings().sleepBetweenDownloadsSeconds; const subsBasePace = staticSleepRequestsSeconds(channelPlatform(opts.channelConfig)); const subsMinGap = platformMinGapSeconds(channelPlatform(opts.channelConfig)); await Promise.all( tofetch.map((url, index) => limit(async () => { if (opts.signal.aborted || opts.drainSignal?.aborted) return; if (firstFailure) return; if (index > 0) { const gapMs = downloadGapMs( subsSleepSeconds, channelPaceSeconds(opts.channelConfig), subsBasePace, { minSeconds: subsMinGap }, ); if (gapMs > 0) { opts.onLog(`Sleeping ${gapMs / 1000}s before the next video...\n`); await abortableSleep(gapMs, opts.signal); if (opts.signal.aborted || opts.drainSignal?.aborted) return; } } const task = opts.tracker?.start({ id: extractVideoId(url) ?? url, label: extractVideoId(url) ?? url, kind: "download", }); const id = extractVideoId(url); try { const stderr = await downloadSubsForUrl( task ? { ...opts, onLog: task.onLog } : opts, root, url, subLangs, ); // yt-dlp exited clean. A subtitle 429 it reported anyway (a WARNING) // is still a deferral; otherwise the subtitles came down and the // video's deferral, if it had one, is forgotten. if (id) { if (hasSubtitleRateLimit(stderr)) { await recordSubtitleDeferral(id, opts.channelSlug, opts.paths).catch(() => {}); } else { await clearSubtitleDeferral(id, opts.paths).catch(() => {}); } } } catch (err) { failed++; // A SUBTITLE 429 IS PER VIDEO (release 17, slice RL): the video is // deferred and the batch goes on to the next one, abortOnError or // not. Any other failure keeps today's abort. const tail = (err as { stderrTail?: string }).stderrTail ?? ""; if (id && classifyDownloadFailure(tail, undefined) === "subs_rate_limit") { const d = await recordSubtitleDeferral(id, opts.channelSlug, opts.paths).catch(() => null); (task?.onLog ?? opts.onLog)( `Subtitles for ${id} were rate-limited (HTTP 429)${d ? ` — ${d.count}×, deferred` : ""}; continuing with the next video.\n`, ); } else if (abortOnError && !firstFailure) { firstFailure = err as Error; } } finally { task?.end(); } processed++; }), ), ); if (firstFailure && !opts.signal.aborted) throw firstFailure; opts.onLog( `Sub download complete: ${processed - failed} succeeded, ${failed} failed (of ${tofetch.length}).\n`, ); } async function downloadOneAudio(opts: RunYtdlpOpts): Promise { if (!opts.singleVideoUrl) { throw new Error("download-one-audio requires singleVideoUrl"); } const root = channelRoot(opts); await mkdir(root, { recursive: true }); const fmt = opts.audioFormatOverride ?? opts.channelConfig.audioFormat ?? "mp3"; // Resolve the download format (override > channel > global) and expand it // against the source platform so the single-video fix path also honors the // per-source default (Odysee -> original). const downloadFormatSelector = resolveDownloadFormatSelector( resolveDownloadFormatPreset({ override: opts.downloadFormatOverride, channel: opts.channelConfig.downloadFormat, global: getSettings().downloadFormat, }), detectPlatform(opts.singleVideoUrl) ?? opts.channelConfig.platform, ); const args: string[] = [ "--ignore-config", "--restrict-filenames", ...outputArgsForUrl(opts.singleVideoUrl), "--write-info-json", "-f", downloadFormatSelector, "-x", "--audio-format", fmt, ...(opts.channelConfig.keepSourceVideo ? ["-k"] : []), ...configArgs( opts.channelConfig, alwaysCookies(resolveRunCookiePolicy(opts, opts.channelConfig)), ), ...(opts.extraYtdlpArgs ?? []), "--", opts.singleVideoUrl, ]; // --write-info-json rewrites the video's metadata.info.json; the history // keeps what moved (lib/metadataHistory-server.ts). Only when the output dir // is pinned — a %(id)s dir is known only after the run. const dirId = dataDirIdForUrl(opts.singleVideoUrl); const run = () => runChildAndStream(opts, root, args); if (dirId) { await withMetadataHistory( path.join(root, "data", dirId), { by: "download-one", onLog: opts.onLog }, run, ); } else { await run(); } await safeBackfillAvailability(opts); } export async function destinationExists( dataDir: string, id: string, handling: ChannelHandling, opts: { replaceAutoSubs?: boolean } = {}, ): Promise { const dir = path.join(dataDir, id); const entries = await readdir(dir).catch(() => [] as string[]); // Our own transcript always counts as "already done". if (entries.includes("transcript.json")) return true; // A VTT normally counts too — channels can be hybrid. The replace-auto-captions // run is the one exception, and only for a track that still sniffs as YouTube // ASR: that is precisely the transcript we are here to replace, so it must not // suppress the audio download the replacement needs. if (entries.includes("transcript.en.vtt")) { const replaceable = opts.replaceAutoSubs === true && (await readVttProvenance(dir, "transcript.en.vtt")) === "asr"; if (!replaceable) return true; } // Transcribe channels treat raw audio as a download in progress so we don't // re-fetch it before whisper runs. YouTube channels expect a .vtt; an audio // file alone shouldn't suppress the next sync. Use isRealAudioFile (the same // predicate the partial-downloads bucket uses) so audio-check snapshots // (audio..part.good/.part.testing), sidecars (audio.info.json, // audio.live_chat.json), and temp files don't get mistaken for finished audio // — otherwise a genuine partial carrying a .part.good snapshot is skipped. if (handling === "transcribe") { return entries.some(isRealAudioFile); } return false; } // How many newest-first playlist entries to enumerate per sync page. Each page // is one cheap flat-playlist call; new videos on it are fetched, and we stop at // the first page that contains an already-archived entry. const SYNC_PAGE_SIZE = 50; // A sync is one of two passes over the same channel listing: // // syncPaged the cheap default. Walks newest-first one 50-entry // flat-playlist spawn at a time and stops at the first page // containing an already-archived entry. // syncFullSweep the periodic deep pass. ONE full enumeration, from which all // three outputs are derived: a refreshed `playlist` file, the // maybe-missing (upstream deletion) diff, and the very same // download set the paged walk would have chosen. // // The cadence gate decides which (see common/jobs/deepSync.ts); a sweep is far // more expensive on a large channel, so it runs at most once per interval. // `forceFullSweep` overrides the gate in both directions. async function sync(opts: RunYtdlpOpts): Promise { const sweep = opts.forceFullSweep ?? (await fullSweepDue(opts)); return sweep ? syncFullSweep(opts) : syncPaged(opts); } // Whether this sync should upgrade itself to a full sweep, resolved from the // per-channel override, the global cadence and the channel's lastFullSweepAt — // and never while the channel's platform is in a rate-limit cooldown: a sweep // is the longest request run a channel makes, and re-trying one into the // window that just 429'd is how the cooldown gets extended. The paged walk // still runs; `forceFullSweep` still overrides. export async function fullSweepDue( opts: Pick, now: number = Date.now(), ): Promise { const scheduler = getSettings().syncScheduler; const interval = resolveFullSweepIntervalMinutes( opts.channelConfig, scheduler, ); if (!isFullSweepDue(opts.channelConfig.lastFullSweepAt, interval, now)) { return false; } const cooling = await platformCooldownRemainingMs( cooldownPlatformKey(opts.channelConfig), opts.paths, ); return cooling <= 0; } // The per-page download filter, shared by both passes so they can never drift: // entries already in the archive are counted as hits (the paged walk's stop // signal), videos SETTLED by the channel's download filter are counted as hits // TOO, defer-mode needs_auth videos are held back for the Needs-cookies bucket, // and everything else is queued for download. // // WHY A SETTLED VIDEO IS A HIT. `archivedHits > 0` is the paged walk's stop // signal — "we have reached content we already have". A settled video is // content we have already decided about, and it deliberately has no archive // line (an archive line would mean "downloaded" to verifyTranscripts and the // missingFromArchive bucket) and, if the metadata scan settled it, no directory // at all. Counting it as new instead would make a filtered channel's daily sync // walk every page of ~1,800 non-matches to find nothing, forever. // // `settled` is a SET, resolved once per run by the caller. It used to be a // per-video sidecar read; the whole point of deriving settlement from the // channel-level metadata-scan store is that this loop does no I/O for it. async function selectDownloadableUrls( pageUrls: ReadonlyArray, archive: Awaited>, dataDir: string, runCookiePolicy: ResolvedCookiePolicy, settled: ReadonlySet, ): Promise<{ newUrls: string[]; archivedHits: number; deferredAuthCount: number; settledCount: number; }> { const newUrls: string[] = []; let archivedHits = 0; let deferredAuthCount = 0; let settledCount = 0; for (const url of pageUrls) { const archiveId = await archiveIdForUrl(url, dataDir); if (archiveId && archive.ids.has(archiveId)) { archivedHits++; continue; } const dirId = extractVideoId(url); // Settled by the channel's download filter: counted as a hit (see above), // and reported separately so the log doesn't claim they were archived. if (dirId && settled.has(dirId)) { archivedHits++; settledCount++; continue; } // Cookie mode "defer": don't re-attempt known auth-gated videos on // every sync — they wait in the Needs-cookies bucket instead. if ( dirId && (await isDeferredAuthExcluded(path.join(dataDir, dirId), runCookiePolicy)) ) { deferredAuthCount++; continue; } newUrls.push(url); } return { newUrls, archivedHits, deferredAuthCount, settledCount }; } async function syncPaged(opts: RunYtdlpOpts): Promise { const root = channelRoot(opts); const dataDir = path.join(root, "data"); const archivePath = path.join(root, "archive"); // Create the channel root (cwd for enumeration) but NOT data/ — the per-video // downloads create their own canonical dirs, so a sync that's cancelled // before any download leaves no data/ behind. await mkdir(root, { recursive: true }); // Seed the roster from the state that exists right now. A channel may go a // long time between sweeps, and the seeding is itself the protection: it // captures today's listed-but-never-fetched entries, with their URLs, before // any later enumeration can overwrite the stored playlist. await seedRosterIfAbsent(opts.paths, opts.channelSlug, { now: new Date().toISOString(), onLog: opts.onLog, }); // Walk the channel newest-first a page at a time. On each page, download the // entries not yet in the archive, then stop once we reach a page that // contains an already-archived entry (we've caught up to a prior sync) or a // short final page. Each video is fetched via downloadOneManaged — same // auth-retry / audio-check / no-subs-fallback handling as every other // download path — and lands directly in its canonical data// dir. Deeper // mid-channel gaps remain the job of "download missing", exactly as before. let archive = await readArchive(archivePath); let totalNew = 0; let totalFailed = 0; let totalSkipped = 0; let firstFailure: Error | null = null; // Resolved once for the whole sync: drives both the defer-mode needs_auth // exclusion below and the managed downloads themselves. const runCookiePolicy = resolveRunCookiePolicy(opts, opts.channelConfig); // Resolved ONCE for the whole walk: the settled set is derived from one // channel-level file plus the current patterns, so a page costs no I/O for it. const settledIds = await settledByTitleFilterIds( opts.paths, opts.channelSlug, opts.channelConfig, ); for ( let page = 0; !opts.signal.aborted && !opts.drainSignal?.aborted; page++ ) { const start = page * SYNC_PAGE_SIZE + 1; const end = start + SYNC_PAGE_SIZE - 1; const pageUrls = await enumeratePlaylistUrls(opts, root, { start, end }); if (pageUrls.length === 0) break; // ADD-ONLY: a page is a real (if partial) sighting of the listing, so the // ids on it belong in the roster — but a paged walk sees only the newest // window, so it may never compute a "missing" set. That rule is what makes // every non-sweep writer safe to call. await mergeRosterFile( opts.paths, opts.channelSlug, observationsFromUrls(pageUrls), new Date().toISOString(), "listing", ); const { newUrls, archivedHits, deferredAuthCount, settledCount } = await selectDownloadableUrls( pageUrls, archive, dataDir, runCookiePolicy, settledIds, ); opts.onLog( `Sync page ${page + 1}: ${pageUrls.length} entries, ${newUrls.length} new, ${archivedHits} already archived` + (settledCount ? ` (${settledCount} settled by the download filter)` : "") + `.\n`, ); if (deferredAuthCount > 0) { opts.onLog( `Cookie mode is "defer": needs_auth deferred=${deferredAuthCount} — run the "Needs cookies" bucket to download them with cookies.\n`, ); } if (newUrls.length > 0) { const res = await runManagedDownloads( opts, newUrls, opts.channelConfig, runCookiePolicy, ); totalNew += res.okCount; totalFailed += res.failedCount; totalSkipped += res.skippedCount; if (res.firstFailure) { firstFailure = res.firstFailure; break; } // Successful downloads appended to the archive; refresh so the next // page's diff sees them. archive = await readArchive(archivePath); } if (archivedHits > 0) break; // reached previously-synced content if (pageUrls.length < SYNC_PAGE_SIZE) break; // last page } if (firstFailure && !opts.signal.aborted) throw firstFailure; opts.onLog( `Sync complete: ${totalNew} new downloaded, ${totalFailed} failed` + (totalSkipped ? `, ${totalSkipped} skipped by filter` : "") + `.\n`, ); await touchLastSync(opts); await safeBackfillAvailability(opts); } // The periodic deep pass. ONE full `--flat-playlist --print url` enumeration // serves all three purposes that used to cost three separate yt-dlp spawns: // refreshing the stored `playlist`, detecting videos that have left the channel // listing, and picking up new uploads. // // The download set is deliberately IDENTICAL to what syncPaged would have // fetched: newest-first entries up to the first SYNC_PAGE_SIZE window // containing an already-archived entry, just applied to a list already in // memory. Downloading everything absent from the archive would be a // download-from-playlist, and on a channel deliberately kept partial (newest // 200 of 5,000) it would drag down the entire back catalogue — including audio // the cleanup sweep has already reclaimed. A sweep changes what we KNOW, never // what we FETCH; the deletion signal comes from the listing diff, not downloads. async function syncFullSweep(opts: RunYtdlpOpts): Promise { const root = channelRoot(opts); const dataDir = path.join(root, "data"); const archivePath = path.join(root, "archive"); await mkdir(root, { recursive: true }); opts.onLog( `Full sweep: re-reading the whole channel listing (refreshes the video list and flags videos that have gone missing).\n`, ); // 1. One enumeration, no range, then the gate: record everything it saw in // the roster (additive, so this is safe unconditionally) and decide // whether the listing itself may be acted on. let urls: string[]; try { urls = await enumeratePlaylistUrls(opts, root); } catch (err) { if (!(err instanceof EnumerationIncompleteError)) throw err; // A 429 part-way through the listing: what was printed is a prefix, not a // listing, so it never reaches acceptEnumeration (no playlist rewrite, no // missing set, no lastFullSweepAt). Record the platform's cooldown the way // a download's 429 does, then do this sync's job the cheap way. // Best-effort, as on the download path: a failed state write must not fail // the sync this branch exists to rescue. try { await opts.onPlatformBackoff?.("rate_limit"); } catch (backoffErr) { opts.onLog( `Warning: could not record the ${err.platform} cooldown (${ backoffErr instanceof Error ? backoffErr.message : String(backoffErr) }); continuing.\n`, ); } opts.onLog( `Full sweep incomplete: 429 at ${ err.pagesReached !== null ? `page ${err.pagesReached}` : "an unknown page" } of the listing, ${err.count} entries — not a listing. Ran the paged walk instead; syncs wait for the ${err.platform} cooldown, then the sweep is retried.\n`, ); if (opts.signal.aborted) return; return syncPaged(opts); } if (opts.signal.aborted) return; const { decision, roster, listedIds, now } = await acceptEnumeration( opts, root, urls, ); opts.onLog(`Full sweep: ${decision.reason}.\n`); // A rejected listing must not truncate the stored playlist, must not flag // videos missing, and must not stamp lastFullSweepAt — so the next sync // retries instead of waiting out the cadence. It DOES still run the download // walk below: downloading entries the listing does contain is purely // additive, and skipping it would stall a channel's downloads for as long as // the listing stays suspect. let maybeMissing: string[] = []; if (decision.accept) { // 2. Refresh the stored playlist. Everything downstream — "download // missing", the snapshot's undownloaded work-list — reads this file, and // before the sweep only an explicit "store playlist" ever rewrote it. await writePlaylistFile(root, urls, opts.onLog); // 3. Derive every set from one place (controller/channelSets). The // maybe-missing record keeps its exact shape and meaning — videos we // have on disk that the listing no longer carries — so buildIndex, // channelSnapshot and verifyBeforeClean need no changes. What's new is // missingNeverFetched: entries the roster knows we were told about, // never fetched, and that are now gone. Without the roster those had no // dir to be noticed by and no record to be noticed in. const entries = await readdir(dataDir, { withFileTypes: true }).catch( () => [] as Awaited> & { length: 0 }, ); const onDiskIds = new Set( (entries as Array<{ isDirectory(): boolean; name: string }>) .filter((e) => e.isDirectory()) .map((e) => e.name), ); // THE SETTLED SET IS SUBTRACTED, and it is not a nicety. A title-filter // rejection leaves no directory now, so without this a filtered-out video // that has since left the listing reads as `missingNeverFetched` — "we were // told about this, never got it, and it is gone" — and the sweep tells the // operator to attempt a direct-link recovery of the very video they // configured us not to fetch. const sweepSettledIds = await settledByTitleFilterIds( opts.paths, opts.channelSlug, opts.channelConfig, ); const sets = deriveChannelSets({ roster, listedIds, onDiskIds, settledIds: sweepSettledIds, }); maybeMissing = sets.missingDownloaded; await writeMaybeMissing(opts.paths, opts.channelSlug, { checkedAt: now, freshPlaylistCount: listedIds.size, ids: maybeMissing, }); opts.onLog( `Full sweep: ${onDiskIds.size} known, ${listedIds.size} in fresh listing, ${maybeMissing.length} maybe-missing.\n`, ); if (sets.missingNeverFetched.length > 0) { opts.onLog( `Full sweep: ${sets.missingNeverFetched.length} video(s) were listed but never downloaded and have now left the listing — see "Never fetched, now gone" on the channel page to attempt a direct-link recovery.\n`, ); } } // 4. Walk the same listing in newest-first windows, applying the paged walk's // filter and stopping rule to slices instead of spawns. let archive = await readArchive(archivePath); let totalNew = 0; let totalFailed = 0; let totalSkipped = 0; let firstFailure: Error | null = null; const runCookiePolicy = resolveRunCookiePolicy(opts, opts.channelConfig); // Resolved ONCE for the whole walk: the settled set is derived from one // channel-level file plus the current patterns, so a page costs no I/O for it. const settledIds = await settledByTitleFilterIds( opts.paths, opts.channelSlug, opts.channelConfig, ); for ( let page = 0; !opts.signal.aborted && !opts.drainSignal?.aborted; page++ ) { const pageUrls = urls.slice( page * SYNC_PAGE_SIZE, (page + 1) * SYNC_PAGE_SIZE, ); if (pageUrls.length === 0) break; const { newUrls, archivedHits, deferredAuthCount, settledCount } = await selectDownloadableUrls( pageUrls, archive, dataDir, runCookiePolicy, settledIds, ); opts.onLog( `Sync page ${page + 1}: ${pageUrls.length} entries, ${newUrls.length} new, ${archivedHits} already archived` + (settledCount ? ` (${settledCount} settled by the download filter)` : "") + `.\n`, ); if (deferredAuthCount > 0) { opts.onLog( `Cookie mode is "defer": needs_auth deferred=${deferredAuthCount} — run the "Needs cookies" bucket to download them with cookies.\n`, ); } if (newUrls.length > 0) { const res = await runManagedDownloads( opts, newUrls, opts.channelConfig, runCookiePolicy, ); totalNew += res.okCount; totalFailed += res.failedCount; totalSkipped += res.skippedCount; if (res.firstFailure) { firstFailure = res.firstFailure; break; } archive = await readArchive(archivePath); } if (archivedHits > 0) break; // reached previously-synced content if (pageUrls.length < SYNC_PAGE_SIZE) break; // last page } if (firstFailure && !opts.signal.aborted) throw firstFailure; opts.onLog( `Sync complete: ${totalNew} new downloaded, ${totalFailed} failed` + (totalSkipped ? `, ${totalSkipped} skipped by filter` : "") + `.\n`, ); // 5. Capped auto-confirm. Resolving a suspect (deleted vs private vs unlisted) // costs one yt-dlp probe each, so it only runs in-line while the suspect // count is small — the common case of a handful of videos disappearing. // Over the cap this is a manual decision, not something an unattended sync // should fan out into. Runs in the same job (one streamed log, one Cancel), // so re-check the abort between phases. await confirmMaybeMissing(opts, maybeMissing); await touchLastSync(opts); // Only an accepted listing counts as a sweep. Leaving the stamp alone is what // makes the next sync retry the enumeration immediately instead of waiting // out the cadence — and, for a genuine mass deletion, what turns the retry // into the confirming second observation. if (decision.accept) await touchLastFullSweep(opts); await safeBackfillAvailability(opts); } // Resolve the sweep's suspects with the targeted per-video availability probe, // when there are few enough to be worth doing unattended. The option set mirrors // checkMaybeMissingAction / verifyBeforeClean: `ignoreShard` is required because // a saved shard-availability.json would otherwise replace `onlyIds` wholesale, // and `skipExpectedAbsent` avoids re-probing videos already known gone. async function confirmMaybeMissing( opts: RunYtdlpOpts, maybeMissing: ReadonlyArray, ): Promise { if (maybeMissing.length === 0) return; if (opts.signal.aborted) return; const cap = getSettings().syncScheduler.fullSweepConfirmMaxSuspects; if (cap <= 0 || maybeMissing.length > cap) { opts.onLog( `Full sweep: ${maybeMissing.length} maybe-missing video(s) flagged, over the auto-confirm cap of ${cap} — run "Check maybe-missing" on the channel to resolve them.\n`, ); return; } opts.onLog( `Full sweep: confirming ${maybeMissing.length} maybe-missing video(s) upstream…\n`, ); try { await runAvailabilityCheck({ channelSlug: opts.channelSlug, paths: opts.paths, mode: "recheck-all", onlyIds: maybeMissing, skipExpectedAbsent: true, ignoreShard: true, concurrency: 1, onLog: opts.onLog, signal: opts.signal, }); } catch (err) { // Best-effort, exactly like the availability backfill: the sweep's real // product is the flagged set, which is already persisted. opts.onLog( `Full sweep: maybe-missing confirmation skipped: ${(err as Error).message}\n`, ); } } async function safeBackfillAvailability(opts: RunYtdlpOpts): Promise { // Best-effort: write availability.json for any new video dir from this // download whose metadata.info.json has an availability field. Don't fail // the download if the sweep errors. try { await backfillAvailabilityFromMetadata({ channelSlug: opts.channelSlug, paths: opts.paths, onLog: opts.onLog, signal: opts.signal, }); } catch (err) { opts.onLog(`Availability backfill skipped: ${(err as Error).message}\n`); } } // Returns the last STDERR_TAIL_BYTES of stderr (and attaches it to the thrown // error as `stderrTail`), so a caller can classify what yt-dlp said. async function runChildAndStream( opts: RunYtdlpOpts, cwd: string, args: string[], ): Promise { opts.onLog(`$ ${opts.paths.ytdlpBin} ${args.join(" ")}\n`); if (opts.requestCounter) opts.requestCounter.n++; const child = execa(opts.paths.ytdlpBin, args, { cwd, cancelSignal: opts.signal, all: true, buffer: false, reject: false, }); child.all?.on("data", (c: Buffer) => opts.onLog(c.toString("utf8"))); let stderrTail = ""; child.stderr?.on("data", (c: Buffer) => { stderrTail = (stderrTail + c.toString("utf8")).slice(-STDERR_TAIL_BYTES); }); const result = await child; // yt-dlp exit code convention: 101 = "break-on-existing" / "max-downloads" // (clean stop, not an error). Treat it the same as 0. if ( result.exitCode !== 0 && result.exitCode !== 101 && !opts.signal.aborted ) { throw Object.assign(new Error(`yt-dlp exited with code ${result.exitCode}`), { stderrTail, }); } if (result.exitCode === 101) { opts.onLog(`yt-dlp stopped on existing entry (exit 101).\n`); } return stderrTail; } // The three sync-state stamps (CHANNEL_SYNC_STATE_KEYS). Each re-reads the // config at the moment it writes, so an edit made while the sync ran survives; // a channel whose config.json is gone or unreadable is not recreated. async function touchLastSync(opts: RunYtdlpOpts): Promise { await patchChannelConfig(opts.paths, opts.channelSlug, { lastSyncedAt: new Date().toISOString(), }); } async function touchLastFullDownload(opts: RunYtdlpOpts): Promise { await patchChannelConfig(opts.paths, opts.channelSlug, { lastFullDownloadAt: new Date().toISOString(), }); } // Stamps when this channel last paid for a full enumeration, which is what the // cadence gate reads to keep the next N syncs on the cheap paged walk. async function touchLastFullSweep(opts: RunYtdlpOpts): Promise { await patchChannelConfig(opts.paths, opts.channelSlug, { lastFullSweepAt: new Date().toISOString(), }); } // The channel's archive file stores yt-dlp's native extractor ids (e.g. // `twitch:vod v123`, `rumble `), which diverge from our canonical id // on every platform except YouTube. A downloaded video records its native id as // `id` inside metadata.info.json, which (post-reconcile) lives in the canonical // dir. Resolve it from there; for not-yet-downloaded videos (no dir on disk) or // YouTube (native == canonical) fall back to the canonical id. export async function archiveIdForUrl( url: string, dataDir: string, ): Promise { const canonical = extractVideoId(url); if (!canonical) return null; if (detectPlatform(url) === "youtube") return canonical; try { const raw = await readFile( path.join(dataDir, canonical, "metadata.info.json"), "utf8", ); const id = (JSON.parse(raw) as { id?: unknown }).id; if (typeof id === "string" && id) return id; } catch { /* not downloaded yet — fall through to canonical */ } return canonical; } // extractVideoId now lives in ../lib/videoId — a true leaf, so the roster store // can canonicalize URLs without importing this module (which would be a cycle). // Re-exported here because every existing caller imports it from this path. export { extractVideoId };