"use server"; import { safeRevalidate } from "../lib/safeRevalidate"; import { listChannelBriefs, type ChannelBrief, } from "yt-dlp-transcript-common/controller/channels"; import type { StreamActionResult } from "yt-dlp-transcript-common/jobs/streamCommand"; import { activeSlugsForKinds } from "yt-dlp-transcript-common/jobs/syncJobs"; import { isValidGroupId } from "yt-dlp-transcript-common/lib/channelGroups"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { getSettings } from "yt-dlp-transcript-common/lib/settings"; import { isGateHeld } from "yt-dlp-transcript-common/lib/pauseGates"; import { getSite, isValidSiteId } from "yt-dlp-transcript-common/lib/site"; import { laneOffFor, slugsInGroup, stationWorkFor, transcribeStationIds, type StationId, } from "yt-dlp-transcript-common/views/channelGroupSections"; import { queueForSlugs } from "./lib/queueForSlugs"; import { downloadMissingAction, syncAction } from "./[slug]/pipelineActions"; import { transcribeAutoSubsBucketAction, transcribeBucketAction, transcribeMissingAction, } from "./[slug]/whisperActions"; import { backfillChannelAction } from "./[slug]/backfillActions"; import { runOperationChannelJob } from "yt-dlp-transcript-common/controller/operationJobs"; import { DIGEST_OPERATION_ID } from "yt-dlp-transcript-common/lib/operations"; // Run one pipeline stage over every channel in one of a site's groups. // // TAKES IDS, NEVER A CLIENT-SUPPLIED SLUG LIST. The membership is re-read from // sites//site.json here, so a tampered call cannot queue work outside the // group, and the fan-out stays correct if membership changed between the render // and the click. // // FAN-OUT IS SAFE: common/lib/queueKeys.ts submits every non-empty queueKey with // concurrency 1, so N transcribeMissingAction calls become N whisper-all records // that run strictly one after another on the shared TRANSCRIPTION_QUEUE — no GPU // thrash. Downloads serialize per platform, backfill shares BACKFILL_QUEUE. The // `await` in the loop awaits ENQUEUE, not completion, so this returns promptly. // Same shape as syncAllChannelsAction and the backfill sweep. export type GroupOpResult = { group: { siteId: string; groupId: string }; queued: string[]; skipped: { slug: string; reason: string }[]; }; // The job kinds each station enqueues, for the "already running" dedupe. The job // registry has no dedupe of its own. Transcribe names THREE: its button queues // one or both bucket jobs (or the scan, for a channel with no report), and the // dedupe spans them all — pressing it again while any one runs, including the // channel page's own bucket job, reads "already running", not a second batch. const KIND_FOR: Record = { sync: ["sync"], download: ["download-missing"], transcribe: [ "whisper-all", "whisper-bucket-downloaded-no-transcript", "whisper-bucket-auto-subs", ], digest: ["digest-channel-local"], speakers: ["backfill-channel"], }; // THE TRANSCRIBE STATION QUEUES WHAT ITS FIGURE COUNTS — `transcribeStationIds`, // the fold `stationWorkFor` sums — BY ID. Two batches, because one cannot cover // both: the no-transcript half is `transcribeBucketAction` over its ids (the // channel page's own bucket control, replayable as `downloadedNoTranscript`), // and the auto-captions half is `transcribeAutoSubsBucketAction`, whose // replace mode needs an ASR track per id. By id matters: a video downloaded // before it went private, members-only or deleted is still on disk and still // in `downloadedNoTranscript`, and a directory SCAN would transcribe it — one // more than the figure said. The channel page runs the same jobs; a combined // job kind would need its own replay spec and /jobs label and mirror nothing. // Handling is not consulted — a youtube channel's captionless download is // whisper work too. // // The SCAN (`transcribeMissingAction`, a whisper-all with no ids) is used ONLY // when the channel has never reported: with no snapshot there are no ids to // name, and the scan finds what there is, as the station always did. // // All three kinds resolve to TRANSCRIPTION_QUEUE, so the "strictly one at a // time" note above stays true. `queued` counts CHANNELS; `jobIds` carries one // id per channel — the no-transcript job's when both halves were queued, and // the other stream is cancelled here exactly as queueForSlugs cancels the one // returned. A half that fails while the other queues cannot ride on an ok // result (StreamActionResult's ok arm has no message), so it is logged. async function transcribeChannel( slug: string, brief: ChannelBrief | undefined, ): Promise { if (!brief?.snapshot) return transcribeMissingAction(slug); const ids = transcribeStationIds(brief.snapshot); const halves = [ ids.missing.length > 0 ? await transcribeBucketAction( slug, ids.missing, undefined, undefined, undefined, "downloadedNoTranscript", ) : null, ids.autoSubs.length > 0 ? await transcribeAutoSubsBucketAction(slug, ids.autoSubs) : null, ].filter((r): r is StreamActionResult => r !== null); // Both empty is skipped upstream as "nothing to do"; say so if reached. if (halves.length === 0) return { ok: false, error: "nothing to do", info: true }; const ok = halves.find((r) => r.ok); if (!ok) return halves[0]; for (const r of halves) { if (r === ok) continue; if (r.ok) void r.stream.cancel(); else console.warn( `[group transcribe] ${slug}: one half queued, the other refused: ${r.error}`, ); } return ok; } const RUN_FOR: Record< StationId, (slug: string, brief: ChannelBrief | undefined) => Promise > = { sync: (slug) => syncAction(slug, undefined), // runPipelineAction already refuses with {ok:false, error} for a paused // install, a per-platform 429 cooldown and low disk, so all three land in // `skipped` for free. download: (slug) => downloadMissingAction(slug), transcribe: (slug, brief) => transcribeChannel(slug, brief), // THE LOCAL LANE, PINNED. The metered lane is behind a settings gate and a // spend cap, so it is never what a group button starts — which is why this // names the lane explicitly instead of letting the dispatcher pick the // configured one. Through runOperationChannelJob so a station that IS one // operation goes through the one runner path. digest: (slug) => runOperationChannelJob({ paths: getPaths(), channelSlug: slug, operation: DIGEST_OPERATION_ID, digest: { lane: "local" }, onDone: () => safeRevalidate([`/channels/${slug}`]), }), // EVERY ENABLED LANE KIND, and so NOT runOperationChannelJob: this station is // the LANE, not one operation on it — KIND_FOR above says the same thing by // mapping it to the lane's job kind. Passing no kindIds is what "all of them" // means to the runner. speakers: (slug) => backfillChannelAction(slug), }; async function runGroupStation( station: StationId, siteId: string, groupId: string, ): Promise { const group = { siteId, groupId }; if (!isValidSiteId(siteId)) { return { group, queued: [], skipped: [{ slug: siteId, reason: "invalid site id" }] }; } if (!isValidGroupId(groupId)) { return { group, queued: [], skipped: [{ slug: groupId, reason: "invalid group id" }] }; } const paths = getPaths(); const settings = getSettings(); // One read, outside the loop. Only throws on an invalid id, guarded above; a // site with no config file parses to one with no channels, so an unknown id // fans out over nothing rather than erroring. const site = getSite(siteId, paths); // Whole-group refusals answer once rather than repeating one identical line // per channel. if (laneOffFor(station, settings)) { return { group, queued: [], skipped: [ { slug: groupId, reason: "the lane is switched off on its operation page", }, ], }; } if (station === "download" && isGateHeld(settings, "download")) { return { group, queued: [], skipped: [ { slug: groupId, reason: "Downloads are paused — resume them from the dashboard.", }, ], }; } const members = new Set(slugsInGroup(site, groupId)); const briefs = (await listChannelBriefs(paths)).filter((b) => members.has(b.slug), ); const briefBySlug = new Map(briefs.map((b) => [b.slug, b])); const active = activeSlugsForKinds(KIND_FOR[station]); const outcome = await queueForSlugs( briefs.map((b) => b.slug), { skip: (slug) => { const brief = briefBySlug.get(slug); if (!brief) return "channel not found"; const work = stationWorkFor(station, brief, settings); if (!work.eligible) return work.reason ?? "not applicable"; if (active.has(slug)) return "already running"; // work === null means the channel has never reported: run it rather // than assume it is finished. Only a report that SAYS zero is a skip — // which is what keeps a 4-channel group from becoming 18 no-op jobs. // // Sync is exempt: it carries no figure at all (syncAction decides for // itself whether the channel is due), so its 0 means "no count", not // "nothing to do", and gating on it would skip every channel. if (station !== "sync" && work.work === 0) return "nothing to do"; return null; }, run: (slug) => RUN_FOR[station](slug, briefBySlug.get(slug)), }, ); return { group, ...outcome }; } export async function syncChannelGroupAction( siteId: string, groupId: string, ): Promise { return runGroupStation("sync", siteId, groupId); } export async function downloadChannelGroupAction( siteId: string, groupId: string, ): Promise { return runGroupStation("download", siteId, groupId); } export async function transcribeChannelGroupAction( siteId: string, groupId: string, ): Promise { return runGroupStation("transcribe", siteId, groupId); } export async function digestChannelGroupAction( siteId: string, groupId: string, ): Promise { return runGroupStation("digest", siteId, groupId); } export async function speakersChannelGroupAction( siteId: string, groupId: string, ): Promise { return runGroupStation("speakers", siteId, groupId); }