commit c968ddf6c36ff18bdc434d87e8bb69487a74d648
parent d6afa774ba81b366cc9f5cdb3e2b27f68789a121
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Fri, 9 Oct 2026 12:09:23 -0400
editor: clip windows get their own queue per platform and are never queued twice
Release 19 slice A5. A fetch-window job ran on the platform's download queue,
and a queue runs one job at a time — a tier orders only the waiting jobs — so
a few seconds of video waited behind a multi-hour persist-videos or sync.
clipWindowQueueKey (lib/queueKeys.ts) puts the single fetch (fetchWindowAction)
and the batches (fetchWindowsAction) on `clips:<platform>`: one window at a
time per platform, the batches' clip-window gap kept, the platform's hold,
cooldown and backoff still shared.
A repeated ask is answered with the job already fetching it:
jobs/windowJobs.ts' windowInFlight finds a queued or running fetch-window
whose window contains the asked one, or a fetch-windows batch with such an
item, same channel, video and height cap, off the job's replay spec. The
single route answers that jobId with `existing: true`; the batch lists such
windows in `inFlight`, starts no job for them, and the ops route adds their
jobs to `jobIds` so --wait follows the windows asked for. Retry is not
deduplicated.
umtool/report-to-video/fetch-via-editor.mjs (its own wait and 10-minute poll
timeout) is not touched: that directory belongs to another session.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Diffstat:
12 files changed, 364 insertions(+), 24 deletions(-)
diff --git a/common/jobs/jobKinds.ts b/common/jobs/jobKinds.ts
@@ -294,10 +294,11 @@ const JOB_KINDS: Record<string, JobKindMeta> = {
needsMedia: true,
},
// ONE WINDOW of a video's source media, fetched into data/<id>/clips/ for a
- // tool that asked for it by name (umtool's clip bench). Platform-queued like
- // every other fetch so it takes its turn behind the channel's own downloads;
- // replayable because the request is a few numbers and a reason, which is
- // exactly what a JobSpec holds.
+ // tool that asked for it by name (umtool's clip bench). On its platform's
+ // CLIP queue (`clips:<platform>`, lib/queueKeys.ts clipWindowQueueKey —
+ // release 19, A5), one window at a time per platform, never behind the
+ // platform's long downloads; replayable because the request is a few numbers
+ // and a reason, which is exactly what a JobSpec holds.
"fetch-window": {
kind: "fetch-window",
label: "Fetch window",
@@ -309,8 +310,9 @@ const JOB_KINDS: Record<string, JobKindMeta> = {
},
// A LIST OF CLIP WINDOWS, one platform's, as one paced job
// (controller/fetchWindows.ts): a site's missing evidence, a umtool
- // manifest's timeline, an explicit list. One job per platform queue, so
- // YouTube and Rumble run side by side, each paced on its own. Drainable: it
+ // manifest's timeline, an explicit list. One job per platform's clip queue
+ // (`clips:<platform>`, as fetch-window), so YouTube and Rumble run side by
+ // side, each paced on its own. Drainable: it
// stops between windows and a re-run fetches only what is still missing,
// which is also why it is replayable. It spans channels and so starts with no
// channelSlug — the per-channel text check `needsText` stands for is made by
diff --git a/common/jobs/windowJobs.test.ts b/common/jobs/windowJobs.test.ts
@@ -0,0 +1,71 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import type { JobRecord } from "./registry";
+import { windowInFlight } from "./windowJobs";
+
+// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/windowJobs.test.ts
+//
+// Which queued or running job will already write a requested clip window
+// (release 19, A5): a single fetch-window whose window contains it, or a
+// fetch-windows batch with such an item — same channel, same video, same
+// height cap — read off the job's replay spec.
+
+function job(
+ id: string,
+ kind: string,
+ status: JobRecord["status"],
+ params: Record<string, unknown> | null,
+ slug = "chan",
+): JobRecord {
+ return {
+ id,
+ kind,
+ queueKey: "clips:youtube",
+ status,
+ queuedAt: 0,
+ logPath: "/dev/null",
+ ...(params ? { spec: { kind, slug, params } } : {}),
+ } as JobRecord;
+}
+
+const want = { slug: "chan", videoId: "v1", from: 100, to: 110 };
+
+test("a single job whose window contains the asked one is in flight; its own window is answered", () => {
+ const hit = windowInFlight([job("A", "fetch-window", "queued", { videoId: "v1", from: 95, to: 130 })], want);
+ assert.deepEqual(hit, { jobId: "A", kind: "fetch-window", status: "queued", from: 95, to: 130 });
+ // Exactly the same window, running.
+ assert.equal(
+ windowInFlight([job("B", "fetch-window", "running", { videoId: "v1", from: 100, to: 110 })], want)?.jobId,
+ "B",
+ );
+});
+
+test("a window that only overlaps, another video, another channel or an ended job is not", () => {
+ const records = [
+ job("overlap", "fetch-window", "queued", { videoId: "v1", from: 105, to: 130 }),
+ job("video", "fetch-window", "queued", { videoId: "v2", from: 0, to: 900 }),
+ job("chan", "fetch-window", "queued", { videoId: "v1", from: 0, to: 900 }, "other"),
+ job("done", "fetch-window", "done", { videoId: "v1", from: 0, to: 900 }),
+ job("nospec", "fetch-window", "queued", null),
+ job("sync", "sync", "running", { videoId: "v1", from: 0, to: 900 }),
+ ];
+ assert.equal(windowInFlight(records, want), null);
+});
+
+test("a batch with a covering item is in flight; the height cap must match", () => {
+ const batch = job("C", "fetch-windows", "running", {
+ items: [
+ { slug: "chan", id: "v9", from: 0, to: 10 },
+ { slug: "chan", id: "v1", from: 90, to: 120 },
+ ],
+ maxHeight: 480,
+ });
+ assert.equal(windowInFlight([batch], want), null, "no cap asked, 480 held");
+ assert.deepEqual(windowInFlight([batch], { ...want, maxHeight: 480 }), {
+ jobId: "C",
+ kind: "fetch-windows",
+ status: "running",
+ from: 90,
+ to: 120,
+ });
+});
diff --git a/common/jobs/windowJobs.ts b/common/jobs/windowJobs.ts
@@ -0,0 +1,71 @@
+import type { JobRecord } from "./registry";
+
+// THE CLIP-WINDOW JOB ALREADY ON ITS WAY (release 19, A5).
+//
+// A tool that asks for a window and does not get it at once asks again — umtool's
+// fetch-via-editor after a timeout, an agent retrying, the same manifest sent
+// twice — and every ask used to queue another job for the same seconds of the
+// same video: the second waited behind the first, found nothing to fetch when
+// it ran, and meanwhile both sat in the platform's queue in front of real work.
+// The disk cache (`findContainingClipWindow`) only answers once a file exists;
+// this answers the window that is still being fetched.
+//
+// A request is IN FLIGHT when a queued or running job will produce a file that
+// covers it: a single `fetch-window` job for the same channel and video whose
+// window CONTAINS the asked one (the cache accepts a wider file, so this does
+// too), or a `fetch-windows` batch with such an item. The source height cap
+// must match: a window capped at 480 is not an answer to a request for 720.
+// Read off the job's replay spec — the request it was started with — so a
+// record without one (an old job) is never a match.
+
+export type WantedWindow = {
+ slug: string;
+ videoId: string;
+ from: number;
+ to: number;
+ maxHeight?: number;
+};
+
+export type WindowInFlight = {
+ jobId: string;
+ kind: "fetch-window" | "fetch-windows";
+ status: "queued" | "running";
+ // The window the job will write, which contains the asked one.
+ from: number;
+ to: number;
+};
+
+const num = (v: unknown): number | undefined =>
+ typeof v === "number" && Number.isFinite(v) ? v : undefined;
+
+export function windowInFlight(
+ records: readonly JobRecord[],
+ want: WantedWindow,
+): WindowInFlight | null {
+ for (const r of records) {
+ if (r.status !== "queued" && r.status !== "running") continue;
+ const p = (r.spec?.params ?? {}) as Record<string, unknown>;
+ if (num(p.maxHeight) !== want.maxHeight) continue;
+ const covers = (from: number | undefined, to: number | undefined) =>
+ from !== undefined && to !== undefined && from <= want.from && to >= want.to;
+ if (r.kind === "fetch-window") {
+ if (r.spec?.slug !== want.slug || p.videoId !== want.videoId) continue;
+ const from = num(p.from);
+ const to = num(p.to);
+ if (covers(from, to)) {
+ return { jobId: r.id, kind: "fetch-window", status: r.status, from: from!, to: to! };
+ }
+ } else if (r.kind === "fetch-windows" && Array.isArray(p.items)) {
+ for (const raw of p.items) {
+ const it = (raw ?? {}) as Record<string, unknown>;
+ if (it.slug !== want.slug || it.id !== want.videoId) continue;
+ const from = num(it.from);
+ const to = num(it.to);
+ if (covers(from, to)) {
+ return { jobId: r.id, kind: "fetch-windows", status: r.status, from: from!, to: to! };
+ }
+ }
+ }
+ }
+ return null;
+}
diff --git a/common/lib/queueKeys.ts b/common/lib/queueKeys.ts
@@ -76,6 +76,30 @@ export function downloadQueueKey(config: ChannelConfig): string {
: queueKeyForUrl(config.url);
}
+// CLIP WINDOWS HAVE THEIR OWN QUEUE PER PLATFORM (release 19, A5):
+// `clips:youtube` beside `platform:youtube`.
+//
+// A window is a few seconds of a video, asked for by a tool whose operator is
+// waiting on it. On the platform's download queue it waited behind whatever
+// held that queue — a multi-hour `persist-videos`, a channel sync — because a
+// queue runs one job at a time and a tier only orders the jobs still waiting,
+// never the one running. So the windows of a platform are serialized against
+// EACH OTHER (one clip stream per platform, the batches' 30–45 s gap between
+// two windows kept by `clip-window:<platform>` in jobs/platformGap.ts), and no
+// longer against the platform's long downloads.
+//
+// THE PLATFORM'S STATE IS STILL SHARED: a hold or a rate-limit cooldown on the
+// platform refuses a window as it refuses a download, and a 429 on either backs
+// the platform off for both. What changes is that one window may be fetched
+// while one download runs on the same platform.
+//
+// Takes the platform queue a window would otherwise have run on
+// (`downloadQueueKey`, `queueKeyForUrl`), so an unrecognized host keeps its
+// per-domain key: `platform:vimeo.com` → `clips:vimeo.com`.
+export function clipWindowQueueKey(platformQueue: string): string {
+ return `clips:${platformQueue.replace(/^platform:/, "")}`;
+}
+
// Single definition of how a UI-supplied queue override resolves against a
// feature's default queue. `undefined` (no control wired up) → the default;
// `""` → immediate (the registry runs `queueKey === ""` right away); any other
diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md
@@ -6,6 +6,7 @@
- **What the editor's pages show, over `pnpm ops`.** `get settings [<key>]` is settings.json as the editor reads it — migrated and defaulted — or one top-level key of it; `get storage` the /storage page (each location, mounted or not, its free space and tiers); `get sites` every site's id, title, public URL, Pages project, audience, whether it is listed, publish policy and channels; `get workers`, `get auto-queue` and `get scheduler` the payloads /workers, the four lanes and /operations/sync poll; `get cleanup <slug>` one channel's /cleanup row — what each sweep would reclaim, what holds the rest, the failed-transcriptions count. All behind the token, and no secret leaves through any of them: a remote worker's token reads `<redacted>`. Needs a restart of the editor.
- **Settings, lanes, workers, single videos and cleanup over `pnpm ops`.** `pnpm ops settings {"patch": {…}}` writes settings.json through the editor's one writer (a block's keys merge one level): an unknown key is refused, as is a key another writer owns (`channelPriority`, `autoQueue`, `workers`, `storage` — the refusal names the command or page), and so is a value the schema would not keep as sent — a clamped number, a dropped choice — named with what would have been saved; the answer is the saved values. `lane` takes the publish lane too (its switch, its hold, and Start, Drain and Stop, as on /operations/publish). `clear-platform-hold {"platform"}` is the lane strip's Clear hold. `workers {"op": "enable" | "disable" | "drain", "ids"}` is /workers' switch and Drain. `transcribe-one {"slug", "id", "file"?}` transcribes one video — a named audio file, or the video page's Transcribe, which fetches the audio first when there is none; `delete-file {"slug", "id", "file"}` is the Files list's delete (a tiered file goes with its bytes on the media drive, and nothing is deleted while that drive is not answering); `do-not-clean {"slug", "id", "keep"?}` sets the marker. `cleanup {"slug", "sweep"}` runs one of the channel's cleanup sweeps as a job (`transcribed`, `extra-formats`, `wrong-format`, `failed-transcriptions`). `relocate` takes `"dryRun": true`: each channel's preview — bytes to copy, free space on both sides — and no move (a channel never tiered is tiered in place first, as the Storage panel's preview does). Needs a restart of the editor.
- **`archilyzer storage report [<slug>…] [--json]`** lists each channel's media (the tierable files), text and clip-window bytes and where its media is — the corpus disk, a storage location, another path, or the retired layout — with totals per place; the corpus disk's total is what moving media would free. It reads the channels' reports and touches no drive, so an unmounted one cannot stall it; a figure a report does not have is shown as unknown, never 0. Needs no editor.
+- **A clip window no longer waits behind a long download, and is never queued twice.** Clip windows — `/api/media/fetch-window`, `pnpm ops fetch-windows`, a site's Fetch missing evidence, umtool's fetches — run on their platform's own clip queue (`clips:youtube`, `clips:rumble`, … on /jobs) instead of the platform's download queue, where a few seconds of video waited behind a multi-hour Persist or Sync. Windows still go one at a time per platform with the batches' gap between them, and the platform's hold, cooldown and backoff still apply to them: a window may now be fetched while one download runs on the same platform. Asking again for a window a queued or running job will already write — the same window, or one inside a wider one, at the same height cap — answers with that job instead of starting another: `/api/media/fetch-window` answers its `jobId` with `"existing": true`, and `fetch-windows` lists such windows under `inFlight` and adds their jobs to `jobIds`, so `--wait` follows them. A Retry from /jobs still re-runs. Needs a restart of the editor.
- **`/api/auto-queue/control` needs the ops token.** It started, stopped and drained lanes for anything that could reach the editor, a form posted from another page in the operator's own browser included. A script now sends `Authorization: Bearer $WORKER_TOKEN` (or uses `pnpm ops lane {"lane", "action"}`, the same three verbs); the Start, Drain and Stop buttons on /operations call server actions and are unchanged to use. Needs a restart of the editor.
- **A curated tag can exist on some sites only.** A tag's new **Sites** field on /tags (`sites` in `transcripts/tags.json`; `pnpm ops tags` takes it in a define) names the sites it exists on. Its rules then fire, and its pins apply, only to videos on those sites' channels, and every other site drops it from its records, its counts and its `/tags.json` — where **Hidden** only hid the chip. Empty is every site, as before. Setting it, or changing the channels of those sites, re-derives the corpus's tags once at the next index update. The Eva tags are what this is for: they belong on Anilyzer alone.
- **The publish lane.** Publishing can run itself: turn it on at **/operations/publish** (the runner's Start, Drain and Stop, the hold, and the lane's settings; or `publish.enabled` in settings) and the lane checks every `checkEveryMinutes` (10) whether the index is stale; when it is — and its last update is at least `refreshEveryMinutes` (360) old — it updates it, then builds every site whose channels changed or whose data the new index moved, one stage at a time on the `publish` queue. What it may do with a site is the site's own — the **Publish policy** on the site's settings form, `site.json` `publish.auto` —: `off` (the default: left alone), `build`, `preview` (built and deployed to the preview branch `publish.previewBranch`) or `production`; the hub and the homepage have `publish.hub` and `publish.homepage`. A private site is only ever built, and a site needs its Cloudflare Pages project before it may deploy. Hold the lane and the stage running finishes and no next one starts; quiet hours (`publish.quietHours`) do the same; Drain finishes the stage and ends the runner. The lane never forces a stage: a stage that finds its target current does nothing. On /jobs every stage of one run reads `run <id> · <target>`, and a stage still queued when the editor restarts is cancelled, never re-queued — the lane works out again what is stale from what is on disk. `archilyzer publish now` runs the same plan from the command line, one stage after another in its own process.
diff --git a/editor/app/api/media/fetch-window/route.test.ts b/editor/app/api/media/fetch-window/route.test.ts
@@ -97,3 +97,41 @@ test("a cached whole recording answers with the height its persist recorded", as
assert.equal(res.status, 200);
assert.equal("height" in body, false);
});
+
+test("an ask for a window a queued job will already write answers that job, and starts none", async () => {
+ // A fabricated fetch-window job for 295–330 of abc123, held queued behind a
+ // running holder on its queue (release 19, A5).
+ const { getRegistry } = await import("yt-dlp-transcript-common/jobs/registry");
+ const reg = getRegistry();
+ const make = (id: string, params: Record<string, unknown>) => {
+ const record = {
+ id,
+ kind: "fetch-window",
+ queueKey: "clips:youtube",
+ channelSlug: SLUG,
+ status: "queued" as const,
+ queuedAt: Date.now(),
+ logPath: "/dev/null",
+ spec: { kind: "fetch-window", slug: SLUG, params },
+ };
+ reg.register(record);
+ reg.enqueue(record, { start: () => {}, onCancel: () => {} });
+ };
+ make("HOLD", { videoId: "zzz", from: 0, to: 1 });
+ make("WIN5", { videoId: "abc123", from: 295, to: 330 });
+ const before = reg.list().length;
+ try {
+ const res = await post({ ...base, from: 300, to: 310 });
+ assert.equal(res.status, 202);
+ const body = (await res.json()) as Record<string, unknown>;
+ assert.equal(body.jobId, "WIN5");
+ assert.equal(body.existing, true);
+ assert.equal(body.from, 295);
+ assert.equal(body.to, 330);
+ assert.match(String(body.file), /clips\/295\.00-330\.00\.mp4$/);
+ assert.equal(reg.list().length, before, "no second job");
+ } finally {
+ reg.cancel("WIN5");
+ reg.forceRelease("HOLD");
+ }
+});
diff --git a/editor/app/api/media/fetch-window/route.ts b/editor/app/api/media/fetch-window/route.ts
@@ -89,6 +89,8 @@ function answer(outcome: FetchMediaOutcome): NextResponse {
{ status: 200 },
);
}
+ // `existing: true` (release 19, A5): the window was already queued or being
+ // fetched, and `jobId` is THAT job — poll it as you would a new one.
return NextResponse.json(
{
cached: false,
@@ -96,6 +98,7 @@ function answer(outcome: FetchMediaOutcome): NextResponse {
file: outcome.file,
from: outcome.from,
to: outcome.to,
+ ...(outcome.existing ? { existing: true } : {}),
},
{ status: 202 },
);
diff --git a/editor/app/api/ops/fetch-windows/route.test.ts b/editor/app/api/ops/fetch-windows/route.test.ts
@@ -125,10 +125,63 @@ test("a dry run answers the cache, groups by platform queue, and names what it c
[["rumble", 2], ["youtube", 1]],
);
// The window's own URL picks the queue: the mix channel's Rumble video joins
- // the Rumble job rather than starting a second one beside it.
- assert.deepEqual(j.groups.map((g) => g.queueKey).sort(), ["platform:rumble", "platform:youtube"]);
+ // the Rumble job rather than starting a second one beside it. The queue is
+ // the platform's CLIP queue (release 19, A5), not its download queue.
+ assert.deepEqual(j.groups.map((g) => g.queueKey).sort(), ["clips:rumble", "clips:youtube"]);
const yt = j.groups.find((g) => g.platform === "youtube")!;
assert.equal(yt.items[0].webpageUrl, "https://www.youtube.com/watch?v=vid1");
assert.equal(j.unresolved.length, 1);
assert.match(j.unresolved[0].error, /Channel "no-such" not found/);
});
+
+test("a window a queued job will already write is answered in inFlight with that job, and joins no group", async () => {
+ // A fabricated queued fetch-window job for 95–130 of vid1, behind a running
+ // holder on its queue so it never starts: its spec is the request it was
+ // started with, which is what the dedupe reads.
+ const { getRegistry } = await import("yt-dlp-transcript-common/jobs/registry");
+ const reg = getRegistry();
+ const make = (id: string, kind: string, params: Record<string, unknown>) => {
+ const record = {
+ id,
+ kind,
+ queueKey: "clips:youtube",
+ channelSlug: "yt-chan",
+ status: "queued" as const,
+ queuedAt: Date.now(),
+ logPath: "/dev/null",
+ spec: { kind, slug: "yt-chan", params },
+ };
+ reg.register(record);
+ reg.enqueue(record, { start: () => {}, onCancel: () => {} });
+ };
+ make("HOLDER", "fetch-window", { videoId: "other", from: 0, to: 5 });
+ make("WIN95", "fetch-window", { videoId: "vid1", from: 95, to: 130 });
+ try {
+ const r = await post({
+ requestedBy: "test",
+ dryRun: true,
+ items: [item({ from: 100, to: 110 }), item({ from: 200, to: 210 })],
+ });
+ assert.equal(r.status, 200);
+ const j = r.json as {
+ inFlight: { item: { from: number }; jobId: string }[];
+ groups: { items: { from: number }[] }[];
+ };
+ assert.deepEqual(
+ j.inFlight.map((f) => [f.item.from, f.jobId]),
+ [[100, "WIN95"]],
+ );
+ assert.deepEqual(j.groups.flatMap((g) => g.items.map((i) => i.from)), [200]);
+ // A different height cap is a different request.
+ const capped = await post({
+ requestedBy: "test",
+ dryRun: true,
+ maxHeight: 480,
+ items: [item({ from: 100, to: 110 })],
+ });
+ assert.deepEqual((capped.json as { inFlight: unknown[] }).inFlight, []);
+ } finally {
+ reg.cancel("WIN95");
+ reg.forceRelease("HOLDER");
+ }
+});
diff --git a/editor/app/api/ops/fetch-windows/route.ts b/editor/app/api/ops/fetch-windows/route.ts
@@ -29,9 +29,9 @@ export const dynamic = "force-dynamic";
// | { items: [{ slug, id, from, to, clipId?, reason?, pad?, webpageUrl? }, …],
// requestedBy, manifest?, maxHeight?, dryRun? }
// dryRun -> { ok, dryRun: true, groups: [{ platform, queueKey, items }],
-// cached, unresolved[, onDisk, unfetchable] }
+// cached, inFlight, unresolved[, onDisk, unfetchable] }
// else -> { ok, dryRun: false, jobs: [{ platform, queueKey, jobId, items }],
-// jobIds, jobId (one job only), refused, cached, unresolved
+// jobIds, jobId (one job only), refused, cached, inFlight, unresolved
// [, onDisk, unfetchable] }
//
// Fetch clip windows through the managed path, ONE PACED JOB PER PLATFORM
@@ -168,7 +168,16 @@ export async function POST(request: Request) {
return NextResponse.json(plan);
}
const { groups: _groups, ...run } = result;
- const jobIds = run.jobs.map((j) => j.jobId);
+ // EVERY JOB THAT WILL WRITE ONE OF THESE WINDOWS: the ones this call
+ // started, then the ones already queued or running for a window it asked
+ // for (`inFlight`, release 19 A5) — so --wait follows the windows asked
+ // for, not only the jobs that happened to be new.
+ const jobIds = [
+ ...new Set([
+ ...run.jobs.map((j) => j.jobId),
+ ...run.inFlight.map((f) => f.jobId),
+ ]),
+ ];
return NextResponse.json({
...run,
jobIds,
diff --git a/editor/app/channels/[slug]/videos/[id]/videoActions.ts b/editor/app/channels/[slug]/videos/[id]/videoActions.ts
@@ -27,9 +27,12 @@ import {
VTT_FILENAME,
} from "yt-dlp-transcript-common/lib/videoStatus";
import {
+ clipWindowQueueKey,
downloadQueueKey,
resolveQueueKey,
} from "yt-dlp-transcript-common/lib/queueKeys";
+import { windowInFlight } from "yt-dlp-transcript-common/jobs/windowJobs";
+import { getRegistry } from "yt-dlp-transcript-common/jobs/registry";
import { readChannelConfig } from "yt-dlp-transcript-common/controller/channels";
import { setDoNotClean } from "yt-dlp-transcript-common/lib/doNotClean-server";
import {
@@ -991,6 +994,10 @@ export type FetchMediaOutcome =
file: string | null;
from: number;
to: number;
+ // The job was ALREADY queued or running for this window (or a wider one
+ // that covers it) — no second job was started (release 19, A5). `from`
+ // and `to` are then that job's window, which `file` will hold.
+ existing?: true;
}
// `stream: true` was asked for: the caller reads the job's log itself, so the
// StreamActionResult is handed over rather than cancelled.
@@ -1101,6 +1108,31 @@ export async function fetchWindowAction(req: {
};
}
+ // ...THEN THE JOBS ALREADY FETCHING IT (release 19, A5): an ask repeated
+ // while its window is still queued or downloading is answered with that
+ // job, not a second one for the same seconds (jobs/windowJobs.ts). Retry
+ // (`stream: true`) is an explicit re-run and is not deduplicated.
+ if (!req.stream) {
+ const inFlight = windowInFlight(getRegistry().list(), {
+ slug,
+ videoId,
+ from,
+ to,
+ maxHeight: req.maxHeight,
+ });
+ if (inFlight) {
+ return {
+ ok: true,
+ cached: false,
+ jobId: inFlight.jobId,
+ file: clipWindowPath(videoDir, inFlight.from, inFlight.to),
+ from: inFlight.from,
+ to: inFlight.to,
+ existing: true,
+ };
+ }
+ }
+
const url =
req.webpageUrl?.trim() ||
(await findVideoSourceUrl(paths, slug, videoId, r.config));
@@ -1148,7 +1180,10 @@ export async function fetchWindowAction(req: {
// control says and not why it gets used.
const res = await runManagedFunction({
kind: "fetch-window",
- queueKey: videoQueueKey(r.config, req.queueKey),
+ // The platform's CLIP queue, not its download queue: a window must not
+ // wait behind a multi-hour persist or sync (lib/queueKeys.ts,
+ // clipWindowQueueKey). An explicit override still wins.
+ queueKey: resolveQueueKey(clipWindowQueueKey(downloadQueueKey(r.config)), req.queueKey),
paths,
channelSlug: slug,
videoId,
diff --git a/editor/app/channels/[slug]/videos/fetchWindowsAction.ts b/editor/app/channels/[slug]/videos/fetchWindowsAction.ts
@@ -5,7 +5,9 @@ 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 { resolveQueueKey } from "yt-dlp-transcript-common/lib/queueKeys";
+import { clipWindowQueueKey, resolveQueueKey } from "yt-dlp-transcript-common/lib/queueKeys";
+import { getRegistry } from "yt-dlp-transcript-common/jobs/registry";
+import { windowInFlight } from "yt-dlp-transcript-common/jobs/windowJobs";
import { detectPlatform, queueKeyForUrl } from "yt-dlp-transcript-common/lib/platform";
import {
MAX_CLIP_WINDOW_SECONDS,
@@ -45,14 +47,17 @@ import { safeRevalidate } from "../../../lib/safeRevalidate";
// The batch form of fetchWindowAction (videos/[id]/videoActions.ts), with the
// same rules at the door: a window of at most MAX_CLIP_WINDOW_SECONDS, a height
// cap in range, a channel that exists, a URL that resolves. What it adds is the
-// fan-out: the list is grouped by the queue of each window's own URL
-// (`queueKeyForUrl`) and each group starts its
-// own job on that queue, so YouTube and Rumble run side by side, each behind
-// its own platform's other downloads — persistVideosAction's shape.
+// fan-out: the list is grouped by the platform of each window's own URL
+// (`queueKeyForUrl`) and each group starts its own job on that platform's CLIP
+// queue (`clipWindowQueueKey`: `clips:youtube`), so YouTube and Rumble run side
+// by side, and neither waits behind its platform's long downloads (release 19,
+// A5 — it used to be the download queue itself, behind a multi-hour persist).
//
// A PLATFORM COOLING DOWN OR HELD is refused at the door for ITS group (the
// sentence comes back in `refused`); the other groups still start. A window
-// already on disk is answered here (`cached`) and joins no job.
+// already on disk is answered here (`cached`) and joins no job, and so is one a
+// queued or running job will already write (`inFlight`, with that job's id —
+// jobs/windowJobs.ts): the same manifest sent twice queues nothing twice.
//
// Like the single fetch, NOT GATED BY THE DOWNLOAD PAUSE: an operator (or a
// tool they are driving) asked for these seconds by hand.
@@ -92,9 +97,13 @@ export type FetchWindowsActionResult =
// Groups refused at the door, with the refusal's own sentence.
refused: { platform: string; error: string; items: number }[];
cached: FetchWindowsItem[];
+ // Windows a queued or running job will already write, and that job.
+ inFlight: FetchWindowsInFlight[];
unresolved: FetchWindowsUnresolved[];
};
+export type FetchWindowsInFlight = { item: FetchWindowsItem; jobId: string };
+
// The window rules the HTTP door enforces, asked again here: a replayed spec is
// a file on disk.
function itemProblem(item: FetchWindowsItem): string | null {
@@ -121,16 +130,20 @@ function itemProblem(item: FetchWindowsItem): string | null {
async function planFetchWindows(
items: FetchWindowsItem[],
queueOverride: string | undefined,
+ maxHeight: number | undefined,
): Promise<{
groups: FetchWindowsGroup[];
cached: FetchWindowsItem[];
+ inFlight: FetchWindowsInFlight[];
unresolved: FetchWindowsUnresolved[];
}> {
const paths = getPaths();
const configs = new Map<string, ChannelConfig | null>();
const groups = new Map<string, FetchWindowsGroup>();
const cached: FetchWindowsItem[] = [];
+ const inFlight: FetchWindowsInFlight[] = [];
const unresolved: FetchWindowsUnresolved[] = [];
+ const live = getRegistry().list();
for (const item of dedupeFetchWindowsItems(items)) {
const problem = itemProblem(item);
if (problem) {
@@ -150,6 +163,17 @@ async function planFetchWindows(
cached.push(item);
continue;
}
+ const running = windowInFlight(live, {
+ slug: item.slug,
+ videoId: item.id,
+ from: item.from,
+ to: item.to,
+ maxHeight,
+ });
+ if (running) {
+ inFlight.push({ item, jobId: running.jobId });
+ continue;
+ }
const url =
item.webpageUrl?.trim() ||
(await findVideoSourceUrl(paths, item.slug, item.id, config));
@@ -166,7 +190,7 @@ async function planFetchWindows(
// whose URL is no platform's (a curated mix) can hold Rumble videos, and
// its queue would put a second Rumble job beside the first, each pacing
// only itself.
- const queueKey = resolveQueueKey(queueKeyForUrl(url), queueOverride);
+ const queueKey = resolveQueueKey(clipWindowQueueKey(queueKeyForUrl(url)), queueOverride);
const platform = detectPlatform(url) ?? "unknown";
// One job per queue. Two platforms sharing a queue (an override) share a
// job too; the controller keys the cooldown per item, so each is honoured.
@@ -174,7 +198,7 @@ async function planFetchWindows(
group.items.push({ ...item, webpageUrl: url });
groups.set(queueKey, group);
}
- return { groups: [...groups.values()], cached, unresolved };
+ return { groups: [...groups.values()], cached, inFlight, unresolved };
}
// The door each group passes before its job is queued: the platform's hold,
@@ -287,8 +311,12 @@ export async function fetchWindowsAction(
error: `maxHeight ${req.maxHeight} is not a source height to cap a fetch at.`,
};
}
- const { groups, cached, unresolved } = await planFetchWindows(req.items, req.queueKey);
- const base = { groups, cached, unresolved, jobs: [], refused: [] };
+ const { groups, cached, inFlight, unresolved } = await planFetchWindows(
+ req.items,
+ req.queueKey,
+ req.maxHeight,
+ );
+ const base = { groups, cached, inFlight, unresolved, jobs: [], refused: [] };
if (req.dryRun) return { ok: true, dryRun: true, ...base };
if (groups.length === 0) return { ok: true, dryRun: false, ...base };
@@ -350,7 +378,7 @@ export async function fetchWindowsAction(
error: refused.map((r) => `${r.platform}: ${r.error}`).join("; "),
};
}
- return { ok: true, dryRun: false, groups, cached, unresolved, jobs, refused };
+ return { ok: true, dryRun: false, groups, cached, inFlight, unresolved, jobs, refused };
}
export type FetchMissingEvidenceResult =
@@ -384,6 +412,7 @@ export async function fetchMissingEvidenceAction(
jobs: [],
refused: [],
cached: [],
+ inFlight: [],
unresolved: [],
...extra,
};
diff --git a/scripts/archilyzer-ops.mjs b/scripts/archilyzer-ops.mjs
@@ -775,7 +775,11 @@ export function usage() {
' ones no fetch can fill ("unfetchable": deleted, off the site) and starts',
' nothing. A platform cooling down or held is refused for its group; a',
' 429, or two 403s in a row, backs the platform off and stops its job.',
- ' Running the same body again resumes — fetched windows are cached.',
+ ' Running the same body again resumes — fetched windows are cached, and a',
+ ' window a queued or running job will already write is answered in',
+ ' "inFlight" with that job, whose id joins "jobIds" (so --wait follows it)',
+ ' and no second job is queued. Windows run on the platform\'s clip queue',
+ ' (clips:youtube), never behind its long downloads.',
"",
'"preview": "<branch>" on deploy-site or build-deploy makes it a Cloudflare',
" Pages PREVIEW instead of production: the same bundle goes to a branch",