import type { StreamActionResult } from "yt-dlp-transcript-common/jobs/streamCommand"; // The one bulk fan-out loop, shared by "Sync every channel" and the per-group // stage buttons on /channels. // // Deliberately NOT in ../actions.ts: that file carries "use server", which means // every non-type export in it must be an async server action — a shared helper // taking callbacks cannot live there. export type QueueOutcome = { queued: string[]; skipped: { slug: string; reason: string }[]; // THE JOB IDS, IN THE SAME ORDER AS `queued`. Added because a fan-out's // response carried no way to FOLLOW what it started: `pnpm ops relocate // --wait` read `{ok, queued:[slug], skipped:[]}`, found no `jobId`, decided // nothing had been started and returned 0 immediately — reporting success // about a 130 GB copy that had not begun. // // A separate array rather than turning `queued` into objects: `queued` is the // documented body shape of two public routes and a UI reads it as slugs. // Additive, never renamed. jobIds: string[]; }; export async function queueForSlugs( slugs: ReadonlyArray, opts: { // Return a reason string to skip this slug, or null to run it. skip?: (slug: string) => string | null | Promise; run: (slug: string) => Promise; }, ): Promise { const queued: string[] = []; const jobIds: string[] = []; const skipped: { slug: string; reason: string }[] = []; for (const slug of slugs) { const reason = opts.skip ? await opts.skip(slug) : null; if (reason) { skipped.push({ slug, reason }); continue; } const result = await opts.run(slug); if (!result.ok) { skipped.push({ slug, reason: result.error }); continue; } queued.push(slug); jobIds.push(result.jobId); // Nobody will read the stream here — cancel it so the buffered chunks // can be GC'd. The job keeps running and writes to its log file via // runManagedFunction's onLog regardless. void result.stream.cancel(); } return { queued, jobIds, skipped }; }