commit 512de151cf395f602530b4bc2434bb3e18f51b24
parent 400adc37096fffe053b191ec8a679f19bd887577
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Sun, 20 Sep 2026 19:19:51 -0400
editor: POST /api/media/fetch-window sources clip media for another tool
umtool asks here instead of running yt-dlp: the fetch takes the channel's
cookie policy and extra args, the per-platform 429 cooldown and its backoff,
the low-disk gate and a job log, and the bytes land in the corpus where the
next tool reuses them.
Same door as /api/worker/* — the shared WORKER_TOKEN — deliberately: this
endpoint spends a source's patience, so it must be no easier to reach than the
executor protocol, and an instance with no token answers 503 rather than being
an open proxy by accident.
It is NOT gated by the download pause. That control exists to stop the lanes
spending the source's patience on their own schedule; this is an operator
asking by hand, through a second tool, for one window they are about to watch.
`requestedBy` is required. The whole point of routing the fetch through here is
that the bytes carry a reason, and an anonymous window is one nobody can
explain in six months.
fetchWindowAction has two entry points because it has two callers: the route
wants a JSON outcome and polls the job, Retry wants the StreamActionResult the
/jobs page streams. They differ only in what they do with the started job.
redownloadToArchiveAction's body is factored into archiveSourceVideo so the
full-source branch adds an origin rather than a second copy of it.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Diffstat:
4 files changed, 684 insertions(+), 1 deletion(-)
diff --git a/editor/app/api/media/fetch-window/[jobId]/route.ts b/editor/app/api/media/fetch-window/[jobId]/route.ts
@@ -0,0 +1,101 @@
+import path from "node:path";
+import { stat } from "node:fs/promises";
+import { NextResponse } from "next/server";
+import { getPaths } from "yt-dlp-transcript-common/lib/paths";
+import { authorizeWorkerRequest } from "yt-dlp-transcript-common/lib/workerToken";
+import { getRegistry } from "yt-dlp-transcript-common/jobs/registry";
+import { readJobMeta } from "yt-dlp-transcript-common/jobs/jobMeta";
+import { getJobEntry, readLogChunk } from "yt-dlp-transcript-common/jobs/listJobs";
+import { clipWindowPath } from "yt-dlp-transcript-common/lib/clipWindow-server";
+import { loadSavedVideo } from "yt-dlp-transcript-common/lib/savedVideo-server";
+import { savedVideoPath } from "yt-dlp-transcript-common/lib/savedVideo";
+
+export const dynamic = "force-dynamic";
+
+// POLL ONE FETCH. The POST that started it returned a jobId; this says whether
+// it finished and, when it did, what file came out.
+//
+// The file is recomputed from the job's own SPEC, never echoed back from the
+// caller: the spec is what the editor itself queued, so a poll cannot be talked
+// into naming a path the fetch never wrote.
+
+const ERROR_TAIL_BYTES = 2000;
+
+export async function GET(
+ request: Request,
+ { params }: { params: Promise<{ jobId: string }> },
+) {
+ const auth = authorizeWorkerRequest(request.headers.get("authorization"));
+ if (!auth.ok) {
+ return NextResponse.json({ error: auth.error }, { status: auth.status });
+ }
+ const { jobId } = await params;
+ if (!/^[A-Za-z0-9_-]+$/.test(jobId)) {
+ return NextResponse.json({ error: "Invalid job id" }, { status: 400 });
+ }
+ const paths = getPaths();
+
+ // The live record first (it has the spec and the true status); then the
+ // sidecar, which survives eviction and a restart and also carries the spec;
+ // then the log-only fallback, which can still answer "what happened" for a
+ // job whose sidecar is gone.
+ const live = getRegistry().get(jobId);
+ const meta = live ? null : await readJobMeta(paths, jobId);
+ const entry = live || meta ? null : await getJobEntry(paths, jobId);
+ if (!live && !meta && !entry) {
+ return NextResponse.json({ error: "no such job" }, { status: 404 });
+ }
+ const kind = live?.kind ?? meta?.kind ?? entry?.kind;
+ if (kind !== "fetch-window" && kind !== "redownload-archive") {
+ return NextResponse.json(
+ { error: `job ${jobId} is a ${kind ?? "?"} job, not a media fetch` },
+ { status: 400 },
+ );
+ }
+ const status = live?.status ?? meta?.status ?? entry?.status ?? "queued";
+ const exitCode = live?.exitCode ?? meta?.exitCode ?? entry?.exitCode;
+ const spec = live?.spec ?? meta?.spec ?? null;
+
+ const out: Record<string, unknown> = { status, jobId };
+ if (exitCode !== undefined) out.exitCode = exitCode;
+
+ if (status === "done") {
+ const slug = live?.channelSlug ?? meta?.channelSlug ?? spec?.slug;
+ const videoId = live?.videoId ?? meta?.videoId;
+ const p = spec?.params ?? {};
+ if (slug && videoId) {
+ const videoDir = path.join(paths.channelsDir, slug, "data", videoId);
+ let file: string | null = null;
+ if (kind === "fetch-window") {
+ const from = Number(p.from);
+ const to = Number(p.to);
+ if (Number.isFinite(from) && Number.isFinite(to)) {
+ file = clipWindowPath(videoDir, from, to);
+ out.from = from;
+ out.to = to;
+ }
+ } else {
+ const pointer = await loadSavedVideo(videoDir);
+ file = pointer ? savedVideoPath(pointer) : null;
+ }
+ if (file) {
+ const st = await stat(file).catch(() => null);
+ if (st?.isFile()) {
+ out.file = file;
+ out.bytes = st.size;
+ }
+ }
+ }
+ }
+
+ if (status === "failed" || status === "cancelled") {
+ // The tail of the job's own log. A caller that cannot see this server's
+ // /jobs page still gets the yt-dlp line that explains the refusal.
+ const logPath =
+ live?.logPath ?? entry?.logPath ?? path.join(paths.jobsDir, `${jobId}.log`);
+ const { content } = await readLogChunk(logPath, 0);
+ out.error = content.slice(-ERROR_TAIL_BYTES).trim();
+ }
+
+ return NextResponse.json(out, { headers: { "cache-control": "no-store" } });
+}
diff --git a/editor/app/api/media/fetch-window/route.ts b/editor/app/api/media/fetch-window/route.ts
@@ -0,0 +1,182 @@
+import { NextResponse } from "next/server";
+import { authorizeWorkerRequest } from "yt-dlp-transcript-common/lib/workerToken";
+import {
+ MAX_WINDOW_SECONDS,
+ fetchFullSourceAction,
+ fetchWindowAction,
+ type FetchMediaOutcome,
+} from "../../../channels/[slug]/videos/[id]/videoActions";
+
+export const dynamic = "force-dynamic";
+
+// SOURCE MEDIA FOR ANOTHER TOOL, fetched through the editor's managed download
+// path.
+//
+// umtool's clip bench ran yt-dlp itself. It now posts here, and the bytes land
+// in the corpus beside the video — cookie policy, per-platform sleeps, 429
+// cooldown and backoff included — where the next tool and the next build reuse
+// them instead of paying for the same seconds twice.
+//
+// SAME DOOR AS /api/worker/*: the shared WORKER_TOKEN. Deliberately, and not
+// laziness — this endpoint spends the source's patience, so it must be no
+// easier to reach than the executor protocol, and an instance with no token is
+// not one by accident (503, not an open proxy).
+//
+// The download PAUSE does not gate it; see fetchWindowAction for why.
+
+const ID_RE = /^[\w.-]+$/;
+
+type Body = {
+ channelSlug?: unknown;
+ videoId?: unknown;
+ webpageUrl?: unknown;
+ from?: unknown;
+ to?: unknown;
+ requestedBy?: unknown;
+ manifest?: unknown;
+ clipId?: unknown;
+ reason?: unknown;
+ pad?: unknown;
+ full?: unknown;
+};
+
+const str = (v: unknown): string | undefined =>
+ typeof v === "string" && v.trim() !== "" ? v.trim() : undefined;
+
+function answer(outcome: FetchMediaOutcome): NextResponse {
+ if (!outcome.ok) {
+ const body: Record<string, unknown> = { error: outcome.error };
+ if (outcome.cooldownMs !== undefined) body.cooldownMs = outcome.cooldownMs;
+ if (outcome.platform !== undefined) body.platform = outcome.platform;
+ return NextResponse.json(body, { status: outcome.status });
+ }
+ // `started` is the Retry shape (a live StreamActionResult); the route never
+ // asks for it, so reaching it here would be a programming error, not a state
+ // to render.
+ if ("started" in outcome) {
+ return NextResponse.json(
+ { error: "internal: streamed outcome on the HTTP path" },
+ { status: 500 },
+ );
+ }
+ if (outcome.cached) {
+ return NextResponse.json(
+ {
+ cached: true,
+ file: outcome.file,
+ from: outcome.from,
+ to: outcome.to,
+ bytes: outcome.bytes,
+ provenance: outcome.provenance,
+ },
+ { status: 200 },
+ );
+ }
+ return NextResponse.json(
+ {
+ cached: false,
+ jobId: outcome.jobId,
+ file: outcome.file,
+ from: outcome.from,
+ to: outcome.to,
+ },
+ { status: 202 },
+ );
+}
+
+export async function POST(request: Request) {
+ const auth = authorizeWorkerRequest(request.headers.get("authorization"));
+ if (!auth.ok) {
+ return NextResponse.json({ error: auth.error }, { status: auth.status });
+ }
+ let body: Body;
+ try {
+ body = (await request.json()) as Body;
+ } catch {
+ return NextResponse.json({ error: "malformed JSON body" }, { status: 400 });
+ }
+
+ const channelSlug = str(body.channelSlug);
+ const videoId = str(body.videoId);
+ // Both land in a filesystem path. Anchored, so `..` and a separator are
+ // refused at the door rather than by whatever the path math happens to do.
+ if (!channelSlug || !ID_RE.test(channelSlug)) {
+ return NextResponse.json(
+ { error: "channelSlug is required and must match /^[\\w.-]+$/" },
+ { status: 400 },
+ );
+ }
+ if (!videoId || !ID_RE.test(videoId)) {
+ return NextResponse.json(
+ { error: "videoId is required and must match /^[\\w.-]+$/" },
+ { status: 400 },
+ );
+ }
+ // WHO ASKED IS NOT OPTIONAL. The whole point of routing the fetch through
+ // here is that the bytes carry a reason; an anonymous window is one nobody
+ // can explain in six months.
+ const requestedBy = str(body.requestedBy);
+ if (!requestedBy) {
+ return NextResponse.json(
+ { error: "requestedBy is required (the tool asking for these bytes)" },
+ { status: 400 },
+ );
+ }
+ const provenance = {
+ requestedBy,
+ manifest: str(body.manifest),
+ clipId: str(body.clipId),
+ reason: str(body.reason),
+ pad: typeof body.pad === "number" && Number.isFinite(body.pad)
+ ? body.pad
+ : undefined,
+ requestedAt: new Date().toISOString(),
+ };
+
+ if (body.full === true) {
+ return answer(
+ await fetchFullSourceAction({
+ slug: channelSlug,
+ videoId,
+ provenance,
+ }),
+ );
+ }
+
+ const from = Number(body.from);
+ const to = Number(body.to);
+ if (!Number.isFinite(from) || !Number.isFinite(to) || from < 0) {
+ return NextResponse.json(
+ { error: "from and to must be finite, non-negative seconds" },
+ { status: 400 },
+ );
+ }
+ if (from >= to) {
+ return NextResponse.json(
+ { error: `from (${from}) must be less than to (${to})` },
+ { status: 400 },
+ );
+ }
+ if (to - from > MAX_WINDOW_SECONDS) {
+ return NextResponse.json(
+ {
+ error:
+ `a window may be at most ${MAX_WINDOW_SECONDS}s ` +
+ `(asked for ${Math.round(to - from)}s); pass full: true for the ` +
+ `whole recording`,
+ },
+ { status: 400 },
+ );
+ }
+
+ return answer(
+ await fetchWindowAction({
+ slug: channelSlug,
+ videoId,
+ webpageUrl: str(body.webpageUrl),
+ from,
+ to,
+ provenance,
+ }),
+ );
+}
diff --git a/editor/app/channels/[slug]/videos/[id]/videoActions.ts b/editor/app/channels/[slug]/videos/[id]/videoActions.ts
@@ -34,7 +34,28 @@ import { pruneFailedTranscriptions } from "yt-dlp-transcript-common/controller/f
import { transcodeAudio } from "yt-dlp-transcript-common/controller/transcode";
import { transcribeWithWorker } from "yt-dlp-transcript-common/controller/transcribeOne";
import { findVideoSourceUrl } from "yt-dlp-transcript-common/controller/undownloadedVideos";
-import { unpersistSavedVideo } from "yt-dlp-transcript-common/lib/savedVideo-server";
+import {
+ loadSavedVideo,
+ unpersistSavedVideo,
+} from "yt-dlp-transcript-common/lib/savedVideo-server";
+import {
+ savedVideoPath,
+ type SavedVideoOrigin,
+} from "yt-dlp-transcript-common/lib/savedVideo";
+import type { ClipProvenance } from "yt-dlp-transcript-common/lib/clipWindow";
+import {
+ clipWindowPath,
+ findContainingClipWindow,
+} from "yt-dlp-transcript-common/lib/clipWindow-server";
+import {
+ fetchWindowManaged,
+ type FetchWindowProvenance,
+} from "yt-dlp-transcript-common/ytdlp/fetchWindowManaged";
+import { detectPlatform } from "yt-dlp-transcript-common/lib/platform";
+import {
+ platformCooldownRemainingMs,
+ recordDownloadBackoff,
+} from "yt-dlp-transcript-common/jobs/downloadBackoff";
import { getSettings } from "yt-dlp-transcript-common/lib/settings";
import { resolveCookiePolicy } from "yt-dlp-transcript-common/lib/cookiePolicy";
import { downloadOneManaged } from "yt-dlp-transcript-common/ytdlp/downloadOneManaged";
@@ -207,6 +228,20 @@ export async function redownloadToArchiveAction(
videoId: string,
queueKey?: string,
): Promise<StreamActionResult> {
+ return archiveSourceVideo(slug, videoId, queueKey);
+}
+
+// The body of both "Persist source video" (the button) and the full-source
+// branch of the clip-window fetch (umtool asking for a whole recording rather
+// than a window). ONE function, because the two differ in exactly one field:
+// who asked. Factored rather than copied — a second copy is a second place for
+// the keepSourceVideoOverride / appendArchive pair to drift.
+async function archiveSourceVideo(
+ slug: string,
+ videoId: string,
+ queueKey: string | undefined,
+ persistOrigin?: SavedVideoOrigin,
+): Promise<StreamActionResult> {
const r = await loadConfigOrError(slug);
if (!r.ok) return r;
const paths = getPaths();
@@ -234,6 +269,9 @@ export async function redownloadToArchiveAction(
kind: "download",
});
try {
+ if (persistOrigin) {
+ onLog(`${requesterLine(persistOrigin)}\n`);
+ }
onLog(`Re-downloading ${videoId} to archive its source video…\n`);
await downloadOneManaged({
channelSlug: slug,
@@ -247,6 +285,7 @@ export async function redownloadToArchiveAction(
globalSkipLiveDownloads: settings.skipLiveDownloads,
appendArchive: true,
keepSourceVideoOverride: true,
+ persistOrigin,
});
revalidatePath(`/channels/${slug}/videos/${videoId}`);
revalidatePath(`/channels/${slug}`);
@@ -643,3 +682,326 @@ export async function unpersistVideoAction(
revalidatePath("/saved-videos");
return { ok: true };
}
+
+// ---------------------------------------------------------------------------
+// Sourcing media FOR ANOTHER TOOL, through this app's managed download path.
+//
+// The operator's rule is that no yt-dlp runs by hand. umtool's clip bench used
+// to fetch its own clip windows; it now asks here, so the fetch inherits the
+// channel's cookie policy and extra args, the per-platform 429 cooldown, the
+// low-disk gate and a job log — and the bytes land in the corpus where the next
+// tool (and the next build) can reuse them, with a note saying who asked.
+// ---------------------------------------------------------------------------
+
+// The one sentence a job log and a video page both show: who wanted this, for
+// what, and why. Built in one place so the two cannot disagree.
+function requesterLine(o: SavedVideoOrigin | FetchWindowProvenance): string {
+ const what = [o.manifest, o.clipId].filter(Boolean).join("/");
+ const head = what ? `${o.requestedBy} · ${what}` : o.requestedBy;
+ return o.reason ? `${head} — ${o.reason}` : head;
+}
+
+export type FetchMediaOutcome =
+ // Already on disk. No job, no bytes, no politeness owed.
+ | {
+ ok: true;
+ cached: true;
+ file: string;
+ from: number;
+ to: number;
+ bytes: number;
+ provenance: ClipProvenance | SavedVideoOrigin | null;
+ }
+ // Queued. `file` is where the window WILL be; for a full source it is not
+ // knowable until the container's real extension is (null until then).
+ | {
+ ok: true;
+ cached: false;
+ jobId: string;
+ file: string | null;
+ from: number;
+ to: number;
+ }
+ // `stream: true` was asked for: the caller reads the job's log itself, so the
+ // StreamActionResult is handed over rather than cancelled.
+ | {
+ ok: true;
+ cached: false;
+ started: StreamActionResult & { ok: true };
+ from: number;
+ to: number;
+ }
+ // `status` is the HTTP status the route should answer with, so the mapping
+ // from a refusal to a code lives with the refusal rather than in a switch
+ // over error strings.
+ | {
+ ok: false;
+ status: number;
+ error: string;
+ cooldownMs?: number;
+ platform?: string;
+ };
+
+// The widest window this will fetch in one go. Fifteen minutes is far past any
+// citation and well short of "you meant to download the video" — which is what
+// `full` is for, and which goes through the saved-video store instead.
+export const MAX_WINDOW_SECONDS = 900;
+
+// A started job's stream has no reader here: the caller polls the job instead.
+// Cancelling the stream is the documented way to say so — it marks the
+// controller closed and leaves the job running into its on-disk log.
+function detach(res: StreamActionResult & { ok: true }): void {
+ void res.stream.cancel().catch(() => {});
+}
+
+// Fetch ONE window of a video's source media into data/<id>/clips/.
+//
+// TWO CALLERS, TWO SHAPES. The HTTP route wants a JSON-able outcome and polls
+// the job; Retry (jobReplayRegistry) wants the StreamActionResult every other
+// action returns, so the /jobs page can stream the log. So the work is one
+// function and the two entry points differ only in what they do with the
+// started job — never in which cache it consulted or which refusal it hit.
+export async function fetchWindowAction(req: {
+ slug: string;
+ videoId: string;
+ // The page URL, when the caller has it. Otherwise it is resolved from the
+ // video's own metadata or the stored playlist, exactly as the archive
+ // re-download does.
+ webpageUrl?: string;
+ from: number;
+ to: number;
+ provenance: FetchWindowProvenance;
+ queueKey?: string;
+ // Hand back the StreamActionResult instead of detaching it. Set by Retry,
+ // which renders the log; the HTTP route leaves it off and polls.
+ stream?: boolean;
+}): Promise<FetchMediaOutcome> {
+ const { slug, videoId, from, to } = req;
+ const r = await loadConfigOrError(slug);
+ if (!r.ok) return { ok: false, status: 404, error: r.error };
+ const paths = getPaths();
+ const videoDir = videoDirOf(slug, videoId);
+
+ // ASK THE CACHE FIRST, and accept a WIDER file: a report cites the same
+ // stream more than once, so a generous fetch for one clip must serve the
+ // neighbour it already covers rather than being downloaded again.
+ const hit = await findContainingClipWindow(videoDir, from, to);
+ if (hit) {
+ return {
+ ok: true,
+ cached: true,
+ file: hit.path,
+ from: hit.from,
+ to: hit.to,
+ bytes: hit.bytes,
+ provenance: hit.provenance,
+ };
+ }
+
+ const url =
+ req.webpageUrl?.trim() ||
+ (await findVideoSourceUrl(paths, slug, videoId, r.config));
+ if (!url) {
+ return {
+ ok: false,
+ status: 404,
+ error:
+ "Could not determine the video URL: no metadata.info.json and the " +
+ "playlist does not contain a matching entry.",
+ };
+ }
+
+ // The same per-platform cooldown a clicked Sync respects. A window is a
+ // handful of seconds, but a 429 is a fact about the SOURCE, not about the
+ // size of the request.
+ const platform = detectPlatform(url) ?? "unknown";
+ const cooldownMs = await platformCooldownRemainingMs(platform, paths);
+ if (cooldownMs > 0) {
+ return {
+ ok: false,
+ status: 409,
+ cooldownMs,
+ platform,
+ error:
+ `${platform} is in a rate-limit cooldown ` +
+ `(${Math.ceil(cooldownMs / 1000)}s remaining).`,
+ };
+ }
+
+ const disk = await lowDiskError(paths, slug);
+ if (disk) return { ok: false, status: 507, error: disk.error };
+
+ const settings = getSettings();
+ // NOT GATED BY THE DOWNLOAD PAUSE. The pause exists to stop the lanes — the
+ // auto-download runner and a channel-wide sync — from spending the source's
+ // patience on their own schedule. This is an operator asking, by hand,
+ // through a second tool, for one window they are about to watch; refusing it
+ // would make "pause downloads" mean "stop working", which is not what the
+ // control says and not why it gets used.
+ const res = await runManagedFunction({
+ kind: "fetch-window",
+ queueKey: videoQueueKey(r.config, req.queueKey),
+ paths,
+ channelSlug: slug,
+ videoId,
+ spec: {
+ kind: "fetch-window",
+ slug,
+ params: {
+ videoId,
+ from,
+ to,
+ webpageUrl: url,
+ queueKey: req.queueKey,
+ ...req.provenance,
+ },
+ },
+ fn: async (onLog, signal) => {
+ onLog(`${requesterLine(req.provenance)}\n`);
+ await fetchWindowManaged({
+ channelSlug: slug,
+ channelConfig: r.config,
+ paths,
+ videoDir,
+ videoId,
+ videoUrl: url,
+ from,
+ to,
+ provenance: req.provenance,
+ cookiePolicy: resolveCookiePolicy(settings, r.config),
+ onLog,
+ signal,
+ onPlatformBackoff: () => recordDownloadBackoff(platform, paths),
+ });
+ revalidatePath(`/channels/${slug}/videos/${videoId}`);
+ },
+ });
+ if (!res.ok) {
+ // An unreachable-media refusal (a relocated channel whose drive is not
+ // mounted, or one mid-relocation) arrives here verbatim from
+ // runManagedFunction's guard, and is passed on verbatim: the caller's
+ // operator is the person who can plug the drive in.
+ return { ok: false, status: 503, error: res.error };
+ }
+ if (req.stream) return { ok: true, cached: false, started: res, from, to };
+ detach(res);
+ return {
+ ok: true,
+ cached: false,
+ jobId: res.jobId,
+ file: clipWindowPath(videoDir, from, to),
+ from,
+ to,
+ };
+}
+
+// Retry, from the /jobs page. A replayed fetch re-derives from the CURRENT
+// disk: a window somebody has since fetched (or that a wider one now covers)
+// replays as a neutral "already there" rather than downloading it twice.
+export async function replayFetchWindowAction(req: {
+ slug: string;
+ videoId: string;
+ webpageUrl?: string;
+ from: number;
+ to: number;
+ provenance: FetchWindowProvenance;
+ queueKey?: string;
+}): Promise<StreamActionResult> {
+ const r = await fetchWindowAction({ ...req, stream: true });
+ if (!r.ok) return { ok: false, error: r.error, info: r.status === 409 };
+ if ("started" in r) return r.started;
+ if (r.cached) {
+ return {
+ ok: false,
+ info: true,
+ error:
+ `${r.from.toFixed(2)}–${r.to.toFixed(2)} is already fetched ` +
+ `(${r.file}); nothing to download.`,
+ };
+ }
+ // Unreachable: a non-cached result with `stream: true` always carries
+ // `started`. Kept total rather than asserted.
+ return { ok: false, info: true, error: "Nothing to do." };
+}
+
+// The WHOLE recording, when a window will not do (a tool that needs to re-cut
+// freely, or a source whose windows would tile the entire runtime). It reuses
+// the saved-video store rather than clips/: that is where big containers
+// already live, with a retention rule that leaves an explicitly-requested one
+// alone.
+export async function fetchFullSourceAction(req: {
+ slug: string;
+ videoId: string;
+ provenance: FetchWindowProvenance;
+ queueKey?: string;
+}): Promise<FetchMediaOutcome> {
+ const { slug, videoId } = req;
+ const r = await loadConfigOrError(slug);
+ if (!r.ok) return { ok: false, status: 404, error: r.error };
+ const paths = getPaths();
+ const videoDir = videoDirOf(slug, videoId);
+
+ const pointer = await loadSavedVideo(videoDir);
+ if (pointer) {
+ return {
+ ok: true,
+ cached: true,
+ file: savedVideoPath(pointer),
+ from: 0,
+ to: 0,
+ bytes: pointer.bytes,
+ provenance: pointer.origin ?? null,
+ };
+ }
+
+ const url = await findVideoSourceUrl(paths, slug, videoId, r.config);
+ if (!url) {
+ return {
+ ok: false,
+ status: 404,
+ error:
+ "Could not determine the video URL: no metadata.info.json and the " +
+ "playlist does not contain a matching entry.",
+ };
+ }
+ const platform = detectPlatform(url) ?? "unknown";
+ const cooldownMs = await platformCooldownRemainingMs(platform, paths);
+ if (cooldownMs > 0) {
+ return {
+ ok: false,
+ status: 409,
+ cooldownMs,
+ platform,
+ error:
+ `${platform} is in a rate-limit cooldown ` +
+ `(${Math.ceil(cooldownMs / 1000)}s remaining).`,
+ };
+ }
+
+ const origin: SavedVideoOrigin = {
+ requestedBy: req.provenance.requestedBy,
+ ...(req.provenance.manifest ? { manifest: req.provenance.manifest } : {}),
+ ...(req.provenance.clipId ? { clipId: req.provenance.clipId } : {}),
+ ...(req.provenance.reason ? { reason: req.provenance.reason } : {}),
+ requestedAt: req.provenance.requestedAt ?? new Date().toISOString(),
+ };
+ const res = await archiveSourceVideo(slug, videoId, req.queueKey, origin);
+ if (!res.ok) {
+ // archiveSourceVideo's own low-disk refusal is the only 507-shaped one it
+ // returns; everything else is a lookup failure or an unreachable drive.
+ const status = res.error.startsWith("Low disk space") ? 507 : 503;
+ return { ok: false, status, error: res.error };
+ }
+ detach(res);
+ return {
+ ok: true,
+ cached: false,
+ jobId: res.jobId,
+ // The container's extension is decided by the format yt-dlp picks, so the
+ // path is not knowable until the job has finished. The poll endpoint reads
+ // it off the saved-video pointer.
+ file: null,
+ from: 0,
+ to: 0,
+ };
+}
diff --git a/editor/app/jobs/jobReplayRegistry.ts b/editor/app/jobs/jobReplayRegistry.ts
@@ -47,6 +47,7 @@ import {
redownloadIncompleteBucketAction,
redownloadShortAudioBucketAction,
} from "../channels/[slug]/incompleteTranscriptActions";
+import { replayFetchWindowAction } from "../channels/[slug]/videos/[id]/videoActions";
export type ReplayHandler = (spec: JobSpec) => Promise<StreamActionResult>;
@@ -79,6 +80,43 @@ function params(spec: JobSpec): {
}
export const JOB_REPLAY_HANDLERS: Record<string, ReplayHandler> = {
+ // 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,
+ });
+ },
// 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).