"use server"; import { safeRevalidate } from "../../lib/safeRevalidate"; 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 { downloadQueueKey, resolveQueueKey, } from "yt-dlp-transcript-common/lib/queueKeys"; import { readChannelConfig } from "yt-dlp-transcript-common/controller/channels"; import { persistKept } from "yt-dlp-transcript-common/controller/persistKept"; import { persistVideos, type PersistReplacePolicy, type PersistVideoItem, type PersistVideosPlan, } from "yt-dlp-transcript-common/controller/persistVideos"; import type { SourceVideoQuality } from "yt-dlp-transcript-common/ytdlp/downloadFormat"; import { runManagedFunction, type StreamActionResult, } from "yt-dlp-transcript-common/jobs/streamCommand"; export type PersistVideosOptions = { // Absent = each channel's effective source-video quality. format?: SourceVideoQuality; replace?: PersistReplacePolicy; gapMs?: number; minFreeMemMb?: number; queueKey?: string; }; export type PersistVideosActionResult = | { ok: false; error: string } | { ok: true; dryRun: true; plan: PersistVideosPlan } | { ok: true; dryRun: false; plan: PersistVideosPlan; jobs: { slug: string; jobId: string }[]; jobIds: string[]; skipped: { slug: string; reason: string }[]; }; // Bulk catch-up: ensure every video in the channel's keep-latest window has its // source container saved to the store. Re-downloads only those not already // saved. Runs on the channel's download queue since it issues real downloads. export async function persistKeptAction( slug: string, queueKey?: string, ): Promise { const paths = getPaths(); const config = await readChannelConfig(paths, slug); if (!config) return { ok: false, error: `Channel "${slug}" not found` }; // The operator clicked this: only the floor applies and the shared hysteresis // latch is left alone (see diskGate). This is a preflight only — persistKept // re-checks per video, because it writes full containers in a loop. const disk = await persistDiskPreflight(); if (disk) return { ok: false, error: disk }; return runManagedFunction({ kind: "persist-kept", queueKey: resolveQueueKey(downloadQueueKey(config), queueKey), paths, channelSlug: slug, spec: { kind: "persist-kept", slug, params: { queueKey } }, fn: async (onLog, signal) => { await persistKept({ paths, channelSlug: slug, channelConfig: config, onLog, signal, }); safeRevalidate([`/channels/${slug}`, "/saved-videos"]); }, }); } // Persist an explicit list of one channel's videos (controller/persistVideos.ts) // as one job on the channel's download queue — persist-kept's queue, since it // issues real downloads. The job carries its ids and options in its spec, so a // replay re-runs the same list, and a re-run fetches only what is still not // saved. export async function persistChannelVideosAction( slug: string, ids: string[], opts: PersistVideosOptions = {}, ): Promise { const wanted = [...new Set(ids)]; if (wanted.length === 0) return { ok: false, error: "No video ids to persist." }; const paths = getPaths(); const config = await readChannelConfig(paths, slug); if (!config) return { ok: false, error: `Channel "${slug}" not found` }; const disk = await persistDiskPreflight(); if (disk) return { ok: false, error: disk }; const { format, replace, gapMs, minFreeMemMb, queueKey } = opts; return runManagedFunction({ kind: "persist-videos", queueKey: resolveQueueKey(downloadQueueKey(config), queueKey), paths, channelSlug: slug, spec: { kind: "persist-videos", slug, params: { queueKey, ids: wanted, format, replace, gapMs, minFreeMemMb }, }, fn: async (onLog, signal, _progress, ctx) => { const result = await persistVideos({ paths, items: wanted.map((id) => ({ slug, id })), format, replace, gapMs, minFreeMemMb, onLog, signal, drainSignal: ctx.drainSignal, }); safeRevalidate([`/channels/${slug}`, "/saved-videos"]); // A run that stopped short or lost a video did not do what it was asked: // the job says so, and the re-run picks up the rest. if (result.stopped === "low-disk" || result.stopped === "rate-limit") { throw new Error( `Stopped (${result.stopped}) with ${result.notAttempted.length} video(s) not attempted — run it again later.`, ); } if (result.failed.length > 0) { throw new Error(`${result.failed.length} video(s) failed to persist.`); } }, }); } // The disk floor a persist asks before it queues (manual mode: the operator // asked). A preflight only — the controllers re-check per video. async function persistDiskPreflight(): Promise { const disk = await diskGate(getPaths(), getSettings(), { mode: "manual" }); if (disk.ok) return null; return ( `Low disk space: ${formatBytes(disk.freeBytes)} free, ` + `${formatBytes(disk.thresholdBytes)} required. Free up space or ` + `lower the floor in Settings.` ); } // Persist a list of videos ACROSS channels. A dry run answers with the buckets // and starts nothing. A real run starts one persist-videos job per channel that // has something to fetch — so each runs on its channel's download queue, behind // that platform's other downloads, under the channel's media guard — and // reports the channels it refused, with the action's own sentence. export async function persistVideosAction( input: { items: PersistVideoItem[]; dryRun?: boolean } & PersistVideosOptions, ): Promise { const { items, dryRun, ...opts } = input; if (items.length === 0) return { ok: false, error: "No videos to persist." }; const paths = getPaths(); const planned = await persistVideos({ paths, items, format: opts.format, replace: opts.replace, dryRun: true, }); const plan = planned.plan; if (dryRun) return { ok: true, dryRun: true, plan }; if (plan.willFetch === 0) { return { ok: true, dryRun: false, plan, jobs: [], jobIds: [], skipped: [] }; } const disk = await persistDiskPreflight(); if (disk) return { ok: false, error: disk }; // Only the channels with work get a job, and each job gets that channel's // WHOLE slice of the list: its own re-run skips what is saved. const due = new Set( [ ...plan.toFetch.items, ...(opts.replace === "above-height" ? plan.wrongHeight.items : []), ].map((i) => i.slug), ); const bySlug = new Map(); for (const { slug, id } of items) { if (!due.has(slug)) continue; bySlug.set(slug, [...(bySlug.get(slug) ?? []), id]); } const jobs: { slug: string; jobId: string }[] = []; const skipped: { slug: string; reason: string }[] = []; for (const [slug, ids] of bySlug) { let result: StreamActionResult; try { result = await persistChannelVideosAction(slug, ids, opts); } catch (e) { skipped.push({ slug, reason: (e as Error).message }); continue; } if (!result.ok) { skipped.push({ slug, reason: result.error }); continue; } // Nobody reads the stream: the job's log is on disk (see the ops layer). void result.stream.cancel(); jobs.push({ slug, jobId: result.jobId }); } if (jobs.length === 0) { return { ok: false, error: skipped.map((s) => `${s.slug}: ${s.reason}`).join("; "), }; } return { ok: true, dryRun: false, plan, jobs, jobIds: jobs.map((j) => j.jobId), skipped, }; }