import path from "node:path"; import { readdir, readFile } from "node:fs/promises"; import { execa } from "execa"; import pLimit from "p-limit"; import { AVAILABILITY_VALUES, EXCLUDED_FROM_DOWNLOAD, availabilityFromJsonField, classifyDownloadFailure, parseUnavailableFromStderr, type Availability, } from "../lib/availability"; import { recordDownloadBackoff } from "../jobs/downloadBackoff"; import { loadAvailability, recordAvailability, resolveEffectiveAvailability, } from "../lib/availability-server"; import { alwaysCookies, cookieArgs, resolveCookiePolicy, } from "../lib/cookiePolicy"; import { getSettings } from "../lib/settings"; // yt-dlp --dump-json emits one JSON object per video, but may print warnings to // stdout first — so the payload is the LAST parseable line. Shared with the // digest registry's claude-code lane, which has the same problem. import { parseStdoutJson } from "../lib/parseStdoutJson"; import { readChannelConfig } from "./channels"; import { channelExtraArgs, pacedPlatformArgs } from "../ytdlp/channelArgs"; import { detectPlatform } from "../lib/platform"; import { resolveShardItems } from "./shard"; import type { Paths } from "../lib/paths"; export type CheckAvailabilityMode = | "resume" | "recheck-non-deleted" | "recheck-all"; export const CHECK_AVAILABILITY_MODES: CheckAvailabilityMode[] = [ "resume", "recheck-non-deleted", "recheck-all", ]; export type CheckAvailabilityOpts = { channelSlug: string; paths: Paths; mode: CheckAvailabilityMode; concurrency?: number; shardTotal?: number; shardIndex?: number; // When set, resolve + persist the shard slice and return immediately without // checking anything. Backs the "Save shard" button. saveShardOnly?: boolean; // Restrict the scan to these data-dir ids (applied before shard slicing). // Backs the "Full-check unexpected" / maybe-missing flow. onlyIds?: ReadonlyArray; // Skip ids already known permanently gone (deleted / private / members_only, // i.e. EXCLUDED_FROM_DOWNLOAD) — they will not change, so don't re-probe. skipExpectedAbsent?: boolean; // Bypass shard resolution entirely. Needed by callers that pass an explicit, // already-narrowed onlyIds set: with no shardTotal/shardIndex but a saved // shard-availability.json on disk, resolveShardItems returns the SAVED slice // and silently discards the onlyIds filter. Set this when the caller's id list // is the authority (the pre-clean gate probes a specific suspect set). ignoreShard?: boolean; onLog?: (msg: string) => void; signal?: AbortSignal; // Records the shared per-platform cooldown when a probe is rate-limited // (see `blocked` below). Defaults to recordDownloadBackoff — the helper and // key a download's 429 and a manual Sync use. A test seam; production // callers pass nothing. onPlatformBackoff?: (platform: string) => void | Promise; }; export type CheckAvailabilityResult = { attempted: number; checked: number; skipped: number; byStatus: Record; // THE CHECK STOPS ON A RATE LIMIT (release 10, L2 review). A probe that // classifies `rate_limit` — HTTP 429, the bot check, YouTube's soft block // "…isn't available, try again later" — records its own `error` observation // and then no further probe is started: every later id is `skipped`, and // the platform cooldown is recorded once. The ids that WERE probed are in // `probedIds`; anything else a caller reads for this run is an OLD record, // and a caller that acts on the result must not judge it (verifyBeforeClean // leaves it unverified). blocked: boolean; blockMessage?: string; probedIds: string[]; }; function emptyByStatus(): Record { return { public: 0, unlisted: 0, private: 0, members_only: 0, needs_auth: 0, deleted: 0, error: 0, }; } async function readWebpageUrl(videoDir: string): Promise { try { const raw = await readFile( path.join(videoDir, "metadata.info.json"), "utf8", ); const parsed = JSON.parse(raw); if (typeof parsed?.webpage_url === "string") return parsed.webpage_url; return null; } catch { return null; } } export async function runAvailabilityCheck({ channelSlug, paths, mode, concurrency, shardTotal, shardIndex, saveShardOnly = false, onlyIds, skipExpectedAbsent = false, ignoreShard = false, onLog, signal, onPlatformBackoff, }: CheckAvailabilityOpts): Promise { const log = onLog ?? ((m: string) => console.log(m)); const channelDir = path.join(paths.channelsDir, channelSlug); const dataDir = path.join(channelDir, "data"); const config = await readChannelConfig(paths, channelSlug); // Probe cookies in "always" mode ONLY. when-required/defer probes stay // cookie-free deliberately, so auth gating keeps being OBSERVED as // needs_auth — defer mode's exclusion + Needs-cookies bucket depend on that // signal. (Always-mode probes may report a cookie-recoverable video as // public; that's the trade-off of prophylactic cookies.) const probeCookies = alwaysCookies( resolveCookiePolicy(getSettings(), config ?? undefined), ); const limit = pLimit(concurrency && concurrency > 0 ? Math.floor(concurrency) : 1); let allDirs = await readdir(dataDir).catch(() => [] as string[]); if (onlyIds) { const onlyIdSet = new Set(onlyIds); allDirs = allDirs.filter((id) => onlyIdSet.has(id)); } const shardResult = ignoreShard ? { items: allDirs, source: "full" as const, config: null } : await resolveShardItems({ paths, slug: channelSlug, op: "availability", fullItems: allDirs, totalShards: shardTotal, shardIndex: shardIndex, onLog: log, }); const videoDirs = shardResult.items; if (saveShardOnly) { log( `Shard availability: save-only — slice ${(shardResult.config?.shardIndex ?? shardIndex ?? 0) + 1}/${shardResult.config?.totalShards ?? shardTotal} (${videoDirs.length} item(s)) persisted; not checking.`, ); return { attempted: 0, checked: 0, skipped: 0, byStatus: emptyByStatus(), blocked: false, probedIds: [], }; } let attempted = 0; let checked = 0; let skipped = 0; const byStatus = emptyByStatus(); // Set by the first rate-limited probe; see CheckAvailabilityResult.blocked. let blocked: { message: string; url: string } | null = null; const probedIds: string[] = []; // Pre-scan: count how many already have a parsed availability.json so the // "skipped" total in the summary is unambiguous (no, we are not skipping // every video that has a transcript — only those with a valid sidecar). let alreadyChecked = 0; await Promise.all( videoDirs.map(async (id) => { const rec = await loadAvailability(path.join(dataDir, id)); if (rec) alreadyChecked++; }), ); const unchecked = videoDirs.length - alreadyChecked; log( `Availability check (${mode}): ${videoDirs.length} video dir(s) under ${dataDir} — ${alreadyChecked} already checked, ${unchecked} unchecked. concurrency=${concurrency && concurrency > 0 ? Math.floor(concurrency) : 1}`, ); await Promise.all( videoDirs.map((id) => limit(async () => { // A rate-limited source is not asked again in this run: each further // probe would be another request into the same refusal. if (signal?.aborted || blocked) { skipped++; return; } const videoDir = path.join(dataDir, id); const existing = await loadAvailability(videoDir); if (mode === "resume" && existing) { skipped++; return; } if ( mode === "recheck-non-deleted" && existing?.availability === "deleted" ) { skipped++; return; } if (skipExpectedAbsent) { const effective = await resolveEffectiveAvailability(videoDir); if (effective && EXCLUDED_FROM_DOWNLOAD.includes(effective)) { skipped++; log(`Skip ${id}: already ${effective}`); return; } } const url = await readWebpageUrl(videoDir); if (!url) { skipped++; log(`Skip ${id}: no webpage_url in metadata.info.json`); return; } attempted++; const args = [ "--dump-json", "--skip-download", "--no-warnings", ...(config ? channelExtraArgs(config, probeCookies) : [ ...cookieArgs(probeCookies), // At the platform's current pace (release 17, slice RL). ...pacedPlatformArgs(detectPlatform(url)), ]), "--", url, ]; const result = await execa(paths.ytdlpBin, args, { cwd: channelDir, cancelSignal: signal, all: false, buffer: true, reject: false, }).catch((err: unknown) => { return { exitCode: 1, stdout: "", stderr: (err as Error)?.message ?? String(err), } as { exitCode: number | null; stdout: string; stderr: string }; }); if (signal?.aborted) { skipped++; return; } let availability: Availability; let error: string | undefined; if (result.exitCode === 0) { const parsed = parseStdoutJson(String(result.stdout ?? "")); if (parsed && typeof parsed === "object") { availability = availabilityFromJsonField( (parsed as { availability?: unknown }).availability, ); } else { availability = "public"; } } else { const stderr = String(result.stderr ?? ""); availability = parseUnavailableFromStderr(stderr); if (availability === "error") { // Keep a short tail of stderr to aid diagnosis without bloating // the sidecar. const trimmed = stderr.trim().split("\n").slice(-3).join("\n"); if (trimmed) error = trimmed; } if ( !blocked && classifyDownloadFailure(stderr, availability) === "rate_limit" ) { blocked = { message: (error ?? stderr.trim()).slice(0, 300), url, }; } } await recordAvailability(videoDir, { availability, observedAt: new Date().toISOString(), source: "check", error, webpageUrl: url, }).catch((err) => { log( `Warn ${id}: failed to record availability history: ${(err as Error)?.message ?? String(err)}`, ); }); checked++; byStatus[availability]++; probedIds.push(id); log(`Check ${id}: ${availability}`); }), ), ); if (blocked) { const b = blocked as { message: string; url: string }; const platform = detectPlatform(config?.url ?? b.url) ?? "unknown"; const unprobed = videoDirs.length - probedIds.length; log( `STOPPED: the source is rate-limiting availability probes (${b.message}). ` + `${unprobed} video(s) were not probed this run; a ${platform} cooldown has been recorded.`, ); try { await (onPlatformBackoff ?? ((pf) => recordDownloadBackoff(pf, paths)))( platform, ); } catch { /* the cooldown is best-effort; the stop above already happened */ } } // Type guard: every key in AVAILABILITY_VALUES has been incremented above. void AVAILABILITY_VALUES; const summary = AVAILABILITY_VALUES.map((v) => `${v}=${byStatus[v]}`).join( " ", ); log( `Availability check done. attempted=${attempted} checked=${checked} skipped=${skipped} (${summary}).`, ); return { attempted, checked, skipped, byStatus, blocked: blocked !== null, ...(blocked ? { blockMessage: (blocked as { message: string }).message } : {}), probedIds, }; }