// The data-driven replay table: maps a stored JobSpec's `kind` to the server // action that re-runs it. Previously this was a hard-coded switch in // runJobSpec.ts; collapsing it to a lookup means adding a replayable kind is // one colocated entry here (plus its jobKinds.ts metadata) instead of a new // switch case. runJobSpec.ts is now a thin dispatcher over this table, which the // retry actions (Retry / Retry all failed) drive. // // Server-only: reached through the retry actions in actions.ts ("use server"), // so it is never bundled to the client. import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { readChannelSnapshot } from "yt-dlp-transcript-common/controller/channelSnapshot"; import type { JobSpec, ReplayBucket } from "yt-dlp-transcript-common/jobs/jobSpec"; import type { StreamActionResult } from "yt-dlp-transcript-common/jobs/streamCommand"; import type { AudioFormat } from "yt-dlp-transcript-common/lib/channelConfig"; import { isSourceVideoQuality } from "yt-dlp-transcript-common/ytdlp/downloadFormat"; import { isPersistReplacePolicy } from "yt-dlp-transcript-common/controller/persistVideos"; import { downloadAction, downloadMissingAction, downloadMissingSubsAction, retryBucketAction, storePlaylistAction, runMetadataScanAction, syncAction, } from "../channels/[slug]/pipelineActions"; import { checkKeptDeletedAction, cleanAudioAction, cleanExtraAudioFormatsAction, diarizeChannelAction, clearFailedTranscriptionsAction, removeWrongFormatAudioAction, purgeSupersededAutoSubsAction, transcribeAutoSubsBucketAction, transcribeBucketAction, transcribeMissingAction, } from "../channels/[slug]/whisperActions"; import { digestChannelAction, type DigestLaneChoice, } from "../channels/[slug]/digestActions"; import { backfillChannelAction } from "../channels/[slug]/backfillActions"; import { persistChannelVideosAction, persistKeptAction, } from "../channels/[slug]/persistActions"; import { capturePostsAction, checkPostAvailabilityAction, fetchPostsAction, } from "../channels/[slug]/socialActions"; import { redownloadIncompleteBucketAction, redownloadShortAudioBucketAction, } from "../channels/[slug]/incompleteTranscriptActions"; import { refreshVideoMetadataAction, replayFetchWindowAction, } from "../channels/[slug]/videos/[id]/videoActions"; import { fetchWindowsJobAction } from "../channels/[slug]/videos/fetchWindowsAction"; import type { FetchWindowsItem } from "yt-dlp-transcript-common/controller/fetchWindows"; import { reportsPrepareAction } from "../sites/lib/reportsPrepareAction"; import { reportsExportAction } from "../sites/lib/reportsExportAction"; import { enqueueStage, stageRequestFromSpec, } from "yt-dlp-transcript-common/publish/publishStages"; export type ReplayHandler = (spec: JobSpec) => Promise; // Bucket jobs re-derive their work from the channel's CURRENT snapshot (not a // frozen id list), so re-running "retry partial downloads" always acts on // whatever is partial right now. async function idsForBucket( slug: string, bucket: ReplayBucket, ): Promise { const snap = await readChannelSnapshot(getPaths(), slug); const ids = snap?.buckets?.[bucket]; return Array.isArray(ids) ? ids : []; } const str = (v: unknown): string | undefined => typeof v === "string" ? v : undefined; const bool = (v: unknown): boolean | undefined => typeof v === "boolean" ? v : undefined; const num = (v: unknown): number | undefined => typeof v === "number" ? v : undefined; // An id scope: the strings of an array, or undefined (= the whole channel). const strings = (v: unknown): string[] | undefined => Array.isArray(v) ? v.filter((k): k is string => typeof k === "string") : undefined; // A fetch-windows spec's items, keeping only the well-formed: the action // re-checks every window, so this only drops what is not even shaped like one. function windowItems(v: unknown): FetchWindowsItem[] { if (!Array.isArray(v)) return []; const out: FetchWindowsItem[] = []; for (const raw of v) { if (typeof raw !== "object" || raw === null) continue; const r = raw as Record; const slug = str(r.slug); const id = str(r.id); const from = num(r.from); const to = num(r.to); if (!slug || !id || from === undefined || to === undefined) continue; out.push({ slug, id, from, to, clipId: str(r.clipId), reason: str(r.reason), pad: num(r.pad), webpageUrl: str(r.webpageUrl), }); } return out; } // Flag-style params + the captured queueKey for a spec. function params(spec: JobSpec): { p: Record; queueKey: string | undefined; } { const p = spec.params ?? {}; return { p, queueKey: str(p.queueKey) }; } // A PUBLISH STAGE (release 18): the spec IS the stage request, so a retry // re-enqueues it as it was — same run id, same preconditions, which the child // re-asks on disk. A spec the stage row would refuse is refused here. const replayPublishStage: ReplayHandler = async (spec) => { const req = stageRequestFromSpec(spec); if (!req) return { ok: false, error: `Not a publish stage this editor can re-run (${spec.kind} ${spec.slug})` }; return enqueueStage(getPaths(), req); }; export const JOB_REPLAY_HANDLERS: Record = { "publish-update-index": replayPublishStage, "publish-build-site": replayPublishStage, "publish-deploy-site": replayPublishStage, "publish-build-hub": replayPublishStage, "publish-deploy-hub": replayPublishStage, "publish-build-homepage": replayPublishStage, "publish-deploy-homepage": replayPublishStage, // A report site's evidence media. The spec's slug is the SITE (the job has // no channel); a re-run re-reads the site's reports and cuts only what its // cache does not already hold. "reports-prepare": (spec) => { const { p } = params(spec); return reportsPrepareAction(str(p.siteId) ?? spec.slug); }, // A report site's exports: the same site, report and formats, from the // reports as they are now. "reports-export": (spec) => { const { p } = params(spec); const reportId = str(p.reportId); const formats = strings(p.formats); return reportsExportAction(str(p.siteId) ?? spec.slug, { ...(reportId ? { reportId } : {}), ...(formats ? { formats } : {}), }); }, // A clip window sourced for another tool. Replay RE-DERIVES from disk like // every bucket job does: if the window (or a wider one covering it) has // arrived since, the retry says so neutrally rather than paying twice. "fetch-window": (spec) => { const { p, queueKey } = params(spec); const videoId = str(p.videoId); if (!videoId) { return Promise.resolve({ ok: false, error: "Job spec is missing its videoId.", }); } const from = num(p.from); const to = num(p.to); if (from === undefined || to === undefined) { return Promise.resolve({ ok: false, error: "Job spec is missing its window.", }); } return replayFetchWindowAction({ slug: spec.slug, videoId, webpageUrl: str(p.webpageUrl), from, to, provenance: { requestedBy: str(p.requestedBy) ?? "unknown", manifest: str(p.manifest), clipId: str(p.clipId), reason: str(p.reason), pad: num(p.pad), requestedAt: str(p.requestedAt), }, queueKey, maxHeight: num(p.maxHeight), }); }, // A batch of windows. Replay re-runs the same list on the same queue; the // controller's cache check skips every window an earlier run fetched. "fetch-windows": (spec) => { const { p, queueKey } = params(spec); const items = windowItems(p.items); if (items.length === 0) { return Promise.resolve({ ok: false, error: "Job spec has no windows." }); } if (queueKey === undefined) { return Promise.resolve({ ok: false, error: "Job spec is missing its queue." }); } return fetchWindowsJobAction({ items, requestedBy: str(p.requestedBy) ?? "unknown", manifest: str(p.manifest), siteId: str(p.siteId), maxHeight: num(p.maxHeight), queueKey, }); }, // Both digest lanes replay through one action; the lane comes from params so a // replayed metered run stays metered (and is refused if the lane has since // been turned off, rather than quietly falling back to local). "digest-channel-local": (spec) => { const { p, queueKey } = params(spec); return digestChannelAction( spec.slug, (str(p.lane) as DigestLaneChoice | undefined) ?? "local", queueKey, str(p.order), num(p.limitCount), bool(p.force), strings(p.ids), ); }, "digest-channel-remote": (spec) => { const { p, queueKey } = params(spec); return digestChannelAction( spec.slug, (str(p.lane) as DigestLaneChoice | undefined) ?? "remote", queueKey, str(p.order), num(p.limitCount), bool(p.force), strings(p.ids), ); }, "whisper-all": (spec) => { const { p, queueKey } = params(spec); return transcribeMissingAction( spec.slug, queueKey, bool(p.reverse), str(p.audioFormat) as AudioFormat | undefined, bool(p.strictAudioFormat), num(p.shardTotal), num(p.shardIndex), ); }, "whisper-bucket-downloaded-no-transcript": async (spec) => { const { p, queueKey } = params(spec); if (!spec.bucket) return { ok: false, error: "Job spec is missing its bucket." }; const ids = await idsForBucket(spec.slug, spec.bucket); if (ids.length === 0) { return { ok: false, error: "Nothing to transcribe right now.", info: true }; } return transcribeBucketAction( spec.slug, ids, queueKey, str(p.audioFormat) as AudioFormat | undefined, bool(p.strictAudioFormat), spec.bucket, ); }, "redownload-incomplete-bucket": async (spec) => { const { queueKey } = params(spec); if (!spec.bucket) return { ok: false, error: "Job spec is missing its bucket." }; const ids = await idsForBucket(spec.slug, spec.bucket); if (ids.length === 0) { return { ok: false, error: "Nothing to re-download right now.", info: true }; } // The same job kind drives both the incompleteTranscript and shortAudio // buckets; route by the captured bucket so the re-run re-stamps the right one. return spec.bucket === "shortAudio" ? redownloadShortAudioBucketAction(spec.slug, ids, queueKey) : redownloadIncompleteBucketAction(spec.slug, ids, queueKey); }, "retry-bucket": async (spec) => { const { p, queueKey } = params(spec); if (!spec.bucket) return { ok: false, error: "Job spec is missing its bucket." }; const ids = await idsForBucket(spec.slug, spec.bucket); if (ids.length === 0) { return { ok: false, error: "Nothing to retry right now.", info: true }; } return retryBucketAction( spec.slug, ids, queueKey, bool(p.abortOnError), str(p.handlingOverride), spec.bucket, Boolean(p.forceCookies), Boolean(p.replaceAutoSubs), ); }, "download-from-playlist": (spec) => { const { p, queueKey } = params(spec); return downloadAction(spec.slug, queueKey, bool(p.abortOnError)); }, "download-missing": (spec) => { const { p, queueKey } = params(spec); return downloadMissingAction( spec.slug, queueKey, bool(p.ignoreArchive), bool(p.abortOnError), num(p.shardTotal), num(p.shardIndex), bool(p.keepSourceVideoOverride), bool(p.extractImmediately), str(p.audioFormatOverride) as AudioFormat | undefined, ); }, sync: (spec) => { const { p, queueKey } = params(spec); return syncAction(spec.slug, queueKey, bool(p.fullSweep)); }, // Re-derives its targets from the channel's CURRENT state, like every other // replay here: a re-run scans whatever is still unscanned now, not the list // the first run was given. "metadata-scan": (spec) => { const { queueKey } = params(spec); return runMetadataScanAction(spec.slug, queueKey); }, // One video's metadata re-read. The action re-resolves the target and // re-asks the platform's hold and cooldown, as a click does. "refresh-metadata": (spec) => { const { p, queueKey } = params(spec); const videoId = str(p.videoId); if (!videoId) { return Promise.resolve({ ok: false, error: "Job spec is missing its videoId.", }); } return refreshVideoMetadataAction(spec.slug, videoId, queueKey); }, "check-post-availability": (spec) => { const { p, queueKey } = params(spec); return checkPostAvailabilityAction( spec.slug, queueKey, str(p.mode) as "stale" | "unchecked" | "all" | undefined, num(p.limit), ); }, "fetch-posts": (spec) => { const { p, queueKey } = params(spec); return fetchPostsAction( spec.slug, queueKey, bool(p.full), num(p.limit), bool(p.older), str(p.floor), bool(p.force), num(p.pages), str(p.from), ); }, // The ids are the spec's own (a capture is OF specific posts, unlike a // bucket); the re-run skips whatever the first run already captured. "capture-posts": (spec) => { const { p, queueKey } = params(spec); return capturePostsAction( spec.slug, strings(p.ids) ?? [], queueKey, bool(p.shots), bool(p.media), bool(p.force), bool(p.articles), ); }, "download-missing-subs": (spec) => { const { p, queueKey } = params(spec); return downloadMissingSubsAction(spec.slug, queueKey, bool(p.abortOnError)); }, "store-playlist": (spec) => { const { queueKey } = params(spec); return storePlaylistAction(spec.slug, queueKey); }, "clear-failed-transcriptions": (spec) => { const { queueKey } = params(spec); return clearFailedTranscriptionsAction(spec.slug, queueKey); }, "clean-extra-audio-formats": (spec) => { const { queueKey } = params(spec); return cleanExtraAudioFormatsAction(spec.slug, queueKey); }, "remove-wrong-format-audio": (spec) => { const { queueKey } = params(spec); return removeWrongFormatAudioAction(spec.slug, queueKey); }, "whisper-bucket-auto-subs": async (spec) => { const { p, queueKey } = params(spec); if (!spec.bucket) return { ok: false, error: "Job spec is missing its bucket." }; const ids = await idsForBucket(spec.slug, spec.bucket); if (ids.length === 0) { return { ok: false, error: "No auto-caption-only videos with audio right now.", info: true, }; } return transcribeAutoSubsBucketAction( spec.slug, ids, queueKey, str(p.audioFormat) as AudioFormat | undefined, bool(p.strictAudioFormat), ); }, "purge-superseded-auto-subs": (spec) => { const { queueKey } = params(spec); return purgeSupersededAutoSubsAction(spec.slug, queueKey); }, "clean-audio-transcribed": (spec) => { const { queueKey } = params(spec); return cleanAudioAction(spec.slug, queueKey); }, "diarize-channel": (spec) => { const { queueKey } = params(spec); return diarizeChannelAction(spec.slug, queueKey); }, // Replays correctly with nothing remembered: the batch re-derives its whole // work-list from disk, so a replay does what is missing NOW rather than what // was missing when the record was written. The kind scope is carried through // so a replayed single-kind run stays single-kind — and the ID scope for the // same reason: a replayed per-video run stays per-video rather than widening // to the whole channel. "backfill-channel": (spec) => { const { p, queueKey } = params(spec); const kindIds = Array.isArray(p.kindIds) ? p.kindIds.filter((k): k is string => typeof k === "string") : undefined; return backfillChannelAction(spec.slug, queueKey, kindIds, strings(p.ids)); }, "check-kept-deleted": (spec) => { const { queueKey } = params(spec); return checkKeptDeletedAction(spec.slug, queueKey); }, "persist-kept": (spec) => { const { queueKey } = params(spec); return persistKeptAction(spec.slug, queueKey); }, // The same list, the same options: a replay fetches whatever of it is still // not saved (the controller skips the rest). "persist-videos": (spec) => { const { p, queueKey } = params(spec); return persistChannelVideosAction(spec.slug, strings(p.ids) ?? [], { queueKey, format: isSourceVideoQuality(p.format) ? p.format : undefined, replace: isPersistReplacePolicy(p.replace) ? p.replace : undefined, gapMs: num(p.gapMs), minFreeMemMb: num(p.minFreeMemMb), }); }, };