"use server"; import path from "node:path"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { getSettings } from "yt-dlp-transcript-common/lib/settings"; import { diskGate } from "yt-dlp-transcript-common/lib/diskSpace"; import { formatBytes } from "yt-dlp-transcript-common/lib/format"; import { SAVED_VIDEO_CLIP_PLATFORM, SAVED_VIDEO_CLIP_QUEUE, clipWindowQueueKey, resolveQueueKey, } from "yt-dlp-transcript-common/lib/queueKeys"; import { getRegistry } from "yt-dlp-transcript-common/jobs/registry"; import { windowInFlight } from "yt-dlp-transcript-common/jobs/windowJobs"; import { detectPlatform, queueKeyForUrl } from "yt-dlp-transcript-common/lib/platform"; import { MAX_CLIP_WINDOW_SECONDS, isFetchMaxHeight, } from "yt-dlp-transcript-common/lib/clipWindow"; import { findContainingClipWindow } from "yt-dlp-transcript-common/lib/clipWindow-server"; import { savedWindowSource } from "yt-dlp-transcript-common/lib/savedVideoWindow-server"; import type { ChannelConfig } from "yt-dlp-transcript-common/lib/channelConfig"; import { isValidChannelSlug, readChannelConfig, } from "yt-dlp-transcript-common/controller/channels"; import { findVideoSourceUrl } from "yt-dlp-transcript-common/controller/undownloadedVideos"; import { fetchWindows, fetchWindowsItemLabel, dedupeFetchWindowsItems, type FetchWindowsItem, } from "yt-dlp-transcript-common/controller/fetchWindows"; import { missingEvidenceWindows, type UnfetchableEvidenceWindow, } from "yt-dlp-transcript-common/publish/reportMedia"; import { isValidSiteId, listSiteIds } from "yt-dlp-transcript-common/lib/site"; import { heldPlatformRefusal, platformCooldownRemainingMs, } from "yt-dlp-transcript-common/jobs/downloadBackoff"; import { runManagedFunction, type StreamActionResult, } from "yt-dlp-transcript-common/jobs/streamCommand"; import { safeRevalidate } from "../../../lib/safeRevalidate"; // FETCH A LIST OF CLIP WINDOWS through the managed path, as ONE JOB PER // PLATFORM QUEUE (controller/fetchWindows.ts walks each one, paced). // // The batch form of fetchWindowAction (videos/[id]/videoActions.ts), with the // same rules at the door: a window of at most MAX_CLIP_WINDOW_SECONDS, a height // cap in range, a channel that exists, a URL that resolves. What it adds is the // fan-out: the list is grouped by the platform of each window's own URL // (`queueKeyForUrl`) and each group starts its own job on that platform's CLIP // queue (`clipWindowQueueKey`: `clips:youtube`), so YouTube and Rumble run side // by side, and neither waits behind its platform's long downloads (release 19, // A5 — it used to be the download queue itself, behind a multi-hour persist). // // A PLATFORM COOLING DOWN OR HELD is refused at the door for ITS group (the // sentence comes back in `refused`); the other groups still start. A window // already on disk is answered here (`cached`) and joins no job, and so is one a // queued or running job will already write (`inFlight`, with that job's id — // jobs/windowJobs.ts): the same manifest sent twice queues nothing twice. // // A WINDOW A SAVED CONTAINER HOLDS (release 21 D2) joins no platform's job: it // goes to one group of its own, platform `saved-video` on the // `clips:saved-video` queue, whose job CUTS each window from the container // (controller/fetchWindows.ts) — no request, no pause, no hold or cooldown to // pass. A pointer whose container cannot be read is `unresolved` with the // media guard's sentence, never sent to the network; one whose container ends // before the window is fetched as before. // // Like the single fetch, NOT GATED BY THE DOWNLOAD PAUSE: an operator (or a // tool they are driving) asked for these seconds by hand. const ID_RE = /^[\w.-]+$/; const isVideoId = (v: string): boolean => ID_RE.test(v) && v !== "." && v !== ".."; export type FetchWindowsRequest = { items: FetchWindowsItem[]; // Who asked: the tool or surface. Recorded beside every window. requestedBy: string; manifest?: string; maxHeight?: number; dryRun?: boolean; // A queue override for every group (resolveQueueKey's rules). queueKey?: string; }; export type FetchWindowsUnresolved = { item: FetchWindowsItem; error: string }; export type FetchWindowsGroup = { // The cooldown key (`detectPlatform(url)`), and the queue the job runs on. // `saved-video` (SAVED_VIDEO_CLIP_PLATFORM) for the windows cut from saved // containers, on SAVED_VIDEO_CLIP_QUEUE. platform: string; queueKey: string; items: FetchWindowsItem[]; }; export type FetchWindowsActionResult = | { ok: false; error: string } | { ok: true; dryRun: boolean; // A real run: one per group that started. jobs: { platform: string; queueKey: string; jobId: string; items: number }[]; // A dry run: what each job would be given. groups: FetchWindowsGroup[]; // Groups refused at the door, with the refusal's own sentence. refused: { platform: string; error: string; items: number }[]; cached: FetchWindowsItem[]; // Windows a queued or running job will already write, and that job. inFlight: FetchWindowsInFlight[]; unresolved: FetchWindowsUnresolved[]; }; export type FetchWindowsInFlight = { item: FetchWindowsItem; jobId: string }; // The window rules the HTTP door enforces, asked again here: a replayed spec is // a file on disk. function itemProblem(item: FetchWindowsItem): string | null { if (!isValidChannelSlug(item.slug)) return `"${item.slug}" is not a channel slug`; if (!isVideoId(item.id)) return `"${item.id}" is not a video id`; const { from, to } = item; if ( !Number.isFinite(from) || !Number.isFinite(to) || from < 0 || from >= to || to - from > MAX_CLIP_WINDOW_SECONDS ) { return ( `${from}–${to} is not a fetchable window ` + `(at most ${MAX_CLIP_WINDOW_SECONDS}s, from < to, from >= 0)` ); } return null; } // Validate, answer the cache, resolve every URL and group by queue. Nothing is // started; nothing touches the network. async function planFetchWindows( items: FetchWindowsItem[], queueOverride: string | undefined, maxHeight: number | undefined, ): Promise<{ groups: FetchWindowsGroup[]; cached: FetchWindowsItem[]; inFlight: FetchWindowsInFlight[]; unresolved: FetchWindowsUnresolved[]; }> { const paths = getPaths(); const configs = new Map(); const groups = new Map(); const cached: FetchWindowsItem[] = []; const inFlight: FetchWindowsInFlight[] = []; const unresolved: FetchWindowsUnresolved[] = []; const live = getRegistry().list(); for (const item of dedupeFetchWindowsItems(items)) { const problem = itemProblem(item); if (problem) { unresolved.push({ item, error: problem }); continue; } if (!configs.has(item.slug)) { configs.set(item.slug, await readChannelConfig(paths, item.slug)); } const config = configs.get(item.slug); if (!config) { unresolved.push({ item, error: `Channel "${item.slug}" not found` }); continue; } const videoDir = path.join(paths.channelsDir, item.slug, "data", item.id); if (await findContainingClipWindow(videoDir, item.from, item.to)) { cached.push(item); continue; } const running = windowInFlight(live, { slug: item.slug, videoId: item.id, from: item.from, to: item.to, maxHeight, }); if (running) { inFlight.push({ item, jobId: running.jobId }); continue; } // The saved container: a window it holds is cut, on a queue no platform // shares; one it cannot be read for is refused here, not fetched. const saved = await savedWindowSource({ paths, slug: item.slug, videoDir, from: item.from, to: item.to, }); if (saved.kind === "unreadable") { unresolved.push({ item, error: saved.error }); continue; } if (saved.kind === "covers") { const queueKey = resolveQueueKey(SAVED_VIDEO_CLIP_QUEUE, queueOverride); const group = groups.get(queueKey) ?? { platform: SAVED_VIDEO_CLIP_PLATFORM, queueKey, items: [], }; group.items.push(item); groups.set(queueKey, group); continue; } const url = item.webpageUrl?.trim() || (await findVideoSourceUrl(paths, item.slug, item.id, config)); if (!url) { unresolved.push({ item, error: "Could not determine the video URL: no metadata.info.json and the " + "playlist does not contain a matching entry.", }); continue; } // THE WINDOW'S OWN URL decides the queue, not the channel's: a channel // whose URL is no platform's (a curated mix) can hold Rumble videos, and // its queue would put a second Rumble job beside the first, each pacing // only itself. const queueKey = resolveQueueKey(clipWindowQueueKey(queueKeyForUrl(url)), queueOverride); const platform = detectPlatform(url) ?? "unknown"; // One job per queue. Two platforms sharing a queue (an override) share a // job too; the controller keys the cooldown per item, so each is honoured. const group = groups.get(queueKey) ?? { platform, queueKey, items: [] }; group.items.push({ ...item, webpageUrl: url }); groups.set(queueKey, group); } return { groups: [...groups.values()], cached, inFlight, unresolved }; } // The door each group passes before its job is queued: the platform's hold, // then its cooldown — fetchWindowAction's checks and sentences. async function groupRefusal(platform: string): Promise { // Cuts from saved containers ask no platform. if (platform === SAVED_VIDEO_CLIP_PLATFORM) return null; const paths = getPaths(); const held = await heldPlatformRefusal(platform, "This batch", paths); if (held) return held; const cooldownMs = await platformCooldownRemainingMs(platform, paths); if (cooldownMs > 0) { return ( `${platform} is in a rate-limit cooldown ` + `(${Math.ceil(cooldownMs / 1000)}s remaining).` ); } return null; } // One platform's windows, as one job. Exported for Retry: the replay hands the // spec's items back here, and the controller's cache check skips whatever an // earlier run fetched. export async function fetchWindowsJobAction(req: { items: FetchWindowsItem[]; requestedBy: string; manifest?: string; siteId?: string; maxHeight?: number; queueKey: string; }): Promise { if (req.items.length === 0) return { ok: false, error: "No windows to fetch." }; if (req.maxHeight !== undefined && !isFetchMaxHeight(req.maxHeight)) { return { ok: false, error: `maxHeight ${req.maxHeight} is not a source height to cap a fetch at.`, }; } for (const item of req.items) { const problem = itemProblem(item); if (problem) return { ok: false, error: `${fetchWindowsItemLabel(item)}: ${problem}` }; } const paths = getPaths(); const requestedAt = new Date().toISOString(); const slugs = [...new Set(req.items.map((i) => i.slug))]; return runManagedFunction({ kind: "fetch-windows", queueKey: req.queueKey, paths, // A batch of ONE channel carries it, so /jobs and the media guard can name // it; one spanning channels carries none, and the controller asks each // channel's text guard itself. ...(slugs.length === 1 ? { channelSlug: slugs[0] } : {}), spec: { kind: "fetch-windows", // A spec needs a slug; the replay reads `params.items`. slug: slugs[0], params: { items: req.items, requestedBy: req.requestedBy, queueKey: req.queueKey, ...(req.manifest ? { manifest: req.manifest } : {}), ...(req.siteId ? { siteId: req.siteId } : {}), ...(req.maxHeight !== undefined ? { maxHeight: req.maxHeight } : {}), }, }, fn: async (onLog, signal, setProgress, ctx) => { const result = await fetchWindows({ paths, items: req.items, provenance: { requestedBy: req.requestedBy, manifest: req.manifest, requestedAt, }, maxHeight: req.maxHeight, onLog, signal, drainSignal: ctx.drainSignal, setProgress, }); safeRevalidate([ ...new Set(req.items.map((i) => `/channels/${i.slug}/videos/${i.id}`)), ]); // A run that stopped short or lost a window did not do what it was // asked: the job says so, and a re-run picks up the rest. if (result.stopped && result.stopped !== "drain" && result.stopped !== "cancel") { throw new Error( `Stopped (${result.stopped}) with ${result.notAttempted.length} window(s) ` + `not attempted — run it again later.`, ); } if (result.failed.length > 0) { throw new Error(`${result.failed.length} window(s) failed to fetch.`); } }, }); } // Fetch a list of windows ACROSS channels and platforms. A dry run answers with // the plan and starts nothing. export async function fetchWindowsAction( req: FetchWindowsRequest & { siteId?: string }, ): Promise { const requestedBy = req.requestedBy.trim(); if (!requestedBy) { return { ok: false, error: "requestedBy is required (who is asking for these bytes)." }; } if (req.maxHeight !== undefined && !isFetchMaxHeight(req.maxHeight)) { return { ok: false, error: `maxHeight ${req.maxHeight} is not a source height to cap a fetch at.`, }; } const { groups, cached, inFlight, unresolved } = await planFetchWindows( req.items, req.queueKey, req.maxHeight, ); const base = { groups, cached, inFlight, unresolved, jobs: [], refused: [] }; if (req.dryRun) return { ok: true, dryRun: true, ...base }; if (groups.length === 0) return { ok: true, dryRun: false, ...base }; // A window is small, but it lands on the corpus disk: the floor an operator's // click asks (manual mode), once, against the channels' text. const paths = getPaths(); const disk = await diskGate(paths, getSettings(), { mode: "manual", dir: paths.channelsDir, }); if (!disk.ok) { return { ok: false, error: `Low disk space: ${formatBytes(disk.freeBytes)} free, ` + `${formatBytes(disk.thresholdBytes)} required. Free up space or ` + `lower the floor in Settings.`, }; } const jobs: { platform: string; queueKey: string; jobId: string; items: number }[] = []; const refused: { platform: string; error: string; items: number }[] = []; for (const g of groups) { const refusal = await groupRefusal(g.platform); if (refusal) { refused.push({ platform: g.platform, error: refusal, items: g.items.length }); continue; } let res: StreamActionResult; try { res = await fetchWindowsJobAction({ items: g.items, requestedBy, manifest: req.manifest, siteId: req.siteId, maxHeight: req.maxHeight, queueKey: g.queueKey, }); } catch (e) { refused.push({ platform: g.platform, error: (e as Error).message, items: g.items.length }); continue; } if (!res.ok) { refused.push({ platform: g.platform, error: res.error, items: g.items.length }); continue; } // Nobody reads the stream: the job's log is on disk. void res.stream.cancel().catch(() => {}); jobs.push({ platform: g.platform, queueKey: g.queueKey, jobId: res.jobId, items: g.items.length, }); } if (jobs.length === 0 && refused.length > 0) { return { ok: false, error: refused.map((r) => `${r.platform}: ${r.error}`).join("; "), }; } return { ok: true, dryRun: false, groups, cached, inFlight, unresolved, jobs, refused }; } export type FetchMissingEvidenceResult = | { ok: false; error: string } | (Extract & { siteId: string; // Cited spans already on disk, and the ones no window fetch can fill. onDisk: number; unfetchable: UnfetchableEvidenceWindow[]; }); // Every window a site's published reports cite and the disk does not hold // (publish/reportMedia.ts, missingEvidenceWindows), fetched as above. What the // Reports tab's "Fetch missing evidence" and `pnpm ops fetch-windows // {"siteId": …}` both run. export async function fetchMissingEvidenceAction( siteId: string, opts: { dryRun?: boolean; maxHeight?: number } = {}, ): Promise { const id = siteId.trim(); if (!isValidSiteId(id)) return { ok: false, error: `"${id}" is not a valid site id` }; const paths = getPaths(); if (!listSiteIds(paths).includes(id)) return { ok: false, error: `No site "${id}"` }; const need = await missingEvidenceWindows(paths, id); const extra = { siteId: id, onDisk: need.onDisk, unfetchable: need.unfetchable }; if (need.missing.length === 0) { return { ok: true, dryRun: opts.dryRun === true, groups: [], jobs: [], refused: [], cached: [], inFlight: [], unresolved: [], ...extra, }; } const r = await fetchWindowsAction({ items: need.missing.map((w) => ({ slug: w.slug, id: w.id, from: w.from, to: w.to, clipId: w.clipId, ...(w.reason ? { reason: w.reason.slice(0, 400) } : {}), })), requestedBy: "reports", manifest: id, siteId: id, maxHeight: opts.maxHeight, dryRun: opts.dryRun, }); if (!r.ok) return r; return { ...r, ...extra }; }