commit f5d2a6e654fa9c7a569642f0c610d52533dc7cdd
parent 1bc6e241fde7b4726d8de5e136b1c4924e110e29
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Fri, 9 Oct 2026 14:52:44 -0400
Merge track A (release 19 A1–A5: ops finds its token, jobs over ops, read side + gated auto-queue control, archival writes, fetch-window queue + dedupe) into r19/integration
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Diffstat:
66 files changed, 4195 insertions(+), 990 deletions(-)
diff --git a/common/bin/archilyzer.ts b/common/bin/archilyzer.ts
@@ -562,6 +562,18 @@ export const COMMANDS: Command[] = [
"--channel <slug> list duplicate and missing transcripts"),
script(["migrate", "channel-priority"], "migrate-channel-priority.ts",
"[--dry-run] the one-shot channel-priority migration (plans/channel-priority.md, S5)"),
+ {
+ path: ["storage", "report"],
+ usage:
+ "[<slug>…] [--json] each channel's media (tierable), text and clip bytes and where its media is — off the reports, no drive touched; totals per place (the corpus disk's is what a move would free)",
+ flags: { json: "boolean" },
+ maxPositionals: 1000,
+ run: async ({ positionals, flags }) =>
+ (await import("./storage-report")).main({
+ slugs: positionals,
+ json: flags.json === true,
+ }),
+ },
script(["storage", "migrate-tier"], "migrate-media-tier.ts",
"<slug>…|--all [--order smallest] [--include-large] [--dry-run] [--reclaim] bring a channel off the retired whole-directory layout onto the media tier, its text home to the corpus disk (editor stopped; --all stops before the three big-text channels)"),
{
diff --git a/common/bin/storage-report.test.ts b/common/bin/storage-report.test.ts
@@ -0,0 +1,93 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import type { ChannelBrief } from "../controller/channels";
+import type { ChannelSnapshot } from "../controller/channelSnapshot";
+import type { StorageLocation } from "../lib/storageLocations";
+import {
+ buildStorageReport,
+ renderStorageReport,
+ storageReportRow,
+} from "./storage-report";
+
+// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test bin/storage-report.test.ts
+//
+// `archilyzer storage report` is a pure fold over the channel briefs: where
+// each channel's media is (corpus, a location, a custom path, legacy), its
+// three tiers' bytes off the report — unknown when the report has no figure,
+// never 0 — and the totals per place.
+
+const LOCATIONS = [
+ { id: "platter", label: "Platter", root: "/mnt/platter" },
+ { id: "archive", label: "Archive", root: "/mnt/platter/archive" },
+] as unknown as StorageLocation[];
+
+function brief(
+ slug: string,
+ config: Record<string, unknown>,
+ snap: Partial<ChannelSnapshot> | null,
+): ChannelBrief {
+ return {
+ slug,
+ config: { handling: "youtube", name: slug, ...config },
+ snapshot: snap ? ({ generatedAt: "2026-10-01T00:00:00.000Z", ...snap } as ChannelSnapshot) : null,
+ } as ChannelBrief;
+}
+
+const briefs = [
+ brief("in-place", {}, { totalMediaBytes: 500, totalTextBytes: 20, totalClipsBytes: 5 }),
+ brief("moved", { mediaDir: "/mnt/platter/archive/moved/media" }, { totalMediaBytes: 900, totalTextBytes: 10 }),
+ brief("custom", { mediaDir: "/srv/elsewhere/custom/media" }, { totalMediaBytes: 100 }),
+ brief("legacy", { dataDir: "/mnt/platter/legacy/data" }, { totalTextBytes: 1 }),
+ brief("old-report", {}, null),
+ brief("posts", { sourceKind: "social", handling: "transcribe", platform: "twitter", postFetcher: "x-gallery-dl" }, {
+ totalMediaBytes: 0,
+ }),
+];
+
+test("a channel's place: corpus, the longest location, a custom path, legacy", () => {
+ const rows = Object.fromEntries(briefs.map((b) => [b.slug, storageReportRow(b, LOCATIONS)]));
+ assert.equal(rows["in-place"].place, "corpus");
+ assert.equal(rows["in-place"].location, null);
+ assert.equal(rows.moved.place, "location");
+ assert.equal(rows.moved.location, "archive");
+ assert.equal(rows.custom.place, "custom");
+ assert.equal(rows.custom.mediaDir, "/srv/elsewhere/custom/media");
+ assert.equal(rows.legacy.place, "legacy");
+ assert.equal(rows.legacy.location, "platter");
+});
+
+test("a figure the report does not have is unknown, never 0", () => {
+ const row = storageReportRow(briefs[4], LOCATIONS);
+ assert.deepEqual(
+ [row.mediaBytes, row.textBytes, row.clipsBytes, row.reportedAt],
+ [null, null, null, null],
+ );
+ assert.equal(storageReportRow(briefs[2], LOCATIONS).textBytes, null);
+});
+
+test("the report leaves out social channels, sorts by media, and totals per place", () => {
+ const report = buildStorageReport(briefs, LOCATIONS);
+ assert.deepEqual(
+ report.rows.map((r) => r.slug),
+ ["moved", "in-place", "custom", "legacy", "old-report"],
+ );
+ assert.deepEqual(report.totals, {
+ "location:archive": { channels: 1, mediaBytes: 900, unknown: 0 },
+ corpus: { channels: 2, mediaBytes: 500, unknown: 1 },
+ custom: { channels: 1, mediaBytes: 100, unknown: 0 },
+ legacy: { channels: 1, mediaBytes: 0, unknown: 1 },
+ });
+ // Narrowed to named channels.
+ assert.deepEqual(
+ buildStorageReport(briefs, LOCATIONS, ["custom"]).rows.map((r) => r.slug),
+ ["custom"],
+ );
+});
+
+test("the table says unknown as a dash and names what a move would free", () => {
+ const text = renderStorageReport(buildStorageReport(briefs, LOCATIONS));
+ assert.match(text, /^channel\s+media \(tierable\)\s+text\s+clips\s+where$/m);
+ assert.match(text, /^old-report\s+—\s+—\s+—\s+corpus disk$/m);
+ assert.match(text, /^moved\s.*location archive$/m);
+ assert.match(text, /on the corpus disk \(a move would free\): 2 channel\(s\), .* media, 1 unknown/);
+});
diff --git a/common/bin/storage-report.ts b/common/bin/storage-report.ts
@@ -0,0 +1,158 @@
+// `archilyzer storage report [<slug>…] [--json]` — WHERE EACH CHANNEL'S BYTES
+// ARE, by tier, off its report (release 19, A4).
+//
+// The tally an operator used to do by hand before freeing the corpus disk: per
+// channel, the MEDIA tier's bytes (what `lib/mediaTier.ts`'s `isTierable`
+// names — the audio and the raw live-chat replay; the bytes a media move
+// carries), the text, and the clip-window cache, and WHERE the media is: on the
+// corpus disk (in place: what a move would free), on a storage location
+// (`config.mediaDir` under its root — a channel is never tagged, it is on L iff
+// that is true), at a path no location names, or `legacy` (the retired
+// whole-`data/` move, held by every guard until `storage migrate-tier`).
+//
+// OFF THE REPORTS, NEVER A WALK — the /storage page's rule. One config and one
+// snapshot per channel; no drive is touched, so an unmounted platter cannot
+// stall it. A channel whose report predates a figure, or whose media drive was
+// not answering when the report was written, has that figure UNKNOWN, never 0
+// (`—` in the table, `null` in --json), and "Refresh report" is the fix.
+// Social channels hold no media and are left out.
+//
+// Needs no editor. Writes nothing.
+
+import { getPaths, type Paths } from "../lib/paths";
+import { getSettings } from "../lib/settings";
+import { formatBytes } from "../lib/format";
+import { isSocialChannel } from "../lib/channelConfig";
+import { locationOfDataDir, type StorageLocation } from "../lib/storageLocations";
+import { listChannelBriefs, type ChannelBrief } from "../controller/channels";
+
+export type StoragePlace = "corpus" | "location" | "custom" | "legacy";
+
+export type StorageReportRow = {
+ slug: string;
+ place: StoragePlace;
+ // The location's id when `place` is "location" (or a legacy channel's, when
+ // its retired dataDir sits under one).
+ location: string | null;
+ // config.mediaDir (or a legacy channel's dataDir): where its media is.
+ mediaDir: string | null;
+ mediaBytes: number | null;
+ textBytes: number | null;
+ clipsBytes: number | null;
+ reportedAt: string | null;
+};
+
+export type StorageReport = {
+ rows: StorageReportRow[];
+ // Per place ("corpus", "location:<id>", "custom", "legacy"): channels, the
+ // media bytes known, and how many channels' media bytes are unknown.
+ totals: Record<string, { channels: number; mediaBytes: number; unknown: number }>;
+};
+
+export function storageReportRow(
+ brief: ChannelBrief,
+ locations: StorageLocation[],
+): StorageReportRow {
+ const { config, snapshot } = brief;
+ const legacyDir = config.dataDir?.trim() || "";
+ const mediaDir = config.mediaDir?.trim() || "";
+ const dir = legacyDir || mediaDir;
+ const loc = dir ? locationOfDataDir(dir, locations) : null;
+ const place: StoragePlace = legacyDir
+ ? "legacy"
+ : !mediaDir
+ ? "corpus"
+ : loc
+ ? "location"
+ : "custom";
+ const n = (v: number | undefined) => (typeof v === "number" ? v : null);
+ return {
+ slug: brief.slug,
+ place,
+ location: loc?.id ?? null,
+ mediaDir: dir || null,
+ mediaBytes: n(snapshot?.totalMediaBytes),
+ textBytes: n(snapshot?.totalTextBytes),
+ clipsBytes: n(snapshot?.totalClipsBytes),
+ reportedAt: snapshot?.generatedAt ?? null,
+ };
+}
+
+export function placeKey(row: StorageReportRow): string {
+ return row.place === "location" ? `location:${row.location}` : row.place;
+}
+
+export function buildStorageReport(
+ briefs: ChannelBrief[],
+ locations: StorageLocation[],
+ only: readonly string[] = [],
+): StorageReport {
+ const wanted = new Set(only);
+ const rows = briefs
+ .filter((b) => !isSocialChannel(b.config))
+ .filter((b) => wanted.size === 0 || wanted.has(b.slug))
+ .map((b) => storageReportRow(b, locations))
+ // Biggest media first; unknowns last, by slug.
+ .sort(
+ (a, b) =>
+ (b.mediaBytes ?? -1) - (a.mediaBytes ?? -1) || a.slug.localeCompare(b.slug),
+ );
+ const totals: StorageReport["totals"] = {};
+ for (const r of rows) {
+ const t = (totals[placeKey(r)] ??= { channels: 0, mediaBytes: 0, unknown: 0 });
+ t.channels++;
+ if (r.mediaBytes === null) t.unknown++;
+ else t.mediaBytes += r.mediaBytes;
+ }
+ return { rows, totals };
+}
+
+const size = (v: number | null) => (v === null ? "—" : formatBytes(v));
+
+export function renderStorageReport(report: StorageReport): string {
+ const head = ["channel", "media (tierable)", "text", "clips", "where"];
+ const body = report.rows.map((r) => [
+ r.slug,
+ size(r.mediaBytes),
+ size(r.textBytes),
+ size(r.clipsBytes),
+ r.place === "location"
+ ? `location ${r.location}`
+ : r.place === "corpus"
+ ? "corpus disk"
+ : `${r.place} ${r.mediaDir ?? ""}`.trim(),
+ ]);
+ const widths = head.map((h, i) => Math.max(h.length, ...body.map((row) => row[i].length)));
+ const line = (cells: string[]) =>
+ cells
+ .map((c, i) => (i >= 1 && i <= 3 ? c.padStart(widths[i]) : c.padEnd(widths[i])))
+ .join(" ")
+ .trimEnd();
+ const totals = Object.entries(report.totals).map(
+ ([place, t]) =>
+ ` ${place === "corpus" ? "on the corpus disk (a move would free)" : place}: ` +
+ `${t.channels} channel(s), ${formatBytes(t.mediaBytes)} media` +
+ (t.unknown ? `, ${t.unknown} unknown (refresh their reports)` : ""),
+ );
+ return [line(head), ...body.map(line), "", "Media by place:", ...totals].join("\n");
+}
+
+export async function main(opts: {
+ slugs: string[];
+ json: boolean;
+ paths?: Paths;
+ out?: (s: string) => void;
+}): Promise<number> {
+ const paths = opts.paths ?? getPaths();
+ const out = opts.out ?? ((s: string) => console.log(s));
+ const briefs = await listChannelBriefs(paths);
+ const known = new Set(briefs.map((b) => b.slug));
+ const missing = opts.slugs.filter((s) => !known.has(s));
+ if (missing.length) {
+ console.error(`storage report: no channel ${missing.join(", ")}`);
+ return 2;
+ }
+ const report = buildStorageReport(briefs, getSettings().storage.locations, opts.slugs);
+ out(opts.json ? JSON.stringify(report, null, 2) : renderStorageReport(report));
+ return 0;
+}
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/listJobs.test.ts b/common/jobs/listJobs.test.ts
@@ -9,6 +9,7 @@ import {
listAllJobs,
getJobEntry,
pruneJobLogs,
+ readLogTail,
} from "./listJobs";
// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test common/jobs/listJobs.test.ts
@@ -176,3 +177,21 @@ test("pruneJobLogs all clears everything", async () => {
assert.equal(await logCount(paths), 0);
});
});
+
+test("readLogTail returns the last N lines, read from the end, never a half line", async () => {
+ const dir = await mkdtemp(path.join(tmpdir(), "log-tail-"));
+ try {
+ const file = path.join(dir, "a.log");
+ await writeFile(file, Array.from({ length: 100 }, (_, i) => `line ${i}`).join("\n") + "\n\n");
+ assert.deepEqual(await readLogTail(file, 3), ["line 97", "line 98", "line 99"]);
+ // A window smaller than the file: the first (partial) line is dropped.
+ const small = await readLogTail(file, 1000, 20);
+ assert.ok(small.length > 0 && small.length < 100);
+ assert.ok(small.every((l) => /^line \d+$/.test(l)));
+ assert.equal(small[small.length - 1], "line 99");
+ assert.deepEqual(await readLogTail(file, 0), []);
+ assert.deepEqual(await readLogTail(path.join(dir, "missing.log"), 5), []);
+ } finally {
+ await rm(dir, { recursive: true, force: true });
+ }
+});
diff --git a/common/jobs/listJobs.ts b/common/jobs/listJobs.ts
@@ -1,5 +1,5 @@
import path from "node:path";
-import { readdir, stat, readFile, rm } from "node:fs/promises";
+import { open, readdir, stat, readFile, rm } from "node:fs/promises";
import type { Paths } from "../lib/paths";
import { mapConcurrent } from "../lib/concurrency";
import { getRegistry, type JobRecord, type JobStatus } from "./registry";
@@ -246,6 +246,38 @@ export async function readLogChunk(
return { content: raw.slice(fromBytes), nextOffset: raw.length };
}
+// THE LAST `lines` LINES OF A LOG, read from its END. A transcription job's log
+// runs to megabytes; an agent asking "what is it doing" wants the last screen
+// of it, not a download of the whole file (`readLogChunk` from 0 is that). At
+// most `maxBytes` are read: a log whose last window holds fewer lines answers
+// with what it has, and the first line of a mid-file window is dropped, so a
+// half line is never handed back as a whole one. Missing file → [].
+export async function readLogTail(
+ logPath: string,
+ lines: number,
+ maxBytes = 256 * 1024,
+): Promise<string[]> {
+ if (lines <= 0) return [];
+ let handle: Awaited<ReturnType<typeof open>> | undefined;
+ try {
+ handle = await open(logPath, "r");
+ const { size } = await handle.stat();
+ const start = Math.max(0, size - maxBytes);
+ const length = size - start;
+ if (length <= 0) return [];
+ const buf = Buffer.alloc(length);
+ await handle.read(buf, 0, length, start);
+ const all = buf.toString("utf8").split("\n");
+ if (start > 0) all.shift();
+ while (all.length && all[all.length - 1] === "") all.pop();
+ return all.slice(-lines);
+ } catch {
+ return [];
+ } finally {
+ await handle?.close().catch(() => {});
+ }
+}
+
// Delete `.log` + `.meta.json` pairs by retention policy. Never touches a job
// the registry currently reports running or queued. `all` clears every finished
// job; otherwise `keepLast` drops the tail past the newest N and `olderThanMs`
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
@@ -4,6 +4,13 @@
- **Reports can be checked, their quotes verified and their video attached from the command line.** `archilyzer reports check <site> [--reports a,b] [--allow-missing-media]` says what the site's compose would say about its reports — each report validated, every cited quote checked against its record, every cited moment's prepared media present and current, the report's video under the publish limit — without a build, and exits 1 with the list; `--reports` checks those reports, drafts not yet in site.json included. `archilyzer reports verify-quotes <report.json> [--json]` runs compose's own quote check over every video, audio and post quote of one report and prints each one's score and the transcript it matched best, and, where the record has an `en-orig` track, that track's score: a quote that matches a served `en` track but not `en-orig` — the words as spoken — is reported, though compose would pass it. `archilyzer reports attach-video <report.json> <video> [--poster <image>] [--caption <line>]` makes the video the report's: an H.264 mp4 under the 24 MiB limit is remuxed, anything else encoded to fit (refused, with its length, when it would be unwatchable), written beside report.json as `video.mp4` with a poster — the one given, the one it had, or a frame of the video — and `video` set in report.json.
- **A file transcription is timed from the engine, and a short one does not queue behind batches.** `pnpm ops transcribe`'s `durationMs` is now the engine's own time — from the moment a worker took the job — and the new `waitedMs` is how long it waited for a free worker; the job's log says both. A file of up to 15 minutes of audio (a window, usually) waits in the worker pool ahead of every parked transcription, a manual batch's next video included, not only ahead of the transcription lane; a longer file keeps its place behind batches queued before it. Nothing running is interrupted. Needs a restart of the editor.
- **Heavy work takes turns, above a memory floor.** A publish stage's `next build` — a site's, the hub's, the homepage's — now waits for the machine's one heavy slot, which every e2e run and any `pnpm heavy -- <cmd>` (a video render) take too, and then until at least 6000 MB is available; the stage's log says whom it waits behind ("waiting for the heavy slot — held by …") or how much memory there is ("waiting for memory — 4210 MB available, the floor is 6000 MB"). Cancel still stops it. `HEAVY_MIN_FREE_MB` moves the floor (`0` turns it off) and `HEAVY=0` skips the gate. Needs a restart of the editor.
+- **`pnpm ops` finds the editor's token itself.** When `WORKER_TOKEN` or `ARCHILYZER_EDITOR_URL` is not set, it reads them from `editor/.env.local` and `editor/.env` of the checkout the script is in — wherever it is run from — and, in a linked worktree, which has no `editor/.env`, from the main worktree's; nothing else is taken from those files, and a variable already set wins. A 401 or 503 now says where the token came from (never what it is). `--wait` asks `GET /api/ops/job/<id>` (behind the token) whether a job it can no longer follow is still there, instead of the UI's live-jobs view; a job the editor has forgotten since (100 later jobs, or a restart) is reported with the status it ended with — it used to print `archived` and exit 1 for a job that finished `done` — and while a job waits, `--wait` prints its place in the queue whenever it changes. `GET /api/ops/job/<id>[?tail=N]` answers one job: its kind, channel, status, times and exit code, from its record or its sidecar; where it waits (`queue: {key, position, queued, head}`); and with `tail` its log's last N lines. Needs a restart of the editor.
+- **The /jobs list and its buttons over `pnpm ops`.** `pnpm ops get jobs` lists jobs newest first — `--active` every queued and running job in queue order with its place, `--failed` the failed ones, `--kind` and `--slug` to narrow either, `--limit` (50, at most 500); a filtered list looks through the newest 2000 jobs. `pnpm ops get job <id> --tail [N]` is one job with the last N lines of its log. `pnpm ops job cancel|drain|promote|force-release|retry <id>…` presses that row's button for each id (`POST /api/ops/job {"verb", "ids"}`) — the same server action, so the same refusals — and answers each id in `results`; an id that could not be acted on (unknown, or a promote of a job already at the front) makes the answer `ok: false`, named, without stopping the others. `job retry-failed` is the page's **Retry all**, and `--wait` follows the jobs a retry started; `job wait <id>…` follows jobs already running, printing each one's place in its queue while it waits. Retry all now returns the new jobs' ids beside its count. Needs a restart of the editor.
+- **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.
- **One index for every site.** The index is updated once and every site, the hub and the homepage are built from it; `archilyzer publish status` says, per site, whether its build is current — "stale: 3 channels changed (a, b, c)" as soon as a download, transcription or digest on one of its channels finishes, before any index runs; "stale: data changed" once the index has run and the site's data moved; "stale: config changed" after its site.json, tags or aliases changed — and whether what is deployed is that build, with a build made by older code marked "code newer" but not stale.
diff --git a/editor/app/api/auto-queue/control/route.ts b/editor/app/api/auto-queue/control/route.ts
@@ -7,14 +7,22 @@ import {
} from "yt-dlp-transcript-common/controller/autoRunner";
import { LANES } from "yt-dlp-transcript-common/lib/autoQueueTypes";
import type { AutoQueueKind } from "yt-dlp-transcript-common/jobs/autoQueueState";
+import { opsAuth } from "../../ops/_lib";
export const dynamic = "force-dynamic";
// Start or stop an auto-queue runner without a server restart. The runner only
// auto-starts at boot (editor/instrumentation.ts) when enabled, so the save flow
-// and this endpoint are how an operator (or an e2e test) brings a runner up after
-// flipping the enable toggle. The editor admin surface is otherwise
-// unauthenticated (trusted self-host), consistent with the rest of the app.
+// and this endpoint are how a script (or an e2e test) brings a runner up after
+// flipping the enable toggle.
+//
+// BEHIND THE OPS TOKEN SINCE RELEASE 19 (A3). It was the one route that starts
+// and stops lanes with no gate at all — reachable by anything that could reach
+// the editor, and by a page in the operator's own browser posting a form at it
+// (the CSRF reasoning in api/test/_guard.ts). The /operations Start, Drain and
+// Stop buttons no longer call it: they are server actions now
+// (operations/actions.ts), which Next checks for origin. `pnpm ops lane
+// {"lane", "action"}` is the same three verbs for a caller with the token.
//
// Body: { kind: <one of LANES>, action: "start" | "stop" | "drain" }.
// "stop" hard-cancels (aborts in-flight); "drain" lets in-flight units finish.
@@ -24,6 +32,8 @@ export const dynamic = "force-dynamic";
// what to do with it is worse than no button. They have one now, and the check
// is against LANES rather than two names so it cannot be forgotten again.
export async function POST(req: Request) {
+ const denied = opsAuth(req);
+ if (denied) return denied;
let body: { kind?: unknown; action?: unknown };
try {
body = (await req.json()) as typeof body;
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/_door.test.ts b/editor/app/api/ops/_door.test.ts
@@ -0,0 +1,144 @@
+import test from "node:test";
+import assert from "node:assert/strict";
+import { readdir, readFile } from "node:fs/promises";
+import path from "node:path";
+import { TOKEN, callGet, callPost, setupOpsCorpus } from "./_testCorpus";
+
+// Run with:
+// pnpm -C editor exec tsx --test "app/api/ops/_door.test.ts"
+//
+// THE DOOR EVERY OPS ROUTE SHARES: the token gate's three answers, unknown
+// keys, and the traversing slug — checked against the routes themselves, in
+// the request/response alone. Moved from e2e/ops-api.spec.ts (release 19,
+// A2b): none of it needs a browser or a job, and none of it writes.
+//
+// The 503 used to need /api/test/worker-token to switch the variable off
+// inside the one test server; here it is this process's own environment,
+// read per call by getWorkerToken().
+
+const corpus = await setupOpsCorpus("title-filter-channel");
+const route = async (name: string) =>
+ (await import(`./${name}/route`)) as {
+ POST?: Parameters<typeof callPost>[0];
+ GET?: Parameters<typeof callGet>[0];
+ };
+test.after(() => corpus.cleanup());
+
+test("the token gate answers 401 for a missing and for a wrong bearer", async () => {
+ const refresh = (await route("refresh-report")).POST!;
+ const channel = (await route("channel/[slug]")).GET!;
+ const cases: Record<string, string>[] = [{}, { authorization: "Bearer wrong" }];
+ for (const headers of cases) {
+ const post = await callPost(refresh, { all: true }, headers);
+ assert.equal(post.status, 401, JSON.stringify(headers));
+ // The READ side is gated by the same token, not merely the write side.
+ const get = await callGet(channel, "http://localhost/api/ops/channel/anything", { slug: "anything" }, headers);
+ assert.equal(get.status, 401, JSON.stringify(headers));
+ }
+});
+
+// UNSET IS OFF, on the ops door and on the worker door alike — one branch,
+// shared, and neither surface may decide for itself that "no token configured"
+// means "let them in".
+test("with no token configured every guarded route answers 503", async () => {
+ const tags = (await route("tags")).GET!;
+ const health = (await import("../worker/health/route")).GET as unknown as Parameters<typeof callGet>[0];
+ const saved = process.env.WORKER_TOKEN;
+ delete process.env.WORKER_TOKEN;
+ try {
+ // The ops door, with the RIGHT token: it is the surface being off that
+ // answers, not the credential being wrong.
+ const off = await callGet(tags);
+ assert.equal(off.status, 503);
+ assert.match(String(off.body.error), /set WORKER_TOKEN to enable/);
+ // And the LAN worker door, which shares the branch.
+ const worker = await callGet(health);
+ assert.equal(worker.status, 503);
+ assert.match(String(worker.body.error), /set WORKER_TOKEN to enable/);
+ // With no token configured a MISSING header is still 503, not 401: there
+ // is nothing to be unauthorized against.
+ assert.equal((await callGet(tags, undefined, {}, {})).status, 503);
+ } finally {
+ process.env.WORKER_TOKEN = saved;
+ }
+ assert.equal((await callGet(tags)).status, 200);
+ assert.equal(saved, TOKEN);
+});
+
+test("an unknown body key is a 400 that names the accepted keys", async () => {
+ // A misspelled key would otherwise get a cheerful { ok: true } and a channel
+ // that did not change.
+ const sync = await callPost((await route("sync")).POST!, { slug: "x", fullSweep: true });
+ assert.equal(sync.status, 400);
+ assert.equal(sync.body.ok, false);
+ assert.match(String(sync.body.error), /unknown key\(s\): fullSweep/);
+ assert.match(String(sync.body.error), /full/);
+ // So is a nested one, on the route whose body carries an object.
+ const patch = await callPost((await route("channel-config")).POST!, {
+ slug: "x",
+ patch: { downloadFilterExcluded: "rerun" },
+ });
+ assert.equal(patch.status, 400);
+ assert.match(String(patch.body.error), /downloadFilterExcluded/);
+});
+
+test("a traversing slug is refused at the door, on every route that takes one", async () => {
+ const channelsDir = path.join(corpus.transcripts, "channels");
+ const before = (await readdir(channelsDir)).sort();
+ assert.deepEqual(before, ["test-filter"]);
+
+ // EVERY SLUG BELOW REACHES A path.join UNDER channelsDir, and the readers
+ // swallow their own errors — so an unchecked traversing segment would fail
+ // SILENTLY (an empty config read as "channel not found") rather than loudly,
+ // and any future writer on that path would land outside the corpus. reqSlug
+ // is one check for all of them; this is the assertion that it is wired to
+ // each.
+ const cases: [string, Record<string, unknown>][] = [
+ ["metadata-scan", { slug: "../../escape" }],
+ ["refresh-metadata", { slug: "../../escape", id: "abc123" }],
+ ["sync", { slug: "../../escape" }],
+ ["download-missing", { slug: "../../escape" }],
+ ["import-video", { slug: "../../escape", url: "https://example.com/v" }],
+ ["retry-bucket", { slug: "../../escape", bucket: "noTranscript" }],
+ ["refresh-report", { slug: "../../escape" }],
+ ["channel-config", { slug: "../../escape", patch: { cookieMode: "always" } }],
+ ["channel-priority", { slugs: ["../../escape"], tier: "paused" }],
+ ["relocate", { slugs: ["../../escape"], root: "/tmp/ops-api-never" }],
+ ["relocate-back", { slugs: ["../../escape"] }],
+ ["fetch-posts", { slug: "../../escape", older: true }],
+ ["capture-posts", { slug: "../../escape", ids: ["1"] }],
+ ["persist-videos", { items: [{ slug: "../../escape", id: "abc123" }] }],
+ ];
+ for (const [action, data] of cases) {
+ const { status, body } = await callPost((await route(action)).POST!, data);
+ assert.equal(status, 400, action);
+ assert.match(String(body.error), /is not a valid channel slug/, action);
+ }
+
+ // The read route takes its slug as a path SEGMENT. A one-segment name
+ // CHANNEL_SLUG_RE still refuses — a leading dot, how a dotfile beside the
+ // channels dir would be named at — reaches the handler and is refused there.
+ const read = await callGet(
+ (await route("channel/[slug]")).GET!,
+ "http://localhost/api/ops/channel/.escape",
+ { slug: ".escape" },
+ );
+ assert.equal(read.status, 400);
+ assert.match(String(read.body.error), /is not a valid channel slug/);
+
+ // NOTHING WAS TOUCHED: the corpus still holds exactly the fixture channel,
+ // and the fixture's own config is byte-identical.
+ assert.deepEqual((await readdir(channelsDir)).sort(), before);
+ assert.deepEqual(
+ JSON.parse(
+ await readFile(path.join(channelsDir, "test-filter", "config.json"), "utf8"),
+ ),
+ {
+ handling: "youtube",
+ name: "Test Title Filter",
+ url: "https://www.youtube.com/@example/videos",
+ downloadFilter: { include: "guest" },
+ },
+ );
+ assert.deepEqual(await corpus.listJobIds(), []);
+});
diff --git a/editor/app/api/ops/_jobs.ts b/editor/app/api/ops/_jobs.ts
@@ -0,0 +1,153 @@
+import { getPaths } from "yt-dlp-transcript-common/lib/paths";
+import {
+ getJobEntry,
+ listAllJobs,
+ readLogTail,
+ type JobListEntry,
+} from "yt-dlp-transcript-common/jobs/listJobs";
+import { getRegistry, type JobRecord } from "yt-dlp-transcript-common/jobs/registry";
+
+// THE JOBS, AS AN AGENT READS THEM — `GET /api/ops/job/<id>` and
+// `GET /api/ops/jobs`.
+//
+// Read-only shaping over the two places a job is known, the same two the /jobs
+// page reads: the live registry (`getRegistry()`) and the `.jobs` directory
+// (`getJobEntry` / `listAllJobs`, which fall back to the `.meta.json`
+// sidecar). No rule lives here, and nothing is constructed: no channel stats,
+// no disk gate, no runner status — which is the difference from the
+// `/api/jobs/active` view, a UI poll that builds the whole live payload (and
+// walks channel stats for progress) on every request.
+
+export type OpsQueueView = {
+ // The scheduler queue the job waits on ("" = runs in parallel, never queued).
+ key: string;
+ // 0 = running at the head, 1.. = queued that many places back.
+ position: number;
+ // How many jobs wait on that queue in all.
+ queued: number;
+ // The job holding the queue's head, when it is not this one.
+ head?: { id: string; kind: string; channelSlug?: string; startedAt?: number };
+};
+
+export type OpsJobView = Omit<JobListEntry, "logPath"> & {
+ // Only while the job is queued or running.
+ queue?: OpsQueueView;
+ // Live jobs only: what the runner last reported.
+ progress?: JobRecord["progress"];
+ tasks?: number;
+ draining?: boolean;
+ background?: boolean;
+};
+
+export function queueOf(id: string): OpsQueueView | undefined {
+ const registry = getRegistry();
+ const record = registry.get(id);
+ if (!record || (record.status !== "queued" && record.status !== "running")) {
+ return undefined;
+ }
+ const position = registry.positionInQueue(id);
+ if (position < 0) return undefined;
+ const snap = registry.listQueues().find((q) => q.name === record.queueKey);
+ const head = snap?.running;
+ return {
+ key: record.queueKey,
+ position,
+ queued: snap?.queued.length ?? 0,
+ ...(head && head.id !== id
+ ? {
+ head: {
+ id: head.id,
+ kind: head.kind,
+ ...(head.channelSlug ? { channelSlug: head.channelSlug } : {}),
+ ...(head.startedAt ? { startedAt: head.startedAt } : {}),
+ },
+ }
+ : {}),
+ };
+}
+
+export function viewOf(entry: JobListEntry): OpsJobView {
+ const { logPath: _logPath, ...rest } = entry;
+ void _logPath;
+ const record = getRegistry().get(entry.id);
+ const queue = queueOf(entry.id);
+ return {
+ ...rest,
+ ...(queue ? { queue } : {}),
+ ...(record?.progress ? { progress: record.progress } : {}),
+ ...(record?.tasks?.length ? { tasks: record.tasks.length } : {}),
+ ...(record?.draining ? { draining: true } : {}),
+ ...(record?.background ? { background: true } : {}),
+ };
+}
+
+// One job, or null when neither the registry nor `.jobs/` knows the id.
+export async function readJob(
+ id: string,
+ tail: number,
+): Promise<{ job: OpsJobView; tail?: string[] } | null> {
+ const entry = await getJobEntry(getPaths(), id);
+ if (!entry) return null;
+ return {
+ job: viewOf(entry),
+ ...(tail > 0 ? { tail: await readLogTail(entry.logPath, tail) } : {}),
+ };
+}
+
+export type JobsFilter = {
+ active?: boolean;
+ failed?: boolean;
+ kind?: string;
+ slug?: string;
+ limit: number;
+};
+
+// How far back a FILTERED list looks. A page of the newest `limit` jobs
+// filtered afterwards would answer "no failed syncs" whenever the last fifty
+// jobs were refresh-reports; reading this many sidecars (32 at a time) is what
+// the /jobs page's own "Load more" costs at its ceiling.
+export const JOBS_SCAN = 2000;
+
+export async function listJobs(
+ filter: JobsFilter,
+): Promise<{ jobs: OpsJobView[]; scanned: number; total: number }> {
+ const match = (e: { status: string; kind?: string; channelSlug?: string }) =>
+ (!filter.failed || e.status === "failed") &&
+ (!filter.kind || e.kind === filter.kind) &&
+ (!filter.slug || e.channelSlug === filter.slug);
+
+ if (filter.active) {
+ // The live head, straight from the registry: queue order, running first.
+ const registry = getRegistry();
+ const live = registry
+ .list()
+ .filter((r) => r.status === "running" || r.status === "queued")
+ .filter(match);
+ const paths = getPaths();
+ const entries = await Promise.all(live.map((r) => getJobEntry(paths, r.id)));
+ const views = entries
+ .filter((e): e is JobListEntry => e !== null)
+ .map(viewOf);
+ views.sort(
+ (a, b) =>
+ (a.queue?.key ?? "").localeCompare(b.queue?.key ?? "") ||
+ (a.queue?.position ?? 0) - (b.queue?.position ?? 0),
+ );
+ return {
+ jobs: views.slice(0, filter.limit),
+ scanned: live.length,
+ total: views.length,
+ };
+ }
+
+ const filtered = Boolean(filter.failed || filter.kind || filter.slug);
+ const page = await listAllJobs(getPaths(), {
+ limit: filtered ? JOBS_SCAN : filter.limit,
+ });
+ const hits = page.entries.filter(match);
+ return {
+ jobs: hits.slice(0, filter.limit).map(viewOf),
+ scanned: page.entries.length,
+ total: page.total,
+ };
+}
diff --git a/editor/app/api/ops/_read.test.ts b/editor/app/api/ops/_read.test.ts
@@ -0,0 +1,146 @@
+import test from "node:test";
+import assert from "node:assert/strict";
+import { writeFile } from "node:fs/promises";
+import path from "node:path";
+import { callGet, callPost, setupOpsCorpus } from "./_testCorpus";
+import { redactSecrets } from "./_read";
+
+// Run with:
+// pnpm -C editor exec tsx --test "app/api/ops/_read.test.ts"
+//
+// The read side (release 19, A3): `get settings|storage|sites|workers|
+// auto-queue|scheduler|cleanup <slug>` against a temp corpus — each behind the
+// token, each the payload its page already builds, no secret in any of them —
+// and the lane control route's new gate.
+
+const corpus = await setupOpsCorpus("one-youtube-channel-with-data");
+await corpus.writeSettings({
+ minFreeDiskGB: 0,
+ workers: [
+ {
+ id: "lan-box",
+ name: "LAN box",
+ kind: "remote",
+ enabled: false,
+ priority: 1,
+ remote: { baseUrl: "http://lan-box.local:3001", token: "never-printed" },
+ },
+ ],
+});
+await corpus.writeSite("demo-site", {
+ siteUrl: "https://demo.example",
+ channels: [{ slug: "test-youtube", groupId: "default" }],
+});
+await writeFile(
+ path.join(corpus.transcripts, "channels", "test-youtube", "snapshot.json"),
+ JSON.stringify({
+ generatedAt: "2026-10-01T00:00:00.000Z",
+ buckets: { downloadedNoTranscript: ["a", "b"], downloadedAutoSubsOnly: [] },
+ cleanupBytes: { transcribedWithAudio: 1000, multipleAudioFormats: 10, foreignAudio: 5 },
+ }),
+);
+
+type Get = Parameters<typeof callGet>[0];
+const route = async (name: string) => (await import(`./${name}/route`)).GET as Get;
+test.after(() => corpus.cleanup());
+
+test("every read route is behind the token", async () => {
+ for (const name of ["settings", "storage", "sites", "workers", "auto-queue", "scheduler"]) {
+ const res = await callGet(await route(name), undefined, {}, {});
+ assert.equal(res.status, 401, name);
+ }
+ const cleanup = await callGet(await route("cleanup/[slug]"), undefined, { slug: "test-youtube" }, {});
+ assert.equal(cleanup.status, 401);
+});
+
+test("get settings is the editor's read, with every secret redacted; ?key= answers one block", async () => {
+ const get = await route("settings");
+ const all = await callGet(get);
+ assert.equal(all.status, 200);
+ const settings = all.body.settings as Record<string, unknown>;
+ assert.equal(settings.minFreeDiskGB, 0);
+ // A default the file never wrote is there: this is getSettings, not the file.
+ assert.ok("autoQueue" in settings);
+ assert.ok(!JSON.stringify(all.body).includes("never-printed"));
+ const workers = settings.workers as { id: string; remote?: { token?: string } }[];
+ assert.equal(workers.find((w) => w.id === "lan-box")?.remote?.token, "<redacted>");
+
+ const one = await callGet(get, "http://localhost/api/ops/settings?key=minFreeDiskGB");
+ assert.deepEqual(one.body, { ok: true, key: "minFreeDiskGB", value: 0 });
+ const nope = await callGet(get, "http://localhost/api/ops/settings?key=minFreeDisk");
+ assert.equal(nope.status, 400);
+ assert.match(String(nope.body.error), /no settings key "minFreeDisk" — known: .*minFreeDiskGB/);
+ const stray = await callGet(get, "http://localhost/api/ops/settings?block=x");
+ assert.equal(stray.status, 400);
+ assert.match(String(stray.body.error), /unknown query key\(s\): block — accepted: key/);
+});
+
+test("redactSecrets replaces secret-named strings at any depth and leaves the rest", () => {
+ assert.deepEqual(
+ redactSecrets({ a: { apiKey: "k", token: "", n: 1 }, list: [{ password: "p", name: "x" }], maxTokens: 5 }),
+ { a: { apiKey: "<redacted>", token: "", n: 1 }, list: [{ password: "<redacted>", name: "x" }], maxTokens: 5 },
+ );
+});
+
+test("get sites lists each site from its site.json", async () => {
+ const res = await callGet(await route("sites"));
+ assert.equal(res.status, 200);
+ assert.deepEqual(res.body.sites, [
+ {
+ siteId: "demo-site",
+ title: "demo-site",
+ siteUrl: "https://demo.example",
+ cloudflareProject: null,
+ audience: "public",
+ listed: true,
+ search: true,
+ publish: null,
+ channels: ["test-youtube"],
+ },
+ ]);
+});
+
+test("get cleanup <slug> is the channel's ledger row, counts for the id lists", async () => {
+ const get = await route("cleanup/[slug]");
+ const res = await callGet(get, undefined, { slug: "test-youtube" });
+ assert.equal(res.status, 200);
+ assert.equal(res.body.slug, "test-youtube");
+ assert.equal(res.body.transcribedBytes, 1000);
+ assert.equal(res.body.extraFormatsBytes, 10);
+ assert.equal(res.body.foreignBytes, 5);
+ assert.equal(res.body.included, true);
+ assert.equal(res.body.measured, false);
+ assert.equal(res.body.downloadedNoTranscript, 2);
+ assert.equal(res.body.reportedAt, "2026-10-01T00:00:00.000Z");
+ assert.equal(res.body.channel, undefined);
+ assert.equal(res.body.snapshot, undefined);
+ const missing = await callGet(get, undefined, { slug: "no-such" });
+ assert.equal(missing.status, 404);
+ const bad = await callGet(get, undefined, { slug: ".escape" });
+ assert.equal(bad.status, 400);
+});
+
+test("get workers, auto-queue, scheduler and storage answer their pages' payloads", async () => {
+ const workers = await callGet(await route("workers"));
+ assert.equal(workers.status, 200, JSON.stringify(workers.body));
+ assert.ok(!JSON.stringify(workers.body).includes("never-printed"));
+ const lanes = await callGet(await route("auto-queue"));
+ assert.equal(lanes.status, 200, JSON.stringify(lanes.body));
+ assert.ok("transcription" in (lanes.body.autoQueue as Record<string, unknown>));
+ const scheduler = await callGet(await route("scheduler"));
+ assert.equal(scheduler.status, 200, JSON.stringify(scheduler.body));
+ const storage = await callGet(await route("storage"));
+ assert.equal(storage.status, 200, JSON.stringify(storage.body));
+});
+
+test("/api/auto-queue/control is behind the ops token now", async () => {
+ const { POST } = await import("../auto-queue/control/route");
+ const anon = await callPost(POST, { kind: "download", action: "stop" }, {});
+ assert.equal(anon.status, 401);
+ const wrong = await callPost(POST, { kind: "download", action: "stop" }, { authorization: "Bearer nope" });
+ assert.equal(wrong.status, 401);
+ // With the token it is the route it always was.
+ const bad = await callPost(POST, { kind: "transcode", action: "stop" });
+ assert.equal(bad.status, 400);
+ assert.match(String(bad.body.error), /kind must be one of/);
+});
diff --git a/editor/app/api/ops/_read.ts b/editor/app/api/ops/_read.ts
@@ -0,0 +1,55 @@
+import { NextResponse } from "next/server";
+import { opsAuth, opsFail } from "./_lib";
+
+// THE READ SIDE'S DOOR (release 19, A3): `GET /api/ops/<noun>` routes that
+// answer what an existing builder already builds for a page — settings, the
+// storage locations, the sites, the workers, the lanes, the sync scheduler, a
+// channel's cleanup row — behind the token. A route here is the gate, the query
+// check and one call; it shapes nothing a page does not already show.
+//
+// UNKNOWN QUERY KEYS ARE A 400, as unknown body keys are on the write side: a
+// misspelled `?key=` would otherwise answer the whole document and look like
+// an answer to the question asked.
+export async function readRoute(
+ request: Request,
+ allowedQuery: readonly string[],
+ run: (q: URLSearchParams) => Promise<NextResponse | Record<string, unknown>>,
+): Promise<NextResponse> {
+ const denied = opsAuth(request);
+ if (denied) return denied;
+ const q = new URL(request.url).searchParams;
+ const unknown = [...q.keys()].filter((k) => !allowedQuery.includes(k));
+ if (unknown.length) {
+ return opsFail(
+ `unknown query key(s): ${unknown.join(", ")} — accepted: ${
+ allowedQuery.length ? allowedQuery.join(", ") : "none"
+ }`,
+ );
+ }
+ try {
+ const out = await run(q);
+ return out instanceof NextResponse ? out : NextResponse.json({ ok: true, ...out });
+ } catch (e) {
+ return opsFail((e as Error).message, 500);
+ }
+}
+
+// A SECRET NEVER LEAVES OVER THIS SURFACE, token or not. A remote worker's
+// outbound `token` lives in settings.json beside everything else; the ops token
+// lets a caller run the editor, not read the credentials of the machines it
+// talks to. Any string under a key that names a secret is replaced by
+// "<redacted>" ("" stays "", so "unset" still reads as unset).
+const SECRET_KEY = /token|secret|password|passwd|api[-_]?key|credential/i;
+
+export function redactSecrets<T>(value: T): T {
+ if (Array.isArray(value)) return value.map(redactSecrets) as T;
+ if (typeof value !== "object" || value === null) return value;
+ const out: Record<string, unknown> = {};
+ for (const [k, v] of Object.entries(value as Record<string, unknown>)) {
+ out[k] =
+ SECRET_KEY.test(k) && typeof v === "string" && v !== ""
+ ? "<redacted>"
+ : redactSecrets(v);
+ }
+ return out as T;
+}
diff --git a/editor/app/api/ops/_testCorpus.ts b/editor/app/api/ops/_testCorpus.ts
@@ -0,0 +1,186 @@
+import { cp, mkdir, mkdtemp, readdir, readFile, rm, writeFile } from "node:fs/promises";
+import os from "node:os";
+import path from "node:path";
+
+// A TEST CORPUS FOR THE OPS ROUTES' UNIT TESTS — imported ONLY by
+// `*.test.ts` files under app/api/ops. Never by a route.
+//
+// The refusals an ops route answers before any job are request/response and
+// nothing else, so they are tested here, in a tsx process, rather than in
+// e2e/ops-api.spec.ts against a running server (release 19, A2b): the same
+// e2e fixtures (`e2e/fixtures/test-transcripts/<name>`), copied into a temp
+// dir; the same default test settings; and the same fake binaries, so a
+// refusal that regresses and starts a job runs a fake, never yt-dlp, next or
+// wrangler.
+//
+// The environment is set by `setupOpsCorpus`, which MUST run before the first
+// import of a route or of anything that calls getPaths() (it caches).
+// Each test file is its own process under `tsx --test`, so each gets its own.
+
+export const TOKEN = "test-token";
+
+const EDITOR = path.resolve(/* turbopackIgnore: true */ import.meta.dirname, "..", "..", "..");
+const FIXTURES = path.join(/* turbopackIgnore: true */ EDITOR, "e2e", "fixtures");
+const BIN = path.join(/* turbopackIgnore: true */ FIXTURES, "bin");
+
+export type OpsCorpus = {
+ root: string;
+ transcripts: string;
+ settingsFile: string;
+ // Settings as e2e's writeSettings writes them: the default test settings
+ // with `patch` laid over, one level deep.
+ writeSettings: (patch?: Record<string, unknown>) => Promise<void>;
+ writeChannelConfig: (slug: string, config: Record<string, unknown>) => Promise<void>;
+ writeSite: (siteId: string, site?: Record<string, unknown>) => Promise<void>;
+ // The job sidecars on disk — "and nothing was queued".
+ listJobIds: () => Promise<string[]>;
+ cleanup: () => Promise<void>;
+};
+
+export async function setupOpsCorpus(fixture: string | null): Promise<OpsCorpus> {
+ const root = await mkdtemp(path.join(os.tmpdir(), "ops-route-"));
+ const transcripts = path.join(root, "test-transcripts");
+ const settingsFile = path.join(root, "test-settings.json");
+ await mkdir(transcripts, { recursive: true });
+ if (fixture) {
+ await cp(path.join(/* turbopackIgnore: true */ FIXTURES, "test-transcripts", fixture), transcripts, {
+ recursive: true,
+ });
+ }
+ Object.assign(process.env, {
+ WORKER_TOKEN: TOKEN,
+ TRANSCRIPTS_DIR: transcripts,
+ EXPORT_PUBLIC_DIR: path.join(transcripts, ".export-public"),
+ SETTINGS_FILE: settingsFile,
+ YTDLP_BIN: path.join(/* turbopackIgnore: true */ BIN, "fake-ytdlp.mjs"),
+ GALLERY_DL_BIN: path.join(/* turbopackIgnore: true */ BIN, "fake-gallery-dl.mjs"),
+ FFMPEG_BIN: path.join(/* turbopackIgnore: true */ BIN, "fake-ffmpeg.mjs"),
+ FFPROBE_BIN: path.join(/* turbopackIgnore: true */ BIN, "fake-ffprobe.mjs"),
+ WHISPER_BIN: path.join(/* turbopackIgnore: true */ BIN, "fake-whisper.mjs"),
+ WHISPER_MODEL: "/dev/null",
+ CHOUGH_BIN: path.join(/* turbopackIgnore: true */ BIN, "fake-chough.mjs"),
+ CHOUGH_MODEL: "/dev/null",
+ PARAKEET_STITCH_BIN: path.join(/* turbopackIgnore: true */ BIN, "fake-parakeet-stitch.mjs"),
+ PARAKEET_CLI: "/dev/null",
+ PARAKEET_MODEL: "/dev/null",
+ DIARIZE_BIN: path.join(/* turbopackIgnore: true */ BIN, "fake-diarize.mjs"),
+ CLAUDE_BIN: path.join(/* turbopackIgnore: true */ BIN, "fake-claude.mjs"),
+ FINDMNT_BIN: path.join(/* turbopackIgnore: true */ BIN, "fake-findmnt.mjs"),
+ UDISKSCTL_BIN: path.join(/* turbopackIgnore: true */ BIN, "fake-udisksctl.mjs"),
+ WRANGLER_BIN: path.join(/* turbopackIgnore: true */ BIN, "fake-wrangler.mjs"),
+ EXPORT_NEXT_BIN: path.join(/* turbopackIgnore: true */ BIN, "fake-next.mjs"),
+ ARCHILYZER_BRANCH: "main",
+ CLOUDFLARE_API_TOKEN: "e2e-fake-token-never-sent",
+ E2E_LIVE_CHECK: "skip",
+ });
+
+ const writeSettings = async (patch: Record<string, unknown> = {}) => {
+ const base = JSON.parse(
+ await readFile(path.join(/* turbopackIgnore: true */ FIXTURES, "test-settings.default.json"), "utf8"),
+ ) as Record<string, unknown>;
+ const merged: Record<string, unknown> = { ...base, ...patch };
+ for (const [k, v] of Object.entries(patch)) {
+ const b = base[k];
+ if (isObject(v) && isObject(b)) merged[k] = { ...b, ...v };
+ }
+ await writeFile(settingsFile, JSON.stringify(merged, null, 2));
+ };
+ await writeSettings({ minFreeDiskGB: 0 });
+
+ return {
+ root,
+ transcripts,
+ settingsFile,
+ writeSettings,
+ async writeChannelConfig(slug, config) {
+ const dir = path.join(transcripts, "channels", slug);
+ await mkdir(dir, { recursive: true });
+ await writeFile(
+ path.join(dir, "config.json"),
+ JSON.stringify({ handling: "youtube", name: slug, ...config }, null, 2),
+ );
+ },
+ async writeSite(siteId, site = {}) {
+ const dir = path.join(transcripts, "sites", siteId);
+ await mkdir(dir, { recursive: true });
+ await writeFile(
+ path.join(dir, "site.json"),
+ JSON.stringify(
+ {
+ siteId,
+ siteTitle: siteId,
+ siteDescription: "",
+ headerTitle: siteId,
+ homeTagline: "",
+ groups: [{ id: "default", name: "All channels", selectedByDefault: true }],
+ defaultGroupId: "default",
+ channels: [],
+ ...site,
+ },
+ null,
+ 2,
+ ),
+ );
+ },
+ async listJobIds() {
+ return (await readdir(path.join(transcripts, ".jobs")).catch(() => []))
+ .filter((f) => f.endsWith(".meta.json"))
+ .sort();
+ },
+ cleanup: () => rm(root, { recursive: true, force: true }),
+ };
+}
+
+function isObject(v: unknown): v is Record<string, unknown> {
+ return typeof v === "object" && v !== null && !Array.isArray(v);
+}
+
+type Handler = (request: Request, ctx?: never) => Promise<Response>;
+export type OpsRes = { status: number; body: Record<string, unknown> & { error?: string } };
+
+// POST a JSON body to a route handler, with the token unless told otherwise.
+export async function callPost(
+ handler: Handler,
+ body: unknown,
+ headers: Record<string, string> = { authorization: `Bearer ${TOKEN}` },
+): Promise<OpsRes> {
+ const res = await handler(
+ new Request("http://localhost/api/ops/x", {
+ method: "POST",
+ headers: { ...headers, "content-type": "application/json" },
+ body: JSON.stringify(body),
+ }),
+ );
+ return { status: res.status, body: (await res.json()) as OpsRes["body"] };
+}
+
+// GET a route handler. `params` for a dynamic segment ([slug], [id]).
+export async function callGet(
+ handler: (request: Request, ctx: { params: Promise<Record<string, string>> }) => Promise<Response>,
+ url = "http://localhost/api/ops/x",
+ params: Record<string, string> = {},
+ headers: Record<string, string> = { authorization: `Bearer ${TOKEN}` },
+): Promise<OpsRes> {
+ const res = await handler(new Request(url, { headers }), {
+ params: Promise.resolve(params),
+ });
+ return { status: res.status, body: (await res.json()) as OpsRes["body"] };
+}
+
+// HOLD THE PUBLISH QUEUE with a fabricated running job, as e2e's
+// holdPublishQueue does: a build or deploy stage a regressed refusal failed to
+// stop can then only QUEUE, never spawn. Writes no sidecar, so listJobIds does
+// not see it.
+export async function holdPublishQueue(): Promise<void> {
+ const { getRegistry, newJobId } = await import("yt-dlp-transcript-common/jobs/registry");
+ const record = {
+ id: newJobId(),
+ kind: "publish-hold",
+ queueKey: "publish",
+ status: "queued" as const,
+ queuedAt: Date.now(),
+ logPath: "/dev/null",
+ };
+ getRegistry().register(record);
+ getRegistry().enqueue(record, { start: () => {}, onCancel: () => {} });
+}
diff --git a/editor/app/api/ops/_writes.test.ts b/editor/app/api/ops/_writes.test.ts
@@ -0,0 +1,168 @@
+import test from "node:test";
+import assert from "node:assert/strict";
+import { mkdtemp, readFile, rm } from "node:fs/promises";
+import os from "node:os";
+import path from "node:path";
+import { callPost, setupOpsCorpus } from "./_testCorpus";
+import { coercedLeaves, settingsPatchProblem } from "./settings/patch";
+
+// Run with:
+// pnpm -C editor exec tsx --test "app/api/ops/_writes.test.ts"
+//
+// The archival writes (release 19, A4), against a temp corpus: the settings
+// patch end to end (its write needs no revalidation), and every other route's
+// refusals — the ones answered before an action revalidates a page, which only
+// a running server can do. Nothing here starts a job.
+
+const corpus = await setupOpsCorpus("one-youtube-channel-with-data");
+type Post = Parameters<typeof callPost>[0];
+const route = async (name: string) => (await import(`./${name}/route`)).POST as Post;
+const SLUG = "test-youtube";
+const VIDEO = "20240101_test1234567";
+const elsewhere = await mkdtemp(path.join(os.tmpdir(), "ops-relocate-root-"));
+test.after(async () => {
+ await corpus.cleanup();
+ await rm(elsewhere, { recursive: true, force: true });
+});
+
+const readSettingsFile = async () =>
+ JSON.parse(await readFile(corpus.settingsFile, "utf8")) as Record<string, unknown>;
+
+test("settings: an unknown key, an owned key and a coerced value are refused, and nothing is written", async () => {
+ const settings = await route("settings");
+ const before = await readSettingsFile();
+
+ const unknown = await callPost(settings, { patch: { minFreeDiskGb: 5 } });
+ assert.equal(unknown.status, 400);
+ assert.match(String(unknown.body.error), /unknown settings key\(s\): minFreeDiskGb/);
+
+ const owned = await callPost(settings, { patch: { channelPriority: {} } });
+ assert.equal(owned.status, 400);
+ assert.match(String(owned.body.error), /"channelPriority" is not patched here: .*pnpm ops channel-priority/);
+
+ // The schema clamps rather than throwing: a value it would not keep as sent
+ // is named with what it would have saved.
+ const clamped = await callPost(settings, { patch: { minFreeDiskGB: -5 } });
+ assert.equal(clamped.status, 400);
+ assert.match(String(clamped.body.error), /minFreeDiskGB: sent -5, would be saved as \d+/);
+
+ const empty = await callPost(settings, { patch: {} });
+ assert.equal(empty.status, 400);
+ const notObject = await callPost(settings, { patch: [] });
+ assert.equal(notObject.status, 400);
+
+ assert.deepEqual(await readSettingsFile(), before);
+});
+
+test("settings: a valid patch is written through saveSettings and answered with the saved value", async () => {
+ const settings = await route("settings");
+ const res = await callPost(settings, {
+ patch: { minFreeDiskGB: 7, syncScheduler: { fullSweepIntervalMinutes: 0 } },
+ });
+ assert.equal(res.status, 200, JSON.stringify(res.body));
+ assert.deepEqual(res.body.changed, ["minFreeDiskGB", "syncScheduler"]);
+ const value = res.body.value as Record<string, unknown>;
+ assert.equal(value.minFreeDiskGB, 7);
+ const file = await readSettingsFile();
+ assert.equal(file.minFreeDiskGB, 7);
+ // One level merged: the block's other keys are kept, not dropped.
+ assert.equal((file.syncScheduler as Record<string, unknown>).fullSweepIntervalMinutes, 0);
+ await corpus.writeSettings({ minFreeDiskGB: 0 });
+});
+
+test("the patch rule: leaves compared, arrays element by element, absent keys not compared", () => {
+ assert.deepEqual(coercedLeaves({ a: 1, b: { c: "x" } }, { a: 1, b: { c: "x", d: 2 } }, "k"), []);
+ assert.deepEqual(coercedLeaves({ b: { c: "x" } }, { b: { c: "y" } }, "k"), [
+ 'k.b.c: sent "x", would be saved as "y"',
+ ]);
+ assert.deepEqual(coercedLeaves([1, 2], [1], "k"), ["k: sent [1,2], would be saved as [1]"]);
+ assert.equal(settingsPatchProblem({ a: 1 }, ["a"], { a: 1 }), null);
+});
+
+test("lane: the publish lane is a lane; an unknown one is refused by the full list", async () => {
+ const lane = await route("lane");
+ const bogus = await callPost(lane, { lane: "transcode", held: true });
+ assert.equal(bogus.status, 400);
+ assert.match(String(bogus.body.error), /transcription, download, digest, backfill, publish/);
+ // Switched off: the publish page's Start refuses, and so does this.
+ const start = await callPost(lane, { lane: "publish", action: "start" });
+ assert.equal(start.status, 400, JSON.stringify(start.body));
+ assert.ok(String(start.body.error).length > 0);
+});
+
+test("clear-platform-hold refuses a name that is not a platform", async () => {
+ const res = await callPost(await route("clear-platform-hold"), { platform: "you tube/../x" });
+ assert.equal(res.status, 400);
+ assert.match(String(res.body.error), /^Not a platform: /);
+ const none = await callPost(await route("clear-platform-hold"), {});
+ assert.equal(none.status, 400);
+});
+
+test("workers: the op and the ids are checked at the door", async () => {
+ const workers = await route("workers");
+ const op = await callPost(workers, { op: "restart", ids: ["a"] });
+ assert.equal(op.status, 400);
+ assert.match(String(op.body.error), /"op" must be one of enable, disable, drain/);
+ const ids = await callPost(workers, { op: "enable", ids: [] });
+ assert.equal(ids.status, 400);
+});
+
+test("transcribe-one, delete-file and do-not-clean refuse before anything is touched", async () => {
+ const one = await route("transcribe-one");
+ const traversing = await callPost(one, { slug: SLUG, id: VIDEO, file: "../audio.mp3" });
+ assert.equal(traversing.status, 400);
+ assert.match(String(traversing.body.error), /is not a file name in the video's directory/);
+ const noChannel = await callPost(one, { slug: "no-such", id: VIDEO });
+ assert.equal(noChannel.status, 400);
+ assert.match(String(noChannel.body.error), /no-such/);
+
+ const del = await route("delete-file");
+ const suspicious = await callPost(del, { slug: SLUG, id: VIDEO, file: "../config.json" });
+ assert.equal(suspicious.status, 400);
+ assert.match(String(suspicious.body.error), /Refusing to delete suspicious filename/);
+ const noFile = await callPost(del, { slug: SLUG, id: VIDEO });
+ assert.equal(noFile.status, 400);
+ assert.match(String(noFile.body.error), /"file" is required/);
+
+ const keep = await route("do-not-clean");
+ const missing = await callPost(keep, { slug: SLUG, id: "not-a-video" });
+ assert.equal(missing.status, 400);
+ assert.equal(missing.body.error, "Video directory not found: not-a-video");
+ const badKeep = await callPost(keep, { slug: SLUG, id: VIDEO, keep: "yes" });
+ assert.equal(badKeep.status, 400);
+ assert.match(String(badKeep.body.error), /"keep" must be a boolean/);
+ assert.deepEqual(await corpus.listJobIds(), []);
+});
+
+test("cleanup refuses an unknown sweep and a traversing slug, before any job", async () => {
+ const cleanup = await route("cleanup");
+ const sweep = await callPost(cleanup, { slug: SLUG, sweep: "everything" });
+ assert.equal(sweep.status, 400);
+ assert.match(
+ String(sweep.body.error),
+ /"sweep" must be one of transcribed, extra-formats, wrong-format, failed-transcriptions/,
+ );
+ const slug = await callPost(cleanup, { slug: "../x", sweep: "transcribed" });
+ assert.equal(slug.status, 400);
+ assert.deepEqual(await corpus.listJobIds(), []);
+});
+
+test("relocate dryRun answers each channel's preview and moves nothing", async () => {
+ const relocate = await route("relocate");
+ const res = await callPost(relocate, { slugs: [SLUG, "no-such"], root: elsewhere, dryRun: true });
+ assert.equal(res.status, 200, JSON.stringify(res.body));
+ assert.equal(res.body.dryRun, true);
+ const [ok, missing] = res.body.previews as { slug: string; preview?: unknown; error?: string }[];
+ assert.equal(ok.slug, SLUG);
+ assert.ok(ok.preview, JSON.stringify(ok));
+ assert.equal(missing.slug, "no-such");
+ assert.ok(missing.error);
+ // No move: the config names no media dir, and no job exists.
+ const config = JSON.parse(
+ await readFile(path.join(corpus.transcripts, "channels", SLUG, "config.json"), "utf8"),
+ ) as { mediaDir?: string };
+ assert.equal(config.mediaDir, undefined);
+ assert.deepEqual(await corpus.listJobIds(), []);
+ const both = await callPost(relocate, { slugs: [SLUG], root: elsewhere, locationId: "x", dryRun: true });
+ assert.equal(both.status, 400);
+});
diff --git a/editor/app/api/ops/auto-queue/route.ts b/editor/app/api/ops/auto-queue/route.ts
@@ -0,0 +1,16 @@
+import { buildAutoQueueStatusPayload } from "../../../operations/status";
+import { readRoute } from "../_read";
+
+export const dynamic = "force-dynamic";
+
+// GET /api/ops/auto-queue
+//
+// The four lanes' status — the /operations payload (`buildAutoQueueStatusPayload`,
+// the same one /api/auto-queue/status polls): per lane its runner (running,
+// draining, stopped), its policy (enabled, held, workers, order), cooldowns,
+// what it picked last and what is pending.
+export async function GET(request: Request) {
+ return readRoute(request, [], async () => ({
+ autoQueue: await buildAutoQueueStatusPayload(),
+ }));
+}
diff --git a/editor/app/api/ops/build-site/route.test.ts b/editor/app/api/ops/build-site/route.test.ts
@@ -0,0 +1,53 @@
+import test from "node:test";
+import assert from "node:assert/strict";
+import { callPost, holdPublishQueue, setupOpsCorpus } from "../_testCorpus";
+
+// Run with:
+// pnpm -C editor exec tsx --test "app/api/ops/build-site/route.test.ts"
+//
+// build-site and build-deploy speak the same body — `siteId` or `siteIds`,
+// never both, never neither — and refuse before any job. Moved from
+// e2e/ops-api.spec.ts (release 19, A2b). The publish queue is held throughout,
+// so a stage a regressed refusal failed to stop could only queue.
+
+const corpus = await setupOpsCorpus("title-filter-channel");
+const buildSite = (await import("./route")).POST;
+const buildDeploy = (await import("../build-deploy/route")).POST;
+await holdPublishQueue();
+test.after(() => corpus.cleanup());
+
+test("build-site and build-deploy each take siteId or siteIds, and refuse both or neither", async () => {
+ for (const [action, handler] of [
+ ["build-site", buildSite],
+ ["build-deploy", buildDeploy],
+ ] as const) {
+ const both = await callPost(handler, { siteId: "a", siteIds: ["a"] });
+ assert.equal(both.status, 400, action);
+ assert.match(String(both.body.error), /not both/, action);
+
+ const neither = await callPost(handler, {});
+ assert.equal(neither.status, 400, action);
+ assert.match(String(neither.body.error), /siteId/, action);
+
+ // `all` is still exclusive of either spelling.
+ const withAll = await callPost(handler, { all: true, siteId: "a" });
+ assert.equal(withAll.status, 400, action);
+ assert.match(String(withAll.body.error), /not both/, action);
+
+ // A MALFORMED ID IS A 400 BEFORE ANY JOB STARTS: an id list of
+ // ["good", "BAD"] used to queue the first build, throw on the second and
+ // answer 500 with no ids at all.
+ const bad = await callPost(handler, { siteIds: ["buildsite", "BAD ID"] });
+ assert.equal(bad.status, 400, action);
+ assert.match(String(bad.body.error), /not a valid site id/, action);
+ assert.equal(bad.body.jobs, undefined, action);
+ }
+ assert.deepEqual(await corpus.listJobIds(), []);
+});
+
+test("build-deploy refuses a preview alongside all rather than building everything", async () => {
+ const { status, body } = await callPost(buildDeploy, { all: true, preview: "tags-exclude" });
+ assert.equal(status, 400);
+ assert.match(String(body.error), /"preview" is not supported with "all"/);
+ assert.deepEqual(await corpus.listJobIds(), []);
+});
diff --git a/editor/app/api/ops/capture-posts/route.test.ts b/editor/app/api/ops/capture-posts/route.test.ts
@@ -71,3 +71,67 @@ test("both halves off is nothing to do unless the articles are asked for by name
// No job was ever written.
assert.deepEqual(await readdir(path.join(ROOT, ".jobs")).catch(() => []), []);
});
+
+// Moved from e2e/ops-api.spec.ts (release 19, A2b): nothing below reaches X —
+// every case is refused before a job exists.
+async function channel(slug: string, config: Record<string, unknown>): Promise<void> {
+ await mkdir(path.join(ROOT, "channels", slug), { recursive: true });
+ await writeFile(path.join(ROOT, "channels", slug, "config.json"), JSON.stringify(config));
+}
+
+test("capture-posts refuses what it cannot capture, and an id not in the archive — before any job", async () => {
+ await channel("test-filter", {
+ handling: "youtube",
+ name: "Test Title Filter",
+ url: "https://www.youtube.com/@example/videos",
+ });
+ await channel("example-bsky", {
+ handling: "transcribe",
+ sourceKind: "social",
+ platform: "bluesky",
+ postFetcher: "bluesky-atproto",
+ socialHandle: "example.bsky.social",
+ name: "Example (Bluesky)",
+ url: "https://bsky.app/profile/example.bsky.social",
+ });
+ // An X channel with an empty posts archive.
+ await channel("example-x", {
+ handling: "transcribe",
+ sourceKind: "social",
+ platform: "twitter",
+ postFetcher: "x-gallery-dl",
+ socialHandle: "example_user",
+ name: "Example (X)",
+ url: "https://x.com/example_user",
+ });
+
+ const video = await post({ slug: "test-filter", ids: ["1"] });
+ assert.equal(video.status, 400);
+ assert.equal(video.error, "test-filter is not a social channel.");
+
+ const bsky = await post({ slug: "example-bsky", ids: ["1"] });
+ assert.equal(bsky.status, 400);
+ assert.match(bsky.error, /cannot capture posts/);
+
+ const neither = await post({ slug: "example-x", ids: ["1"], shots: false, media: false });
+ assert.equal(neither.status, 400);
+ assert.match(neither.error, /both the screenshot and the media are turned off/);
+
+ // The channel's posts archive is empty: every id is a stray, named.
+ const stray = await post({ slug: "example-x", ids: ["111", "222"] });
+ assert.equal(stray.status, 400);
+ assert.equal(stray.error, "2 id(s) not in example-x's posts archive: 111, 222");
+
+ // The body's shape.
+ const noIds = await post({ slug: "example-x" });
+ assert.equal(noIds.status, 400);
+ assert.match(noIds.error, /"ids" is required/);
+ const badFlag = await post({ slug: "example-x", ids: ["1"], shots: "yes" });
+ assert.equal(badFlag.status, 400);
+ assert.match(badFlag.error, /"shots" must be a boolean/);
+ const unknown = await post({ slug: "example-x", ids: ["1"], limit: 5 });
+ assert.equal(unknown.status, 400);
+ assert.match(unknown.error, /unknown key\(s\): limit/);
+
+ assert.deepEqual(await readdir(path.join(ROOT, ".jobs")).catch(() => []), []);
+});
diff --git a/editor/app/api/ops/cleanup/[slug]/route.ts b/editor/app/api/ops/cleanup/[slug]/route.ts
@@ -0,0 +1,52 @@
+import { getPaths } from "yt-dlp-transcript-common/lib/paths";
+import {
+ isValidChannelSlug,
+ readChannelConfig,
+ readChannelSnapshot,
+} from "yt-dlp-transcript-common/controller/channels";
+import { loadCleanupRow } from "../../../../cleanup/lib/loadCleanup";
+import { opsFail } from "../../_lib";
+import { readRoute } from "../../_read";
+
+export const dynamic = "force-dynamic";
+
+// GET /api/ops/cleanup/<slug>
+//
+// One channel's row of the /cleanup ledger (`loadCleanupRow`), off its report:
+// what each sweep would reclaim — `transcribedBytes` (audio of videos that have
+// a transcript: "Clean audio"), `extraFormatsBytes`, `foreignBytes` (audio not
+// in the target format) — which OVERLAP and are never summed; what holds the
+// rest (`heldBytes`/`heldCounts`: no transcript, keep-latest, do-not-clean,
+// awaiting diarization); whether the channel is in the ledger's total
+// (`included`); and the failed-transcriptions list's length. `measured: false`
+// means the report predates the accounting: its zeros are unknowns.
+//
+// The sweeps themselves are `POST /api/ops/cleanup`.
+export async function GET(
+ request: Request,
+ { params }: { params: Promise<{ slug: string }> },
+) {
+ const { slug } = await params;
+ return readRoute(request, [], async () => {
+ if (!isValidChannelSlug(slug)) {
+ return opsFail(`"${slug}" is not a valid channel slug`);
+ }
+ const paths = getPaths();
+ const config = await readChannelConfig(paths, slug);
+ if (!config) return opsFail(`Channel "${slug}" not found`, 404);
+ const snapshot = await readChannelSnapshot(paths, slug);
+ const row = await loadCleanupRow(paths, { slug, config, snapshot });
+ if (!row) return opsFail(`Channel "${slug}" not found`, 404);
+ const { channel: _c, snapshot: _s, transcribeIds, autoSubsIds, ...rest } = row;
+ void _c;
+ void _s;
+ return {
+ slug,
+ reportedAt: snapshot?.generatedAt ?? null,
+ ...rest,
+ // Counts, not the id lists: a bucket can hold thousands of ids.
+ downloadedNoTranscript: transcribeIds.length,
+ downloadedAutoSubsOnly: autoSubsIds.length,
+ };
+ });
+}
diff --git a/editor/app/api/ops/cleanup/route.ts b/editor/app/api/ops/cleanup/route.ts
@@ -0,0 +1,43 @@
+import {
+ cleanAudioAction,
+ cleanExtraAudioFormatsAction,
+ clearFailedTranscriptionsAction,
+ removeWrongFormatAudioAction,
+} from "../../../channels/[slug]/whisperActions";
+import { jobResponse, oneOf, ops, reqSlug } from "../_lib";
+
+export const dynamic = "force-dynamic";
+
+// POST { slug, sweep: "transcribed" | "extra-formats" | "wrong-format" |
+// "failed-transcriptions" }
+//
+// One of a channel's cleanup sweeps — the channel page's Cleanup buttons, as a
+// job on the channel's queue ({ ok, jobId }; --wait follows it):
+// transcribed "Clean audio": delete the audio of every video that
+// has a transcript (cleanAudioAction) — the sweep the
+// /cleanup total counts;
+// extra-formats delete a video's audio in formats other than the
+// channel's target, where the target is on disk;
+// wrong-format delete every finalized audio file not in the target
+// format, including failed-extract orphans;
+// failed-transcriptions clear the channel's failed-transcriptions list, so
+// "Transcribe missing" retries them (deletes no media).
+// Each is the page's own action, so it keeps what the page's keeps: every
+// sweep skips a do-not-clean video; "transcribed" also keeps the channel's
+// keep-latest window and audio still awaiting diarization; every one deletes
+// through the media tier's remover (removeMediaFile), never a bare link.
+// `get cleanup <slug>` says what each would reclaim.
+const SWEEPS = {
+ transcribed: cleanAudioAction,
+ "extra-formats": cleanExtraAudioFormatsAction,
+ "wrong-format": removeWrongFormatAudioAction,
+ "failed-transcriptions": clearFailedTranscriptionsAction,
+} as const;
+
+export async function POST(request: Request) {
+ return ops(request, ["slug", "sweep"], async (body) => {
+ const slug = reqSlug(body, "slug");
+ const sweep = oneOf(body, "sweep", Object.keys(SWEEPS) as (keyof typeof SWEEPS)[]);
+ return jobResponse(await SWEEPS[sweep](slug));
+ });
+}
diff --git a/editor/app/api/ops/clear-platform-hold/route.ts b/editor/app/api/ops/clear-platform-hold/route.ts
@@ -0,0 +1,21 @@
+import { NextResponse } from "next/server";
+import { clearPlatformHoldAction } from "../../../operations/pacingActions";
+import { ops, opsFail, reqString } from "../_lib";
+
+export const dynamic = "force-dynamic";
+
+// POST { platform: "youtube" | "rumble" | … }
+//
+// The lane strip's "Clear hold": a held platform's hold, its backoff and its
+// raised pace all go, and the next failure starts the escalation from the
+// bottom. The action runs it as a one-step job (kind `clear-platform-hold`, on
+// its own queue) and waits for it, so the answer carries the job's sentence:
+// { ok, message } — "youtube had no hold, backoff or raised pace to clear." when
+// there was nothing to clear.
+export async function POST(request: Request) {
+ return ops(request, ["platform"], async (body) => {
+ const res = await clearPlatformHoldAction(reqString(body, "platform"));
+ if (!res.ok) return opsFail(res.error);
+ return NextResponse.json({ ok: true, message: res.message });
+ });
+}
diff --git a/editor/app/api/ops/delete-file/route.ts b/editor/app/api/ops/delete-file/route.ts
@@ -0,0 +1,26 @@
+import { NextResponse } from "next/server";
+import { deleteVideoFileAction } from "../../../channels/[slug]/videos/[id]/videoActions";
+import { ops, opsFail, reqSlug, reqString, reqVideoId } from "../_lib";
+
+export const dynamic = "force-dynamic";
+
+// POST { slug, id, file }
+//
+// The video page's Files list delete: one file of one video's directory. It is
+// deleteVideoFileAction, so it is `removeMediaFile` — a tiered media file goes
+// with its bytes on the media drive, never as a bare link — and it refuses what
+// the button refuses: a name that is not one entry of the directory, a file on
+// a media drive that is not mounted or not answering (nothing is deleted), a
+// channel whose drive is not answering, something that is not a regular file.
+// The channel's report is refreshed after. No undo.
+export async function POST(request: Request) {
+ return ops(request, ["slug", "id", "file"], async (body) => {
+ const res = await deleteVideoFileAction(
+ reqSlug(body, "slug"),
+ reqVideoId(body, "id"),
+ reqString(body, "file"),
+ );
+ if (!res.ok) return opsFail(res.error);
+ return NextResponse.json({ ok: true });
+ });
+}
diff --git a/editor/app/api/ops/deploy-homepage/route.test.ts b/editor/app/api/ops/deploy-homepage/route.test.ts
@@ -0,0 +1,47 @@
+import test from "node:test";
+import assert from "node:assert/strict";
+import { callPost, holdPublishQueue, setupOpsCorpus } from "../_testCorpus";
+
+// Run with:
+// pnpm -C editor exec tsx --test "app/api/ops/deploy-homepage/route.test.ts"
+//
+// The homepage's deploy path refuses BEFORE a job exists, like deploy-hub's.
+// Moved from e2e/ops-api.spec.ts (release 19, A2b). Nothing here can reach
+// wrangler: "built" is the homepage's stamp in the corpus's .export-builds
+// (there is none), the missing-build call asks for a PREVIEW, and the publish
+// queue is held, so a stage a regressed refusal failed to stop only queues.
+
+const corpus = await setupOpsCorpus("empty");
+const deployHomepage = (await import("./route")).POST;
+const buildHomepage = (await import("../build-homepage/route")).POST;
+await holdPublishQueue();
+test.after(() => corpus.cleanup());
+
+test("deploy-homepage refuses a bad preview name and a missing build; build-homepage refuses a preview without a deploy", async () => {
+ // A bad preview name is judged before any build is looked at.
+ const production = await callPost(deployHomepage, { preview: "main" });
+ assert.equal(production.status, 400);
+ assert.equal(production.body.error, '"main" is the production branch; a preview needs another name.');
+ const shape = await callPost(deployHomepage, { preview: "Not_Valid" });
+ assert.equal(shape.status, 400);
+ assert.match(String(shape.body.error), /not a valid preview branch name/);
+ // And on the build-then-deploy, before anything is built.
+ const buildDeploy = await callPost(buildHomepage, { deploy: true, preview: "main" });
+ assert.equal(buildDeploy.status, 400);
+ assert.match(String(buildDeploy.body.error), /is the production branch/);
+
+ // build-homepage: a preview without a deploy is a mistake, not a build.
+ const previewOnly = await callPost(buildHomepage, { preview: "home-check" });
+ assert.equal(previewOnly.status, 400);
+ assert.match(String(previewOnly.body.error), /"preview" needs "deploy": true/);
+
+ // No build: the homepage's stamp is the corpus's.
+ const unbuilt = await callPost(deployHomepage, { preview: "home-check" });
+ assert.equal(unbuilt.status, 400);
+ assert.match(
+ String(unbuilt.body.error),
+ /^no build of _homepage in .*\.export-builds\/_homepage — archilyzer publish homepage$/,
+ );
+
+ assert.deepEqual(await corpus.listJobIds(), []);
+});
diff --git a/editor/app/api/ops/deploy-hub/route.test.ts b/editor/app/api/ops/deploy-hub/route.test.ts
@@ -0,0 +1,56 @@
+import test from "node:test";
+import assert from "node:assert/strict";
+import { mkdir, writeFile } from "node:fs/promises";
+import path from "node:path";
+import { callPost, holdPublishQueue, setupOpsCorpus } from "../_testCorpus";
+
+// Run with:
+// pnpm -C editor exec tsx --test "app/api/ops/deploy-hub/route.test.ts"
+//
+// The hub's deploy path refuses BEFORE a job exists, which is what lets a
+// runbook's `pnpm ops deploy-hub --wait` fail fast instead of queueing a deploy
+// that can only fail. Moved from e2e/ops-api.spec.ts (release 19, A2b), with
+// the publish queue held as it was there.
+
+const corpus = await setupOpsCorpus("title-filter-channel");
+const deployHub = (await import("./route")).POST;
+const buildHub = (await import("../build-hub/route")).POST;
+await holdPublishQueue();
+test.after(() => corpus.cleanup());
+
+test("deploy-hub refuses no project, the homepage's project, and a bundle that is not the hub", async () => {
+ // The fixture has no homepage.json at all: no project.
+ const none = await callPost(deployHub, {});
+ assert.equal(none.status, 400);
+ assert.equal(
+ none.body.error,
+ "The hub has no Cloudflare Pages project configured — set it on /sites under Hub.",
+ );
+
+ // The homepage's project is refused by name: a hub deployed there would
+ // replace the software's own site.
+ const dir = path.join(corpus.transcripts, "sites", "_homepage");
+ await mkdir(dir, { recursive: true });
+ const hubFile = path.join(dir, "homepage.json");
+ await writeFile(hubFile, JSON.stringify({ cloudflareProject: "archilyzer" }));
+ const homepages = await callPost(deployHub, {});
+ assert.equal(homepages.status, 400);
+ assert.match(String(homepages.body.error), /"archilyzer", which is the Archilyzer homepage's/);
+
+ // A real project, but no hub build: the hub's bundle is the corpus's own
+ // .export-builds/_hub, stamped built.json.
+ await writeFile(hubFile, JSON.stringify({ cloudflareProject: "archilyzer-hub" }));
+ const unbuilt = await callPost(deployHub, { preview: "hub-check" });
+ assert.equal(unbuilt.status, 400);
+ assert.match(
+ String(unbuilt.body.error),
+ /^no build of _hub in .*\.export-builds\/_hub — archilyzer publish hub$/,
+ );
+
+ // build-hub: a preview without a deploy is a mistake, not a build.
+ const previewOnly = await callPost(buildHub, { preview: "hub-check" });
+ assert.equal(previewOnly.status, 400);
+ assert.match(String(previewOnly.body.error), /"preview" needs "deploy": true/);
+
+ assert.deepEqual(await corpus.listJobIds(), []);
+});
diff --git a/editor/app/api/ops/deploy-site/route.test.ts b/editor/app/api/ops/deploy-site/route.test.ts
@@ -0,0 +1,102 @@
+import test from "node:test";
+import assert from "node:assert/strict";
+import { callPost, holdPublishQueue, setupOpsCorpus } from "../_testCorpus";
+
+// Run with:
+// pnpm -C editor exec tsx --test "app/api/ops/deploy-site/route.test.ts"
+//
+// The deploy routes' refusals — preview names, a site with no Pages project, a
+// site never built, the body's spelling — all decided BEFORE a job exists,
+// which is the whole of what they are. Moved from e2e/ops-api.spec.ts
+// (release 19, A2b). The publish queue is held, and wrangler and next are the
+// e2e fakes, so a regressed refusal could only queue a stage.
+
+const corpus = await setupOpsCorpus("title-filter-channel");
+await corpus.writeSite("previewsite", { cloudflareProject: "proj" });
+// NO Cloudflare project, on purpose (see the valid-preview test).
+await corpus.writeSite("noproj", {});
+const deploySite = (await import("./route")).POST;
+const buildDeploy = (await import("../build-deploy/route")).POST;
+await holdPublishQueue();
+test.after(() => corpus.cleanup());
+
+const BOTH = [
+ ["deploy-site", deploySite],
+ ["build-deploy", buildDeploy],
+] as const;
+
+test("deploy-site and build-deploy refuse a preview name that is not one, before any job", async () => {
+ for (const [action, handler] of BOTH) {
+ // "main" is the production branch: deploying there is not a preview, it is
+ // the live site, and the refusal says so rather than shipping it.
+ const main = await callPost(handler, { siteId: "previewsite", preview: "main" });
+ assert.equal(main.status, 400, action);
+ assert.match(String(main.body.error), /production branch/, action);
+
+ // Uppercase is refused rather than lowercased FOR the caller: silently
+ // rewriting "Main" into "main" would deploy to production.
+ const upper = await callPost(handler, { siteId: "previewsite", preview: "Main" });
+ assert.equal(upper.status, 400, action);
+ assert.match(String(upper.body.error), /not a valid preview branch name/, action);
+ assert.doesNotMatch(String(upper.body.error), /production branch/, action);
+
+ const spaced = await callPost(handler, { siteId: "previewsite", preview: "bad name" });
+ assert.equal(spaced.status, 400, action);
+ assert.match(String(spaced.body.error), /not a valid preview branch name/, action);
+
+ const tooLong = await callPost(handler, { siteId: "previewsite", preview: "a".repeat(29) });
+ assert.equal(tooLong.status, 400, action);
+ assert.match(String(tooLong.body.error), /at most 28 characters/, action);
+ }
+ assert.deepEqual(await corpus.listJobIds(), []);
+});
+
+test("a valid preview is accepted and reaches the action, on both deploy routes", async () => {
+ // The action's own refusal is the last gate before wrangler, so a 400 saying
+ // THAT — rather than "unknown key: preview" or a branch complaint — proves
+ // the name passed validation and the route got all the way to the action.
+ for (const [action, handler] of BOTH) {
+ const { status, body } = await callPost(handler, { siteId: "noproj", preview: "tags-exclude" });
+ assert.equal(status, 400, action);
+ assert.match(String(body.error), /no Cloudflare Pages project/, action);
+ assert.doesNotMatch(String(body.error), /unknown key/, action);
+ }
+ assert.deepEqual(await corpus.listJobIds(), []);
+});
+
+test("deploy-site refuses a site that was never built, in the deploy stage's words, before any job", async () => {
+ // A deploy ships the site's OWN bundle (export/.export-builds/<id>/out,
+ // stamped built.json) — never "whatever export/out holds".
+ const { status, body } = await callPost(deploySite, { siteId: "previewsite" });
+ assert.equal(status, 400);
+ assert.match(
+ String(body.error),
+ /^no build of previewsite in .*\.export-builds\/previewsite — archilyzer publish build previewsite$/,
+ );
+ // A preview is refused for the same reason and just as early.
+ const preview = await callPost(deploySite, { siteId: "previewsite", preview: "tags-exclude" });
+ assert.equal(preview.status, 400);
+ assert.match(String(preview.body.error), /no build of previewsite/);
+ assert.deepEqual(await corpus.listJobIds(), []);
+});
+
+test("deploy-site takes siteId or siteIds, and has no all", async () => {
+ const both = await callPost(deploySite, { siteId: "a", siteIds: ["a"] });
+ assert.equal(both.status, 400);
+ assert.match(String(both.body.error), /not both/);
+
+ const neither = await callPost(deploySite, {});
+ assert.equal(neither.status, 400);
+ assert.match(String(neither.body.error), /siteId/);
+
+ const bad = await callPost(deploySite, { siteIds: ["deploysite", "BAD ID"] });
+ assert.equal(bad.status, 400);
+ assert.match(String(bad.body.error), /not a valid site id/);
+
+ // Deploy-only has no all-sites form — build-deploy owns that — so `all` is
+ // an unknown key here rather than a second spelling of it.
+ const all = await callPost(deploySite, { all: true });
+ assert.equal(all.status, 400);
+ assert.match(String(all.body.error), /unknown key/);
+ assert.deepEqual(await corpus.listJobIds(), []);
+});
diff --git a/editor/app/api/ops/do-not-clean/route.ts b/editor/app/api/ops/do-not-clean/route.ts
@@ -0,0 +1,23 @@
+import { NextResponse } from "next/server";
+import { toggleDoNotCleanAction } from "../../../channels/[slug]/videos/[id]/videoActions";
+import { ops, opsFail, optBool, reqSlug, reqVideoId } from "../_lib";
+
+export const dynamic = "force-dynamic";
+
+// POST { slug, id, keep?: boolean }
+//
+// One video's "Do not clean" marker (do-not-clean.json), set to the value
+// given — `keep` defaults to true; false removes it. Set, not toggled: a retry
+// of the same body is the same answer. The cleanup sweeps skip a marked
+// video's audio. `keep-videos` is the bulk form, by title/description match.
+export async function POST(request: Request) {
+ return ops(request, ["slug", "id", "keep"], async (body) => {
+ const res = await toggleDoNotCleanAction(
+ reqSlug(body, "slug"),
+ reqVideoId(body, "id"),
+ optBool(body, "keep") ?? true,
+ );
+ if (!res.ok) return opsFail(res.error);
+ return NextResponse.json({ ok: true });
+ });
+}
diff --git a/editor/app/api/ops/fetch-posts/route.test.ts b/editor/app/api/ops/fetch-posts/route.test.ts
@@ -79,3 +79,61 @@ test("pages must be a positive whole number", async () => {
assert.equal(res.status, 400);
assert.match(res.error, /pages/);
});
+
+// Moved from e2e/ops-api.spec.ts (release 19, A2b): every case is refused
+// before a job exists, so it needs no server.
+async function channel(slug: string, config: Record<string, unknown>): Promise<void> {
+ await mkdir(path.join(ROOT, "channels", slug), { recursive: true });
+ await writeFile(path.join(ROOT, "channels", slug, "config.json"), JSON.stringify(config));
+}
+
+test("fetch-posts refuses a channel that is not social, full with older, and an older fetch its fetcher cannot do — before any job", async () => {
+ await channel("test-filter", {
+ handling: "youtube",
+ name: "Test Title Filter",
+ url: "https://www.youtube.com/@example/videos",
+ });
+ // A social channel whose fetcher has no older-posts walk.
+ await channel("example-bsky", {
+ handling: "transcribe",
+ sourceKind: "social",
+ platform: "bluesky",
+ postFetcher: "bluesky-atproto",
+ socialHandle: "example.bsky.social",
+ name: "Example (Bluesky)",
+ url: "https://bsky.app/profile/example.bsky.social",
+ });
+
+ // The sentence the action gives for a video channel.
+ const video = await post({ slug: "test-filter" });
+ assert.equal(video.status, 400);
+ assert.equal(video.error, "test-filter is not a social channel.");
+ const videoOlder = await post({ slug: "test-filter", older: true });
+ assert.equal(videoOlder.status, 400);
+ assert.equal(videoOlder.error, "test-filter is not a social channel.");
+
+ // Two walks at once, on any channel.
+ const both = await post({ slug: "example-bsky", full: true, older: true });
+ assert.equal(both.status, 400);
+ assert.match(both.error, /different walks — run one at a time/);
+ const bothVideo = await post({ slug: "test-filter", full: true, older: true });
+ assert.equal(bothVideo.status, 400);
+ assert.match(bothVideo.error, /different walks/);
+
+ // A fetcher with no older walk, a floor that is not a date, a floor
+ // without older, and a limit that is not a count.
+ const bsky = await post({ slug: "example-bsky", older: true });
+ assert.equal(bsky.status, 400);
+ assert.match(bsky.error, /cannot fetch older posts/);
+ const floor = await post({ slug: "example-bsky", older: true, floor: "2020-13-01" });
+ assert.equal(floor.status, 400);
+ assert.match(floor.error, /is not a date/);
+ const floorAlone = await post({ slug: "example-bsky", floor: "2020-01-01" });
+ assert.equal(floorAlone.status, 400);
+ assert.match(floorAlone.error, /only to an older-posts fetch/);
+ const limit = await post({ slug: "example-bsky", limit: 0 });
+ assert.equal(limit.status, 400);
+ assert.match(limit.error, /"limit" must be a whole number above zero/);
+
+ assert.deepEqual(await readdir(path.join(ROOT, ".jobs")).catch(() => []), []);
+});
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/api/ops/job/[id]/route.ts b/editor/app/api/ops/job/[id]/route.ts
@@ -0,0 +1,52 @@
+import { NextResponse } from "next/server";
+import { opsAuth, opsFail } from "../../_lib";
+import { readJob } from "../../_jobs";
+
+export const dynamic = "force-dynamic";
+
+// GET /api/ops/job/<id>[?tail=<lines>]
+//
+// ONE JOB'S STATE, behind the token: its kind, channel, status, times, exit
+// code and replay flag — from the registry while it is live, from its
+// `.meta.json` sidecar after — and, while it is queued or running, WHERE it
+// waits: `queue: {key, position, queued, head}` (0 = running at the head; the
+// head is the job holding the queue when it is not this one). `?tail=N` adds the
+// log's last N lines (at most 500), read from the end of the file.
+//
+// What `pnpm ops --wait` asks when the log endpoint stops answering, in place of
+// the `/api/jobs/active` view: that is a UI poll — ungated, a rewrite onto
+// /api/view/activeJobs that builds the whole live payload (channel stats, the
+// disk gate, every runner's status) to answer "is job X still there?". This
+// reads one record and one sidecar, and it can say how the job ended after the
+// registry forgot it, which the active list never could.
+//
+// An id neither the registry nor `.jobs/` knows is a 404 `{ ok: false }`.
+const MAX_TAIL = 500;
+
+export async function GET(
+ request: Request,
+ { params }: { params: Promise<{ id: string }> },
+) {
+ const denied = opsAuth(request);
+ if (denied) return denied;
+ const { id } = await params;
+ // The log route's own id rule: the id names a file under .jobs/.
+ if (!/^[A-Za-z0-9_-]+$/.test(id)) return opsFail(`"${id}" is not a job id`);
+ const url = new URL(request.url);
+ const unknown = [...url.searchParams.keys()].filter((k) => k !== "tail");
+ if (unknown.length) {
+ return opsFail(`unknown query key(s): ${unknown.join(", ")} — accepted: tail`);
+ }
+ const rawTail = url.searchParams.get("tail");
+ let tail = 0;
+ if (rawTail !== null) {
+ tail = Number(rawTail);
+ if (!Number.isInteger(tail) || tail <= 0) {
+ return opsFail('"tail" must be a whole number of lines above zero');
+ }
+ tail = Math.min(tail, MAX_TAIL);
+ }
+ const found = await readJob(id, tail);
+ if (!found) return opsFail(`no job "${id}"`, 404);
+ return NextResponse.json({ ok: true, ...found });
+}
diff --git a/editor/app/api/ops/job/route.test.ts b/editor/app/api/ops/job/route.test.ts
@@ -0,0 +1,212 @@
+import test from "node:test";
+import assert from "node:assert/strict";
+import { mkdir, mkdtemp, rm, writeFile } from "node:fs/promises";
+import os from "node:os";
+import path from "node:path";
+
+// Run with:
+// pnpm -C editor exec tsx --test "app/api/ops/job/route.test.ts"
+//
+// The jobs over ops, against a temp `.jobs/` and the process's own registry:
+// one job's record (live and archived), where a queued job waits, the log tail,
+// and the gate. Every answer here is a read — nothing revalidates.
+
+const ROOT = await mkdtemp(path.join(os.tmpdir(), "ops-job-route-"));
+process.env.WORKER_TOKEN = "test-token";
+process.env.TRANSCRIPTS_DIR = ROOT;
+process.env.SETTINGS_FILE = path.join(ROOT, "settings.json");
+const JOBS = path.join(ROOT, ".jobs");
+await mkdir(JOBS, { recursive: true });
+
+const { getRegistry } = await import("yt-dlp-transcript-common/jobs/registry");
+const one = (await import("./[id]/route")).GET;
+const verbs = (await import("./route")).POST;
+const many = (await import("../jobs/route")).GET;
+
+async function list(query = ""): Promise<Res> {
+ const res = await many(
+ new Request(`http://localhost/api/ops/jobs${query}`, {
+ headers: { authorization: "Bearer test-token" },
+ }),
+ );
+ return { status: res.status, json: (await res.json()) as Record<string, unknown> };
+}
+
+async function post(body: unknown): Promise<Res> {
+ const res = await verbs(
+ new Request("http://localhost/api/ops/job", {
+ method: "POST",
+ headers: { authorization: "Bearer test-token", "content-type": "application/json" },
+ body: JSON.stringify(body),
+ }),
+ );
+ return { status: res.status, json: (await res.json()) as Record<string, unknown> };
+}
+test.after(() => rm(ROOT, { recursive: true, force: true }));
+
+type Res = { status: number; json: Record<string, unknown> };
+async function get(id: string, query = "", token = "test-token"): Promise<Res> {
+ const res = await one(
+ new Request(`http://localhost/api/ops/job/${id}${query}`, {
+ headers: token ? { authorization: `Bearer ${token}` } : {},
+ }),
+ { params: Promise.resolve({ id }) },
+ );
+ return { status: res.status, json: (await res.json()) as Record<string, unknown> };
+}
+
+// An archived job: a log and a sidecar, no registry record.
+async function archived(id: string, status: string, lines: number): Promise<void> {
+ await writeFile(
+ path.join(JOBS, `${id}.log`),
+ Array.from({ length: lines }, (_, i) => `line ${i}`).join("\n") + "\n",
+ );
+ await writeFile(
+ path.join(JOBS, `${id}.meta.json`),
+ JSON.stringify({
+ id,
+ kind: "sync",
+ queueKey: "youtube",
+ channelSlug: "demo",
+ status,
+ queuedAt: 1,
+ startedAt: 2,
+ endedAt: 3,
+ exitCode: status === "done" ? 0 : 1,
+ }),
+ );
+}
+
+// A live job on queue `key`, registered and submitted like a real one.
+function live(id: string, key: string): void {
+ const record = {
+ id,
+ kind: "fetch-window",
+ queueKey: key,
+ channelSlug: "demo",
+ status: "queued" as const,
+ queuedAt: Date.now(),
+ logPath: path.join(JOBS, `${id}.log`),
+ };
+ getRegistry().register(record);
+ getRegistry().enqueue(record, { start: () => {}, onCancel: () => {} });
+}
+
+test("the gate: no token is a 401, and the id is checked before anything is read", async () => {
+ assert.equal((await get("abc", "", "")).status, 401);
+ assert.equal((await get("abc", "", "wrong")).status, 401);
+ const bad = await get("..%2Fx");
+ assert.equal(bad.status, 400);
+ assert.match(String(bad.json.error), /is not a job id/);
+});
+
+test("an unknown id is a 404 { ok: false }, which --wait reads as 'gone'", async () => {
+ const res = await get("NOSUCHJOB");
+ assert.equal(res.status, 404);
+ assert.equal(res.json.ok, false);
+ assert.equal(res.json.error, 'no job "NOSUCHJOB"');
+});
+
+test("an archived job answers with its sidecar's status — not 'archived'", async () => {
+ await archived("ARCH1", "done", 3);
+ const res = await get("ARCH1");
+ assert.equal(res.status, 200);
+ const job = res.json.job as Record<string, unknown>;
+ assert.equal(job.status, "done");
+ assert.equal(job.kind, "sync");
+ assert.equal(job.inRegistry, false);
+ assert.equal(job.queue, undefined);
+ // The log's path on the server is not part of the answer.
+ assert.equal(job.logPath, undefined);
+ assert.equal(res.json.tail, undefined);
+});
+
+test("?tail=N adds the log's last N lines, and refuses a bad N or an unknown key", async () => {
+ await archived("ARCH2", "failed", 50);
+ const res = await get("ARCH2", "?tail=3");
+ assert.deepEqual(res.json.tail, ["line 47", "line 48", "line 49"]);
+ assert.equal((await get("ARCH2", "?tail=0")).status, 400);
+ assert.equal((await get("ARCH2", "?tail=x")).status, 400);
+ const stray = await get("ARCH2", "?tial=3");
+ assert.equal(stray.status, 400);
+ assert.match(String(stray.json.error), /unknown query key\(s\): tial/);
+});
+
+test("a queued job says where it waits and who holds the head", async () => {
+ live("HEAD1", "rumble");
+ live("WAIT1", "rumble");
+ live("WAIT2", "rumble");
+ const head = (await get("HEAD1")).json.job as Record<string, unknown>;
+ assert.equal(head.status, "running");
+ assert.deepEqual(head.queue, { key: "rumble", position: 0, queued: 2 });
+ const waiting = (await get("WAIT2")).json.job as Record<string, unknown>;
+ assert.equal(waiting.status, "queued");
+ const queue = waiting.queue as Record<string, unknown>;
+ assert.equal(queue.position, 2);
+ assert.equal(queue.queued, 2);
+ assert.deepEqual(
+ { ...(queue.head as Record<string, unknown>), startedAt: undefined },
+ { id: "HEAD1", kind: "fetch-window", channelSlug: "demo", startedAt: undefined },
+ );
+});
+
+test("get jobs --active is the live head in queue order, each with its place", async () => {
+ const res = await list("?active=1");
+ assert.equal(res.status, 200);
+ const jobs = res.json.jobs as { id: string; queue?: { position: number } }[];
+ const rumble = jobs.filter((j) => ["HEAD1", "WAIT1", "WAIT2"].includes(j.id));
+ assert.deepEqual(
+ rumble.map((j) => [j.id, j.queue?.position]),
+ [
+ ["HEAD1", 0],
+ ["WAIT1", 1],
+ ["WAIT2", 2],
+ ],
+ );
+ // The archived jobs are not live.
+ assert.ok(!jobs.some((j) => j.id.startsWith("ARCH")));
+});
+
+test("get jobs filters the history by status, kind and channel", async () => {
+ const failed = (await list("?failed=1")).json.jobs as { id: string }[];
+ assert.deepEqual(failed.map((j) => j.id), ["ARCH2"]);
+ const syncs = (await list("?kind=sync&slug=demo")).json.jobs as { id: string }[];
+ assert.deepEqual(syncs.map((j) => j.id).sort(), ["ARCH1", "ARCH2"]);
+ const none = (await list("?kind=sync&slug=other")).json;
+ assert.deepEqual(none.jobs, []);
+ assert.equal(typeof none.scanned, "number");
+ const limited = (await list("?limit=1")).json.jobs as unknown[];
+ assert.equal(limited.length, 1);
+});
+
+test("get jobs refuses what matches nothing or is not a filter", async () => {
+ const both = await list("?active=1&failed=1");
+ assert.equal(both.status, 400);
+ assert.match(String(both.json.error), /a failed job is not active/);
+ assert.equal((await list("?slug=../x")).status, 400);
+ assert.equal((await list("?kind=Sync%20All")).status, 400);
+ assert.equal((await list("?limit=-1")).status, 400);
+ const stray = await list("?status=failed");
+ assert.equal(stray.status, 400);
+ assert.match(String(stray.json.error), /unknown query key\(s\): status/);
+});
+
+test("job verbs: an unknown verb, missing ids, ids on retry-failed and a bad id are refused before anything", async () => {
+ const unknown = await post({ verb: "kill", ids: ["WAIT1"] });
+ assert.equal(unknown.status, 400);
+ assert.match(String(unknown.json.error), /"verb" must be one of cancel, drain, promote/);
+ const noIds = await post({ verb: "cancel" });
+ assert.equal(noIds.status, 400);
+ assert.match(String(noIds.json.error), /"ids" is required/);
+ const withIds = await post({ verb: "retry-failed", ids: ["WAIT1"] });
+ assert.equal(withIds.status, 400);
+ assert.match(String(withIds.json.error), /takes no "ids"/);
+ const bad = await post({ verb: "cancel", ids: ["WAIT1", "../x"] });
+ assert.equal(bad.status, 400);
+ assert.equal(bad.json.error, "not a job id: ../x");
+ const stray = await post({ verb: "cancel", ids: ["WAIT1"], force: true });
+ assert.equal(stray.status, 400);
+ assert.match(String(stray.json.error), /unknown key\(s\): force/);
+ // Nothing was cancelled by any of them.
+ assert.equal(getRegistry().get("WAIT1")?.status, "queued");
+});
diff --git a/editor/app/api/ops/job/route.ts b/editor/app/api/ops/job/route.ts
@@ -0,0 +1,113 @@
+import { NextResponse } from "next/server";
+import {
+ cancelJobAction,
+ drainJobAction,
+ forceReleaseJobAction,
+ promoteJobAction,
+ retryAllFailedAction,
+ retryJobAction,
+} from "../../../jobs/actions";
+import { OpsInputError, oneOf, ops, opsFail, reqStringArray } from "../_lib";
+
+export const dynamic = "force-dynamic";
+
+// POST { verb: "cancel" | "drain" | "promote" | "force-release" | "retry",
+// ids: [<jobId>, …] }
+// | { verb: "retry-failed" }
+//
+// The buttons on a /jobs row, over HTTP — each verb is the row's server action
+// and nothing more (`editor/app/jobs/actions.ts`), so a job cancelled here and
+// one cancelled by a click are the same act:
+//
+// cancel a queued job leaves its queue; a running one is stopped.
+// drain a running batch starts no new sub-operation and ends `done`
+// when the ones in flight finish; a queued job is cancelled.
+// promote a QUEUED job jumps to the front of its queue.
+// force-release the wedged-slot escape hatch: frees the job's queue slot
+// whatever its record says, SIGKILLing any child.
+// retry re-runs a job from its replay descriptor, ahead of the
+// queue; the new job's id is returned (`jobIds`).
+// retry-failed every failed job still in the registry that has one —
+// the /jobs page's Retry all. Takes no ids.
+//
+// Every id is acted on in turn and answered in `results` ({id, ok, jobId?,
+// error?}); one that could not be (an unknown id, a job no longer queued for a
+// promote) makes the answer `ok: false` — a 400 naming them — but never stops
+// the others, exactly as clicking each row would not. `jobIds` lists the jobs a
+// retry started, so `--wait` follows them.
+const VERBS = [
+ "cancel",
+ "drain",
+ "promote",
+ "force-release",
+ "retry",
+ "retry-failed",
+] as const;
+type Verb = (typeof VERBS)[number];
+
+type Result = { id: string; ok: boolean; jobId?: string; error?: string };
+
+// What a row action's `false` means, per verb — the action itself only says no.
+const REFUSED: Record<Exclude<Verb, "retry" | "retry-failed">, string> = {
+ cancel: "not queued or running (or unknown)",
+ drain: "not queued or running (or unknown)",
+ promote: "not queued, or already at the front of its queue",
+ "force-release": "",
+};
+
+const DONE: Record<Exclude<Verb, "retry-failed">, string> = {
+ cancel: "cancelled",
+ drain: "drained",
+ promote: "promoted",
+ "force-release": "released",
+ retry: "retried",
+};
+
+async function act(verb: Verb, id: string): Promise<Result> {
+ if (verb === "retry") {
+ const res = await retryJobAction(id);
+ if (!res.ok) return { id, ok: false, error: res.error };
+ void res.stream.cancel();
+ return { id, ok: true, jobId: res.jobId };
+ }
+ const run = {
+ cancel: cancelJobAction,
+ drain: drainJobAction,
+ promote: promoteJobAction,
+ "force-release": forceReleaseJobAction,
+ }[verb as Exclude<Verb, "retry" | "retry-failed">];
+ const { ok } = await run(id);
+ return ok ? { id, ok } : { id, ok, error: REFUSED[verb as keyof typeof REFUSED] };
+}
+
+export async function POST(request: Request) {
+ return ops(request, ["verb", "ids"], async (body) => {
+ const verb = oneOf(body, "verb", VERBS);
+ if (verb === "retry-failed") {
+ if (body.ids !== undefined) {
+ throw new OpsInputError('"retry-failed" takes no "ids" — it retries every failed job');
+ }
+ const { count, jobIds } = await retryAllFailedAction();
+ return NextResponse.json({ ok: true, verb, count, jobIds });
+ }
+ const ids = [...new Set(reqStringArray(body, "ids"))];
+ const bad = ids.filter((id) => !/^[A-Za-z0-9_-]+$/.test(id));
+ if (bad.length) {
+ throw new OpsInputError(`not a job id: ${bad.join(", ")}`);
+ }
+ const results: Result[] = [];
+ for (const id of ids) results.push(await act(verb, id));
+ const jobIds = results.flatMap((r) => (r.jobId ? [r.jobId] : []));
+ const failed = results.filter((r) => !r.ok);
+ if (failed.length) {
+ return opsFail(
+ `${failed.length} of ${ids.length} not ${DONE[verb]}: ${failed
+ .map((r) => `${r.id} (${r.error})`)
+ .join("; ")}`,
+ 400,
+ { verb, results, jobIds },
+ );
+ }
+ return NextResponse.json({ ok: true, verb, results, jobIds });
+ });
+}
diff --git a/editor/app/api/ops/jobs/route.ts b/editor/app/api/ops/jobs/route.ts
@@ -0,0 +1,63 @@
+import { NextResponse } from "next/server";
+import { isValidChannelSlug } from "yt-dlp-transcript-common/controller/channels";
+import { opsAuth, opsFail } from "../_lib";
+import { listJobs } from "../_jobs";
+
+export const dynamic = "force-dynamic";
+
+// GET /api/ops/jobs[?active=1 | ?failed=1][&kind=<kind>][&slug=<slug>][&limit=N]
+//
+// The /jobs list, for a caller with no browser. Newest first, `limit` rows
+// (default 50, at most 500), each row `GET /api/ops/job/<id>`'s answer without
+// the tail:
+//
+// active=1 the live head — every queued and running job, in queue order
+// (running first), each with where it waits. From the registry.
+// failed=1 jobs that ended `failed`.
+// kind, slug narrow either list to one job kind / one channel.
+//
+// Without `active` the list is the `.jobs/` directory merged with the registry,
+// as the page's tail is; a filtered list looks through the newest 2000 jobs
+// (`scanned` says how many it read), so "no failed syncs" is not an artefact of
+// the last fifty being refresh-reports. `active` and `failed` together are a
+// 400: a failed job is not active.
+const KEYS = ["active", "failed", "kind", "slug", "limit"];
+const MAX_LIMIT = 500;
+
+export async function GET(request: Request) {
+ const denied = opsAuth(request);
+ if (denied) return denied;
+ const q = new URL(request.url).searchParams;
+ const unknown = [...q.keys()].filter((k) => !KEYS.includes(k));
+ if (unknown.length) {
+ return opsFail(`unknown query key(s): ${unknown.join(", ")} — accepted: ${KEYS.join(", ")}`);
+ }
+ const flag = (k: string) => {
+ const v = q.get(k);
+ return v !== null && v !== "0" && v !== "false";
+ };
+ const active = flag("active");
+ const failed = flag("failed");
+ if (active && failed) {
+ return opsFail('"active" and "failed" together match nothing — a failed job is not active');
+ }
+ const kind = q.get("kind") ?? undefined;
+ if (kind !== undefined && !/^[a-z0-9-]+$/.test(kind)) {
+ return opsFail(`"${kind}" is not a job kind`);
+ }
+ const slug = q.get("slug") ?? undefined;
+ if (slug !== undefined && !isValidChannelSlug(slug)) {
+ return opsFail(`"${slug}" is not a valid channel slug`);
+ }
+ let limit = 50;
+ const rawLimit = q.get("limit");
+ if (rawLimit !== null) {
+ limit = Number(rawLimit);
+ if (!Number.isInteger(limit) || limit <= 0) {
+ return opsFail('"limit" must be a whole number above zero');
+ }
+ limit = Math.min(limit, MAX_LIMIT);
+ }
+ const page = await listJobs({ active, failed, kind, slug, limit });
+ return NextResponse.json({ ok: true, ...page });
+}
diff --git a/editor/app/api/ops/lane/route.ts b/editor/app/api/ops/lane/route.ts
@@ -1,19 +1,29 @@
import { NextResponse } from "next/server";
-import { LANES, type AutoQueueKind } from "yt-dlp-transcript-common/lib/autoQueueTypes";
+import {
+ LANES,
+ PIPELINE_LANES,
+ type AutoQueueKind,
+} from "yt-dlp-transcript-common/lib/autoQueueTypes";
import { getSettings } from "yt-dlp-transcript-common/lib/settings";
-import { drainAutoRunner } from "yt-dlp-transcript-common/controller/autoRunner";
import {
+ drainAutoQueueAction,
pauseLaneAction,
resumeLaneAction,
saveAutoQueueAction,
startAutoQueueAction,
stopAutoQueueAction,
} from "../../../operations/actions";
+import { savePublishSettingsAction } from "../../../operations/settingsActions";
+import {
+ drainPublishLaneAction,
+ startPublishLaneAction,
+ stopPublishLaneAction,
+} from "../../../sites/lib/publishActions";
import { OpsInputError, okResponse, ops, oneOf, optBool } from "../_lib";
export const dynamic = "force-dynamic";
-// POST { lane: transcription|download|digest|backfill,
+// POST { lane: transcription|download|digest|backfill|publish,
// held?: boolean, enabled?: boolean, action?: "start"|"stop"|"drain" }
//
// One lane, up to three independent changes, applied in the order the operator
@@ -26,9 +36,35 @@ export const dynamic = "force-dynamic";
//
// `action` mirrors /api/auto-queue/control's body — same three verbs, same
// meanings — so a caller that already drives that route needs nothing new.
+//
+// THE PUBLISH LANE (release 19, A4) is a pipeline lane, not a LANES entry: its
+// switch is `settings.publish.enabled` (the /operations/publish form's, saved
+// through that form's own action with every other field as stored), its gate
+// `publish.held`, and Start, Drain and Stop are that page's three buttons.
+// A Start the lane refuses (switched off, or quiet hours) is a 400 with the
+// page's sentence.
+const ALL_LANES = [...LANES, ...PIPELINE_LANES] as const;
+
+// The /operations/publish form, filled from the stored block with only the
+// switch changed — so the form's own validation runs on what is stored.
+function publishForm(enabled: boolean): FormData {
+ const p = getSettings().publish;
+ const fd = new FormData();
+ if (enabled) fd.set("publishEnabled", "on");
+ fd.set("publishCheckEveryMinutes", String(p.checkEveryMinutes));
+ fd.set("publishRefreshEveryMinutes", String(p.refreshEveryMinutes));
+ fd.set("publishQuietStart", p.quietHours ? String(p.quietHours.start) : "");
+ fd.set("publishQuietEnd", p.quietHours ? String(p.quietHours.end) : "");
+ fd.set("publishRunner", p.runner);
+ fd.set("publishPreviewBranch", p.previewBranch);
+ fd.set("publishHub", p.hub);
+ fd.set("publishHomepage", p.homepage);
+ return fd;
+}
+
export async function POST(request: Request) {
return ops(request, ["lane", "held", "enabled", "action"], async (body) => {
- const lane = oneOf(body, "lane", LANES) as AutoQueueKind;
+ const lane = oneOf(body, "lane", ALL_LANES);
const enabled = optBool(body, "enabled");
const held = optBool(body, "held");
const action = body.action;
@@ -46,19 +82,24 @@ export async function POST(request: Request) {
);
}
if (enabled !== undefined) {
- // EVERY OTHER POLICY FIELD IS CARRIED FROM THE STORED POLICY, including
- // `root`: saveAutoQueueAction refuses a tree edit while the priority
- // model compiles the roots, and handing it back the stored tree is what
- // makes this a pure enable/disable rather than a tree write.
- const policy = getSettings().autoQueue[lane];
- const saved = await saveAutoQueueAction(lane, {
- enabled,
- maxWorkers: policy.maxWorkers,
- replaceAutoSubs: policy.replaceAutoSubs === true,
- order: policy.order ?? "listed",
- root: policy.root,
- });
- if (!saved.ok) return okResponse(saved);
+ if (lane === "publish") {
+ const saved = await savePublishSettingsAction(undefined, publishForm(enabled));
+ if (!saved.ok) return okResponse(saved);
+ } else {
+ // EVERY OTHER POLICY FIELD IS CARRIED FROM THE STORED POLICY, including
+ // `root`: saveAutoQueueAction refuses a tree edit while the priority
+ // model compiles the roots, and handing it back the stored tree is what
+ // makes this a pure enable/disable rather than a tree write.
+ const policy = getSettings().autoQueue[lane];
+ const saved = await saveAutoQueueAction(lane, {
+ enabled,
+ maxWorkers: policy.maxWorkers,
+ replaceAutoSubs: policy.replaceAutoSubs === true,
+ order: policy.order ?? "listed",
+ root: policy.root,
+ });
+ if (!saved.ok) return okResponse(saved);
+ }
}
if (held !== undefined) {
const gated = held
@@ -66,14 +107,17 @@ export async function POST(request: Request) {
: await resumeLaneAction(lane);
if (!gated.ok) return okResponse(gated);
}
- if (action === "start") {
- const started = await startAutoQueueAction(lane);
- if (!started.ok) return okResponse(started);
- } else if (action === "stop") {
- const stopped = await stopAutoQueueAction(lane);
- if (!stopped.ok) return okResponse(stopped);
- } else if (action === "drain") {
- drainAutoRunner(lane);
+ if (action !== undefined) {
+ const run =
+ lane === "publish"
+ ? { start: startPublishLaneAction, stop: stopPublishLaneAction, drain: drainPublishLaneAction }[action]
+ : {
+ start: () => startAutoQueueAction(lane as AutoQueueKind),
+ stop: () => stopAutoQueueAction(lane as AutoQueueKind),
+ drain: () => drainAutoQueueAction(lane as AutoQueueKind),
+ }[action];
+ const done = await run();
+ if (!done.ok) return okResponse(done);
}
return NextResponse.json({ ok: true, lane });
});
diff --git a/editor/app/api/ops/persist-videos/route.test.ts b/editor/app/api/ops/persist-videos/route.test.ts
@@ -0,0 +1,73 @@
+import test from "node:test";
+import assert from "node:assert/strict";
+import { callPost, setupOpsCorpus } from "../_testCorpus";
+
+// Run with:
+// pnpm -C editor exec tsx --test "app/api/ops/persist-videos/route.test.ts"
+//
+// persist-videos' dry run, its body's shape and its disk floor — all answered
+// before any job. Moved from e2e/ops-api.spec.ts (release 19, A2b).
+
+const corpus = await setupOpsCorpus("one-youtube-channel-with-data");
+const { POST } = await import("./route");
+test.after(() => corpus.cleanup());
+
+const slug = "test-youtube";
+const video = "20240101_test1234567";
+const items = [
+ { slug, id: video },
+ { slug, id: "missing12345" },
+ { slug: "no-such-channel", id: "abc123" },
+];
+
+type Bucket = { count: number; items: { slug: string; id: string }[] };
+type Plan = Record<"saved" | "wrongHeight" | "toFetch" | "noUrl" | "unknown", Bucket> & {
+ willFetch: number;
+};
+
+test("persist-videos buckets a list in a dry run, and refuses a bad body and a low disk — before any job", async () => {
+ const dry = await callPost(POST, { items, format: "video_720", dryRun: true });
+ assert.equal(dry.status, 200, JSON.stringify(dry.body));
+ assert.equal(dry.body.dryRun, true);
+ const plan = dry.body.plan as Plan;
+ assert.deepEqual(plan.toFetch.items, [{ slug, id: video }]);
+ assert.equal(plan.unknown.count, 2);
+ assert.equal(plan.saved.count, 0);
+ assert.equal(plan.willFetch, 1);
+
+ // Nothing to fetch is an answer, not a job.
+ const none = await callPost(POST, { items: [{ slug: "no-such-channel", id: "abc123" }] });
+ assert.equal(none.status, 200, JSON.stringify(none.body));
+ assert.deepEqual(none.body.jobIds, []);
+
+ // The body's shape.
+ const cases: [Record<string, unknown>, RegExp][] = [
+ [{}, /"items" is required/],
+ [{ items: [] }, /"items" is required/],
+ [{ items: [{ slug, id: video, height: 720 }] }, /unknown key\(s\): height/],
+ [{ items: [{ slug, id: "../escape" }] }, /is not a video id/],
+ [{ items, format: "1080p" }, /"format" must be one of original, video_720/],
+ [{ items, replace: "always" }, /"replace" must be one of never, above-height/],
+ [{ items, gapMs: -1 }, /"gapMs" must be a whole number, zero or above/],
+ [{ items, minFreeMemMb: "4096" }, /"minFreeMemMb" must be a whole number/],
+ [{ items, dryRun: "yes" }, /"dryRun" must be a boolean/],
+ [{ items, ids: ["x"] }, /unknown key\(s\): ids/],
+ ];
+ for (const [data, error] of cases) {
+ const res = await callPost(POST, data);
+ assert.equal(res.status, 400, JSON.stringify(data));
+ assert.match(String(res.body.error), error, JSON.stringify(data));
+ }
+
+ // A real run asks the disk floor before it queues anything.
+ await corpus.writeSettings({ minFreeDiskGB: 1_000_000 });
+ try {
+ const low = await callPost(POST, { items });
+ assert.equal(low.status, 400, JSON.stringify(low.body));
+ assert.match(String(low.body.error), /^Low disk space/);
+ } finally {
+ await corpus.writeSettings({ minFreeDiskGB: 0 });
+ }
+
+ assert.deepEqual(await corpus.listJobIds(), []);
+});
diff --git a/editor/app/api/ops/publish/route.test.ts b/editor/app/api/ops/publish/route.test.ts
@@ -0,0 +1,68 @@
+import test from "node:test";
+import assert from "node:assert/strict";
+import { callGet, callPost, holdPublishQueue, setupOpsCorpus } from "../_testCorpus";
+
+// Run with:
+// pnpm -C editor exec tsx --test "app/api/ops/publish/route.test.ts"
+//
+// The publish route's refusals before any job, and its GET — the status.
+// Moved from e2e/ops-api.spec.ts (release 19, A2b); the runs each verb
+// enqueues are still asserted there, against the running server's queue.
+
+const corpus = await setupOpsCorpus("title-filter-channel");
+await corpus.writeSite("pubsite", { cloudflareProject: "pubproj" });
+const { POST, GET } = await import("./route");
+await holdPublishQueue();
+test.after(() => corpus.cleanup());
+
+test("publish: refusals before any job — a never-built deploy, a preview that is not one, a key the verb does not take", async () => {
+ const unbuilt = await callPost(POST, { verb: "deploy", siteId: "pubsite", preview: "r18" });
+ assert.equal(unbuilt.status, 400);
+ assert.match(String(unbuilt.body.error), /^no build of pubsite in .* — archilyzer publish build pubsite$/);
+
+ const main = await callPost(POST, { verb: "deploy", siteId: "pubsite", preview: "main" });
+ assert.equal(main.status, 400);
+ assert.match(String(main.body.error), /production branch/);
+
+ const verb = await callPost(POST, { verb: "launch" });
+ assert.equal(verb.status, 400);
+ assert.match(
+ String(verb.body.error),
+ /"verb" must be one of index, build, deploy, hub, homepage, now, stale/,
+ );
+
+ const extra = await callPost(POST, { verb: "index", siteId: "pubsite" });
+ assert.equal(extra.status, 400);
+ assert.match(String(extra.body.error), /verb "index" takes no other keys/);
+
+ const hubPreview = await callPost(POST, { verb: "hub", preview: "r18" });
+ assert.equal(hubPreview.status, 400);
+ assert.match(String(hubPreview.body.error), /need "deploy": true/);
+
+ assert.deepEqual(await corpus.listJobIds(), []);
+});
+
+test("GET publish is the publish status: the index, the lane, a row per target", async () => {
+ const res = await callGet(GET as unknown as Parameters<typeof callGet>[0]);
+ assert.equal(res.status, 200);
+ const status = res.body as unknown as {
+ ok: boolean;
+ index: { chip: { text: string } };
+ lane: { held: boolean; enabled: boolean };
+ sites: { target: string; chips: { built: { text: string } } }[];
+ hub: { target: string };
+ homepage: { target: string };
+ plan: { steps: { kind: string }[] };
+ };
+ assert.equal(status.ok, true);
+ assert.equal(status.index.chip.text, "no index yet");
+ assert.equal(status.lane.held, false);
+ assert.equal(status.lane.enabled, false);
+ assert.deepEqual(status.sites.map((x) => x.target), ["pubsite"]);
+ assert.equal(status.sites[0].chips.built.text, "update the index first");
+ assert.deepEqual([status.hub.target, status.homepage.target], ["_hub", "_homepage"]);
+ assert.deepEqual(status.plan.steps.map((x) => x.kind), ["update-index"]);
+
+ const anon = await callGet(GET as unknown as Parameters<typeof callGet>[0], undefined, {}, {});
+ assert.equal(anon.status, 401);
+});
diff --git a/editor/app/api/ops/relocate/route.ts b/editor/app/api/ops/relocate/route.ts
@@ -1,7 +1,12 @@
+import { NextResponse } from "next/server";
+import { getPaths } from "yt-dlp-transcript-common/lib/paths";
+import { readChannelConfig } from "yt-dlp-transcript-common/controller/channels";
import { bulkRelocateChannelMediaAction } from "../../../channels/bulkStorageActions";
+import { previewRelocationAction } from "../../../channels/[slug]/storageActions";
import {
OpsInputError,
ops,
+ optBool,
optString,
queueResponse,
reqSlugs,
@@ -9,7 +14,7 @@ import {
export const dynamic = "force-dynamic";
-// POST { slugs: string[], locationId?: string, root?: string }
+// POST { slugs: string[], locationId?: string, root?: string, dryRun?: boolean }
// -> { ok: true, queued, skipped }
//
// THE DESTINATION IS A LOCATION ID WHEREVER POSSIBLE — the root is resolved on
@@ -17,23 +22,43 @@ export const dynamic = "force-dynamic";
// cannot aim a batch somewhere a re-point has moved. `root` is the one-off
// escape hatch the panel also offers. A skip is not a failure: every slug that
// did not queue comes back with the same sentence the bulk bar shows.
+//
+// `dryRun: true` (release 19, A4) is the Storage panel's preview per channel
+// (previewRelocationAction) and starts nothing: { ok, dryRun, previews:
+// [{slug, preview} | {slug, error}] } — the bytes to copy and both sides'
+// free space. As on the panel, a preview of a CLASSIC channel (its big files
+// still real files in data/) first tiers it in place — same-disk renames into
+// its own media/, the step every move's preflight takes (release 17's ruling);
+// no byte leaves the corpus disk.
export async function POST(request: Request) {
- return ops(request, ["slugs", "locationId", "root"], async (body) => {
+ return ops(request, ["slugs", "locationId", "root", "dryRun"], async (body) => {
const slugs = reqSlugs(body, "slugs");
const locationId = optString(body, "locationId");
const root = optString(body, "root");
+ const dryRun = optBool(body, "dryRun") === true;
if ((locationId ? 1 : 0) + (root ? 1 : 0) !== 1) {
throw new OpsInputError(
'send exactly one of "locationId" (a location configured on /storage) or "root" (an absolute path)',
);
}
- return queueResponse(
- await bulkRelocateChannelMediaAction(
- slugs,
- locationId
- ? { kind: "location", locationId }
- : { kind: "custom", root: root as string },
- ),
- );
+ const dest = locationId
+ ? ({ kind: "location", locationId } as const)
+ : ({ kind: "custom", root: root as string } as const);
+ if (dryRun) {
+ const previews = [];
+ for (const slug of slugs) {
+ // The bulk move's own skip for a slug with no channel: the panel's
+ // preview is only ever asked about a channel that exists, and a
+ // zero-byte preview of a typo is the silent success this layer refuses.
+ if (!(await readChannelConfig(getPaths(), slug))) {
+ previews.push({ slug, error: "channel not found" });
+ continue;
+ }
+ const res = await previewRelocationAction(slug, dest);
+ previews.push(res.ok ? { slug, preview: res.preview } : { slug, error: res.error });
+ }
+ return NextResponse.json({ ok: true, dryRun: true, previews });
+ }
+ return queueResponse(await bulkRelocateChannelMediaAction(slugs, dest));
});
}
diff --git a/editor/app/api/ops/retry-bucket/route.test.ts b/editor/app/api/ops/retry-bucket/route.test.ts
@@ -0,0 +1,45 @@
+import test from "node:test";
+import assert from "node:assert/strict";
+import { writeFile } from "node:fs/promises";
+import path from "node:path";
+import { callPost, setupOpsCorpus } from "../_testCorpus";
+
+// Run with:
+// pnpm -C editor exec tsx --test "app/api/ops/retry-bucket/route.test.ts"
+//
+// retry-bucket and transcribe-bucket narrow a bucket with "ids", and every id
+// must be in it. Moved from e2e/ops-api.spec.ts (release 19, A2b): refused
+// before any job, from a report written by hand.
+
+const corpus = await setupOpsCorpus("title-filter-channel");
+const SLUG = "test-filter";
+// A report by hand, so the bucket is known without a regen.
+await writeFile(
+ path.join(corpus.transcripts, "channels", SLUG, "snapshot.json"),
+ JSON.stringify({ buckets: { downloadedNoTranscript: ["in-bucket-1"] } }),
+);
+const retry = (await import("./route")).POST;
+const transcribe = (await import("../transcribe-bucket/route")).POST;
+test.after(() => corpus.cleanup());
+
+test("retry-bucket and transcribe-bucket ids must be in the bucket: a stray is refused, named, and nothing runs", async () => {
+ const stray = await callPost(retry, {
+ slug: SLUG,
+ bucket: "downloadedNoTranscript",
+ ids: ["stray-aaa", "stray-bbb"],
+ });
+ assert.equal(stray.status, 400);
+ assert.match(String(stray.body.error), /2 of "ids" not in .*: stray-aaa, stray-bbb/);
+ // transcribe-bucket narrows the same bucket by the same rule.
+ const strayT = await callPost(transcribe, { slug: SLUG, ids: ["stray-ccc"] });
+ assert.equal(strayT.status, 400);
+ assert.match(String(strayT.body.error), /1 of "ids" not in .*: stray-ccc/);
+ // An empty list is not "the whole bucket".
+ const empty = await callPost(retry, {
+ slug: SLUG,
+ bucket: "downloadedNoTranscript",
+ ids: [],
+ });
+ assert.equal(empty.status, 400);
+ assert.deepEqual(await corpus.listJobIds(), []);
+});
diff --git a/editor/app/api/ops/scheduler/route.ts b/editor/app/api/ops/scheduler/route.ts
@@ -0,0 +1,16 @@
+import { buildSchedulerStatusPayload } from "../../../scheduler/status";
+import { readRoute } from "../_read";
+
+export const dynamic = "force-dynamic";
+
+// GET /api/ops/scheduler
+//
+// The sync scheduler — the /operations/sync payload
+// (`buildSchedulerStatusPayload`, the same one /api/scheduler/status polls):
+// each channel's resolved schedule and next due time, the recent tick log, and
+// the effective scheduler settings.
+export async function GET(request: Request) {
+ return readRoute(request, [], async () => ({
+ scheduler: await buildSchedulerStatusPayload(),
+ }));
+}
diff --git a/editor/app/api/ops/settings/patch.ts b/editor/app/api/ops/settings/patch.ts
@@ -0,0 +1,66 @@
+// WHAT `POST /api/ops/settings` REFUSES, as one pure function — testable
+// without a settings file, a schema or a server.
+
+// Keys another writer owns. Writing one raw would skip what that writer keeps
+// true, so the refusal names the writer:
+// - channelPriority: the priority write compiles the four lanes' trees in the
+// SAME write (saveSettings.ts' header); a raw one leaves them disagreeing.
+// - autoQueue: a lane's tree is its policy editor's, or compiled from the
+// channel priorities — and the editor refuses a tree edit while they rule.
+// - workers: /workers validates each worker (templates, URLs) on save.
+// - storage: a location is probed and its volume recorded when added or
+// re-pointed (/storage); a raw root would skip both.
+export const OWNED_KEYS: Readonly<Record<string, string>> = {
+ channelPriority: "it is written with the lanes' trees — use `pnpm ops channel-priority`",
+ autoQueue:
+ "a lane's switch, hold and runner are `pnpm ops lane`; its tree is the policy editor's on /operations, or compiled from channel priorities",
+ workers: "workers are saved on /workers; `pnpm ops workers` enables and disables them",
+ storage: "locations are added and re-pointed on /storage, which probes them; a channel's media moves with `pnpm ops relocate`",
+};
+
+function isObject(v: unknown): v is Record<string, unknown> {
+ return typeof v === "object" && v !== null && !Array.isArray(v);
+}
+
+// Every leaf of `sent` whose value in `saved` is not the same, as
+// "path: sent X, would be saved as Y". An object in `sent` is compared key by
+// key (a key it leaves out is not compared); an array must match in length and
+// element by element.
+export function coercedLeaves(sent: unknown, saved: unknown, at: string): string[] {
+ if (isObject(sent)) {
+ if (!isObject(saved)) return [`${at}: sent an object, would be saved as ${JSON.stringify(saved)}`];
+ return Object.keys(sent).flatMap((k) => coercedLeaves(sent[k], saved[k], `${at}.${k}`));
+ }
+ if (Array.isArray(sent)) {
+ if (!Array.isArray(saved) || saved.length !== sent.length) {
+ return [`${at}: sent ${JSON.stringify(sent)}, would be saved as ${JSON.stringify(saved)}`];
+ }
+ return sent.flatMap((v, i) => coercedLeaves(v, saved[i], `${at}[${i}]`));
+ }
+ return Object.is(sent, saved)
+ ? []
+ : [`${at}: sent ${JSON.stringify(sent)}, would be saved as ${JSON.stringify(saved)}`];
+}
+
+// The sentence to refuse `patch` with, or null. `known` is the schema's
+// top-level keys; `parsed` is the schema's reading of the merged settings.
+export function settingsPatchProblem(
+ patch: Record<string, unknown>,
+ known: readonly string[],
+ parsed: Record<string, unknown>,
+): string | null {
+ const keys = Object.keys(patch);
+ const unknown = keys.filter((k) => !known.includes(k));
+ if (unknown.length) {
+ return `unknown settings key(s): ${unknown.join(", ")} — SETTINGS.md lists them`;
+ }
+ const owned = keys.filter((k) => Object.hasOwn(OWNED_KEYS, k));
+ if (owned.length) {
+ return owned.map((k) => `"${k}" is not patched here: ${OWNED_KEYS[k]}`).join("; ");
+ }
+ const coerced = keys.flatMap((k) => coercedLeaves(patch[k], parsed[k], k));
+ if (coerced.length) {
+ return `not saved — the schema would not keep ${coerced.length === 1 ? "this value" : "these values"} as sent: ${coerced.join("; ")}`;
+ }
+ return null;
+}
diff --git a/editor/app/api/ops/settings/route.ts b/editor/app/api/ops/settings/route.ts
@@ -0,0 +1,75 @@
+import { NextResponse } from "next/server";
+import { getSettings, siteSettingsSchema } from "yt-dlp-transcript-common/lib/settings";
+import { mergeSettingsPatch, saveSettings } from "../../../settings/saveSettings";
+import { OpsInputError, ops, opsFail } from "../_lib";
+import { readRoute, redactSecrets } from "../_read";
+import { settingsPatchProblem } from "./patch";
+
+export const dynamic = "force-dynamic";
+
+// GET /api/ops/settings[?key=<top-level key>]
+//
+// settings.json as the editor reads it — through getSettings, so migrated and
+// defaulted, exactly the values every action sees (not the raw file: a key the
+// file leaves out answers with its default). `?key=autoQueue` answers one block.
+// Secrets are redacted (`_read.ts`). SETTINGS.md is the key table.
+export async function GET(request: Request) {
+ return readRoute(request, ["key"], async (q) => {
+ const settings = redactSecrets(getSettings()) as unknown as Record<string, unknown>;
+ const key = q.get("key");
+ if (key === null) return { settings };
+ if (!Object.hasOwn(settings, key)) {
+ return opsFail(`no settings key "${key}" — known: ${Object.keys(settings).sort().join(", ")}`);
+ }
+ return { key, value: settings[key] };
+ });
+}
+
+// POST { patch: { <top-level key>: <value>, … } }
+//
+// THE EDITOR'S ONE SETTINGS WRITER, over HTTP: `saveSettings(patch)` — the
+// function every settings form calls — with its merge rule (a block's keys
+// merge one level; an array, a scalar, or an object nested in a block
+// replaces). SETTINGS.md is the key table.
+//
+// REFUSED, before anything is written (`settings/patch.ts`):
+// - an unknown top-level key, named;
+// - a key another writer owns, because writing it raw breaks what that
+// writer keeps true (channelPriority, autoQueue, workers, storage — each
+// refusal names the command or page that writes it);
+// - a value the schema would not keep as sent. The schema never throws: it
+// coerces (clamps a number, drops an unknown enum, fills a default), so a
+// `{ ok: true }` would otherwise be the answer to a value that was never
+// saved. Every leaf the patch names is compared with what the schema makes
+// of it, and a difference is a 400 naming the path, what was sent and what
+// would be saved. A leaf the patch leaves out may still be filled with its
+// default — which is the merge rule, and is what the answer's `value` shows.
+//
+// Answers { ok, changed: [keys], value: { <key>: <saved value> } } — read back
+// after the write, secrets redacted.
+export async function POST(request: Request) {
+ return ops(request, ["patch"], async (body) => {
+ const patch = body.patch;
+ if (typeof patch !== "object" || patch === null || Array.isArray(patch)) {
+ throw new OpsInputError('"patch" is required and must be an object of top-level settings keys');
+ }
+ const keys = Object.keys(patch);
+ if (keys.length === 0) throw new OpsInputError('"patch" is empty — nothing to change');
+ const current = getSettings();
+ const problem = settingsPatchProblem(
+ patch as Record<string, unknown>,
+ Object.keys(siteSettingsSchema.shape),
+ siteSettingsSchema.parse(
+ mergeSettingsPatch(current, patch as Partial<typeof current>),
+ ) as unknown as Record<string, unknown>,
+ );
+ if (problem) throw new OpsInputError(problem);
+ await saveSettings(patch as Partial<typeof current>);
+ const saved = redactSecrets(getSettings()) as unknown as Record<string, unknown>;
+ return NextResponse.json({
+ ok: true,
+ changed: keys,
+ value: Object.fromEntries(keys.map((k) => [k, saved[k]])),
+ });
+ });
+}
diff --git a/editor/app/api/ops/sites/route.ts b/editor/app/api/ops/sites/route.ts
@@ -0,0 +1,28 @@
+import { getPaths } from "yt-dlp-transcript-common/lib/paths";
+import { listSites } from "yt-dlp-transcript-common/lib/site";
+import { isListedSite, isPrivateSite } from "yt-dlp-transcript-common/lib/siteSchema";
+import { readRoute } from "../_read";
+
+export const dynamic = "force-dynamic";
+
+// GET /api/ops/sites
+//
+// Every site, one row each, from its site.json: id, title, public URL
+// (`siteUrl`), Pages project, audience, whether the family lists it, whether it
+// publishes search, its publish policy and the channels it carries. The deep
+// status of a site's builds and deploys is `get publish`.
+export async function GET(request: Request) {
+ return readRoute(request, [], async () => ({
+ sites: listSites(getPaths()).map((s) => ({
+ siteId: s.siteId,
+ title: s.siteTitle,
+ siteUrl: s.siteUrl ?? null,
+ cloudflareProject: s.cloudflareProject ?? null,
+ audience: isPrivateSite(s) ? "private" : "public",
+ listed: isListedSite(s),
+ search: s.search !== false,
+ publish: s.publish ?? null,
+ channels: s.channels.map((c) => c.slug),
+ })),
+ }));
+}
diff --git a/editor/app/api/ops/storage/route.ts b/editor/app/api/ops/storage/route.ts
@@ -0,0 +1,15 @@
+import { buildStorage } from "../../../storage/buildStorage";
+import { readRoute } from "../_read";
+
+export const dynamic = "force-dynamic";
+
+// GET /api/ops/storage
+//
+// The /storage page's payload (`buildStorage`): every storage location with
+// whether it is mounted, its free space, its channels and the bytes of each
+// tier they hold, the corpus volume's text and clip tiers, and the saved-video
+// store. Probes are the page's own (memoised 10 s); sizes come off the
+// channels' reports, never a walk.
+export async function GET(request: Request) {
+ return readRoute(request, [], async () => ({ storage: await buildStorage() }));
+}
diff --git a/editor/app/api/ops/tag-videos/route.test.ts b/editor/app/api/ops/tag-videos/route.test.ts
@@ -0,0 +1,46 @@
+import test from "node:test";
+import assert from "node:assert/strict";
+import { access } from "node:fs/promises";
+import path from "node:path";
+import { callPost, setupOpsCorpus } from "../_testCorpus";
+
+// Run with:
+// pnpm -C editor exec tsx --test "app/api/ops/tag-videos/route.test.ts"
+//
+// The refusals of the curated-tag batch writer, before tags.json is touched.
+// Moved from e2e/ops-api.spec.ts (release 19, A2b); the writes it makes are
+// still driven over HTTP there.
+
+const corpus = await setupOpsCorpus("one-youtube-channel-with-data");
+const { POST } = await import("./route");
+test.after(() => corpus.cleanup());
+
+test("tag-videos refuses a traversing slug, a bad op and an unknown key", async () => {
+ const traversing = await callPost(POST, {
+ tag: "x",
+ op: "add",
+ videos: [{ slug: "../escape", id: "v" }],
+ });
+ assert.equal(traversing.status, 400);
+ assert.match(String(traversing.body.error), /not a valid channel slug/);
+
+ const badOp = await callPost(POST, {
+ tag: "x",
+ op: "pin",
+ videos: [{ slug: "test-youtube", id: "v" }],
+ });
+ assert.equal(badOp.status, 400);
+ assert.match(String(badOp.body.error), /add, remove, suppress, unsuppress/);
+
+ const unknown = await callPost(POST, {
+ tag: "x",
+ op: "add",
+ videos: [{ slug: "test-youtube", id: "v" }],
+ provenance: "me",
+ });
+ assert.equal(unknown.status, 400);
+ assert.match(String(unknown.body.error), /unknown key\(s\): provenance/);
+
+ // Nothing was written by any of the three.
+ await assert.rejects(access(path.join(corpus.transcripts, "tags.json")));
+});
diff --git a/editor/app/api/ops/transcribe-one/route.ts b/editor/app/api/ops/transcribe-one/route.ts
@@ -0,0 +1,36 @@
+import {
+ transcribeOneAction,
+ whisperVideoAction,
+} from "../../../channels/[slug]/videos/[id]/videoActions";
+import { jobResponse, ops, opsFail, optString, reqSlug, reqVideoId } from "../_lib";
+
+export const dynamic = "force-dynamic";
+
+// POST { slug, id, file? }
+//
+// One video, transcribed — the video page's two buttons:
+// - with "file" (an audio file in the video's dir, e.g. "audio.mp3"): that
+// file, through the worker pool (transcribeOneAction — the per-file
+// Transcribe on the Files list);
+// - without: the page's Transcribe (whisperVideoAction), which uses the audio
+// on disk, or first downloads it through the channel's managed path when
+// there is none (the disk floor applies to that download only).
+// A job on the video's platform queue: { ok, jobId }, so --wait follows it.
+// The file name is one path segment; a refused one is the action's sentence.
+export async function POST(request: Request) {
+ return ops(request, ["slug", "id", "file"], async (body) => {
+ const slug = reqSlug(body, "slug");
+ const id = reqVideoId(body, "id");
+ const file = optString(body, "file");
+ // One path segment, checked at the door as reqVideoId checks an id: the
+ // name is joined under the video's directory.
+ if (file !== undefined && (file === "" || /[/\\\0]/.test(file) || file.startsWith("."))) {
+ return opsFail(`"${file}" is not a file name in the video's directory`);
+ }
+ return jobResponse(
+ file !== undefined
+ ? await transcribeOneAction(slug, id, file)
+ : await whisperVideoAction(slug, id),
+ );
+ });
+}
diff --git a/editor/app/api/ops/workers/route.ts b/editor/app/api/ops/workers/route.ts
@@ -0,0 +1,48 @@
+import { NextResponse } from "next/server";
+import { buildWorkersPayload } from "../../../workers/buildWorkers";
+import {
+ disableWorkerAction,
+ drainWorkerAction,
+ enableWorkerAction,
+} from "../../../workers/actions";
+import { oneOf, ops, opsFail, reqStringArray } from "../_lib";
+import { readRoute, redactSecrets } from "../_read";
+
+export const dynamic = "force-dynamic";
+
+// GET /api/ops/workers
+//
+// The /workers page's live payload (`buildWorkersPayload`, the same one
+// /api/workers polls): every configured transcription worker — enabled or not,
+// busy or idle, what it is running — and whether the pool is paused. Behind the
+// token, with any remote worker's credential redacted.
+export async function GET(request: Request) {
+ return readRoute(request, [], async () => ({
+ workers: redactSecrets(buildWorkersPayload()),
+ }));
+}
+
+// POST { op: "enable" | "disable" | "drain", ids: [<worker id>, …] }
+//
+// The /workers row's switch and its Drain, per id (release 19, A4): enable
+// and disable are the switch — the pool's LIVE state, as the page sets it: a
+// restart brings back the launch default (/workers' "Set as default"); drain
+// stops giving the worker new transcriptions and lets the one in flight
+// finish. An unknown id is named and does not stop the others; the answer is
+// `ok: false` (400) when any id was refused, with every id's result.
+const OPS = ["enable", "disable", "drain"] as const;
+
+export async function POST(request: Request) {
+ return ops(request, ["op", "ids"], async (body) => {
+ const op = oneOf(body, "op", OPS);
+ const ids = [...new Set(reqStringArray(body, "ids"))];
+ const run = { enable: enableWorkerAction, disable: disableWorkerAction, drain: drainWorkerAction }[op];
+ const results: { id: string; ok: boolean; error?: string }[] = [];
+ for (const id of ids) results.push({ id, ...(await run(id)) });
+ const failed = results.filter((r) => !r.ok);
+ if (failed.length) {
+ return opsFail(failed.map((r) => r.error ?? r.id).join("; "), 400, { op, results });
+ }
+ return NextResponse.json({ ok: true, op, results });
+ });
+}
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/editor/app/cleanup/lib/loadCleanup.ts b/editor/app/cleanup/lib/loadCleanup.ts
@@ -124,6 +124,19 @@ export function heldTotalOf(held: HeldAudio): number {
);
}
+// ONE channel's ledger row — the row loadCleanupSummary builds for it, from the
+// same two reads (its brief, its failed-transcriptions list). For
+// `GET /api/ops/cleanup/<slug>`: asking about one channel must not pay for
+// every channel's snapshot. Null when the channel does not exist.
+export async function loadCleanupRow(
+ paths: Paths,
+ brief: ChannelBrief | null,
+): Promise<CleanupRow | null> {
+ if (!brief) return null;
+ const failedT = await loadFailedTranscriptions(paths, brief.slug);
+ return rowOf(brief, brief.snapshot, failedT.length);
+}
+
export async function loadCleanupSummary(
paths: Paths,
): Promise<CleanupSummary> {
diff --git a/editor/app/jobs/actions.ts b/editor/app/jobs/actions.ts
@@ -127,24 +127,31 @@ export async function retryJobAction(id: string): Promise<StreamActionResult> {
}
// Retry every currently-listed failed job that has a spec. Returns how many were
-// re-launched (kinds without a replay descriptor are skipped).
-export async function retryAllFailedAction(): Promise<{ count: number }> {
+// re-launched (kinds without a replay descriptor are skipped) and the new jobs'
+// ids — the button reads the count; `pnpm ops job retry-failed --wait` follows
+// the ids. The new job's stream is cancelled here: nothing consumes it (the
+// button never did — a dropped stream leaves the on-disk log running).
+export async function retryAllFailedAction(): Promise<{
+ count: number;
+ jobIds: string[];
+}> {
const registry = getRegistry();
const failed = registry
.list()
.filter((j): j is typeof j & { spec: JobSpec } =>
Boolean(j.status === "failed" && j.spec),
);
- let count = 0;
+ const jobIds: string[] = [];
for (const j of failed) {
const res = await runJobSpec(j.spec);
if (res.ok) {
registry.promote(res.jobId);
- count++;
+ void res.stream.cancel();
+ jobIds.push(res.jobId);
}
}
revalidatePath("/jobs");
- return { count };
+ return { count: jobIds.length, jobIds };
}
// ARM A LANE: write its scope as a tree, switch it on, and start its runner.
diff --git a/editor/app/operations/actions.ts b/editor/app/operations/actions.ts
@@ -13,6 +13,7 @@ import {
} from "yt-dlp-transcript-common/lib/pauseGates";
import { getWorkerPool } from "yt-dlp-transcript-common/jobs/workerPool";
import {
+ drainAutoRunner,
startAutoRunner,
stopAutoRunner,
} from "yt-dlp-transcript-common/controller/autoRunner";
@@ -120,8 +121,9 @@ export async function saveAutoQueueAction(
return { ok: true };
}
-// Explicit start/stop for the Start/Stop buttons (also reachable as
-// /api/auto-queue/control for e2e). Start is a no-op when the policy is disabled.
+// Explicit start/stop/drain for the lane header's Start, Stop and Drain buttons
+// (also reachable, behind the ops token, as /api/auto-queue/control and as
+// `pnpm ops lane`). Start is a no-op when the policy is disabled.
export async function startAutoQueueAction(
kind: AutoQueueKind,
): Promise<SaveResult> {
@@ -141,6 +143,15 @@ export async function stopAutoQueueAction(
return { ok: true };
}
+// Drain: the unit in flight finishes, no next one starts, the runner ends.
+export async function drainAutoQueueAction(
+ kind: AutoQueueKind,
+): Promise<SaveResult> {
+ drainAutoRunner(kind);
+ revalidateOperations();
+ return { ok: true };
+}
+
// Idle a runner until `untilMs` (epoch ms) without stopping it, or wake it now
// with null. Deliberately NOT part of saveAutoQueueAction: snoozing is a
// one-click operational act, and routing it through the policy form would mean a
diff --git a/editor/app/operations/components/RunnerOperationView.tsx b/editor/app/operations/components/RunnerOperationView.tsx
@@ -18,6 +18,11 @@ import { PolicyTreeEditor } from "./PolicyTreeEditor";
import { SnoozeControl } from "./SnoozeControl";
import { type Channel, formatClock, formatCooldown, leafOrder } from "./dispatch";
import { clearPlatformHoldAction } from "../pacingActions";
+import {
+ drainAutoQueueAction,
+ startAutoQueueAction,
+ stopAutoQueueAction,
+} from "../actions";
// The idle reasons that mean "paused or gated", not "waiting for work".
const LANE_PAUSE_REASONS: ReadonlySet<string> = new Set([
@@ -85,11 +90,13 @@ export function RunnerOperationView({
async (action: "start" | "stop" | "drain") => {
setBusy(true);
try {
- await fetch("/api/auto-queue/control", {
- method: "POST",
- headers: { "content-type": "application/json" },
- body: JSON.stringify({ kind, action }),
- });
+ // Server actions, not /api/auto-queue/control: that route is behind
+ // the ops token since release 19, and a page holds no token.
+ await (action === "start"
+ ? startAutoQueueAction(kind)
+ : action === "stop"
+ ? stopAutoQueueAction(kind)
+ : drainAutoQueueAction(kind));
await onRefresh();
} finally {
setBusy(false);
diff --git a/editor/e2e/auto-queue.spec.ts b/editor/e2e/auto-queue.spec.ts
@@ -153,6 +153,8 @@ async function startRunner(
kind: "transcription" | "download" = "transcription",
) {
const res = await request.post(`${baseUrl}/api/auto-queue/control`, {
+ // Behind the ops token since release 19 (A3): it starts and stops lanes.
+ headers: { authorization: "Bearer test-worker-token" },
data: { kind, action: "start" },
});
expect(res.ok()).toBeTruthy();
@@ -164,6 +166,8 @@ async function stopRunner(
kind: "transcription" | "download" = "transcription",
) {
await request.post(`${baseUrl}/api/auto-queue/control`, {
+ // Behind the ops token since release 19 (A3): it starts and stops lanes.
+ headers: { authorization: "Bearer test-worker-token" },
data: { kind, action: "stop" },
});
}
@@ -534,6 +538,8 @@ test("/jobs: the runner shows in its own labeled section, not Other", async ({
// Start the perpetual runner; it stays "running" (active) even once idle.
await page.request.post(`${baseUrl}/api/auto-queue/control`, {
+ // Behind the ops token since release 19 (A3): it starts and stops lanes.
+ headers: { authorization: "Bearer test-worker-token" },
data: { kind: "transcription", action: "start" },
});
diff --git a/editor/e2e/auto-subs-replace.spec.ts b/editor/e2e/auto-subs-replace.spec.ts
@@ -370,6 +370,8 @@ test("the auto-transcribe runner transcribes over auto-captions when opted in",
});
const started = await request.post(`${baseUrl}/api/auto-queue/control`, {
+ // Behind the ops token since release 19 (A3): it starts and stops lanes.
+ headers: { authorization: "Bearer test-worker-token" },
data: { kind: "transcription", action: "start" },
});
expect(started.ok()).toBeTruthy();
@@ -386,6 +388,8 @@ test("the auto-transcribe runner transcribes over auto-captions when opted in",
);
} finally {
await request.post(`${baseUrl}/api/auto-queue/control`, {
+ // Behind the ops token since release 19 (A3): it starts and stops lanes.
+ headers: { authorization: "Bearer test-worker-token" },
data: { kind: "transcription", action: "stop" },
});
}
diff --git a/editor/e2e/disk-space.spec.ts b/editor/e2e/disk-space.spec.ts
@@ -152,6 +152,8 @@ test("the unattended auto-download runner idles below the floor", async ({
});
const start = await request.post(`${baseUrl}/api/auto-queue/control`, {
+ // Behind the ops token since release 19 (A3): it starts and stops lanes.
+ headers: { authorization: "Bearer test-worker-token" },
data: { kind: "download", action: "start" },
});
expect(start.ok()).toBeTruthy();
@@ -173,6 +175,8 @@ test("the unattended auto-download runner idles below the floor", async ({
).toBe(false);
} finally {
await request.post(`${baseUrl}/api/auto-queue/control`, {
+ // Behind the ops token since release 19 (A3): it starts and stops lanes.
+ headers: { authorization: "Bearer test-worker-token" },
data: { kind: "download", action: "stop" },
});
}
diff --git a/editor/e2e/lane-runner.spec.ts b/editor/e2e/lane-runner.spec.ts
@@ -169,6 +169,8 @@ async function control(
action: "start" | "stop" | "drain",
): Promise<{ started?: boolean; blocked?: string }> {
const res = await request.post(`${baseUrl}/api/auto-queue/control`, {
+ // Behind the ops token since release 19 (A3): it starts and stops lanes.
+ headers: { authorization: "Bearer test-worker-token" },
data: { kind, action },
});
expect(res.ok()).toBeTruthy();
diff --git a/editor/e2e/ops-api.spec.ts b/editor/e2e/ops-api.spec.ts
@@ -7,26 +7,16 @@
// title-filter.spec.ts reads off the form, the busy-channel refusal the Storage
// panel shows — and about the files on disk, not about the routes' own shapes.
//
-// THE TOKEN GATE HAS THREE ANSWERS, AND ALL THREE ARE PINNED HERE.
-//
-// `WORKER_TOKEN` unset => 503, wrong or missing => 401, matching => the route
-// runs. The 503 is the consequential one: it is what stops an instance that
-// never set the variable being an open transcription server, and 401 in its
-// place would tell a scanner "there is a secret here, guess it".
-//
-// It used to be unreachable from a spec — the test server boots with
-// WORKER_TOKEN=test-worker-token (editor/package.json, dev:test) and one server
-// serves the whole suite, so nothing could observe the endpoint disabled. It is
-// reachable now because /api/test/worker-token turns the variable off and back
-// on INSIDE that server, which works only because `getWorkerToken()` reads
-// process.env per call (common/lib/workerToken.test.ts pins that, and covers
-// the branch table itself).
-//
-// ⚠️ THE VARIABLE IS PROCESS-WIDE AND THE SERVER OUTLIVES THE SPEC. Restoring
-// it is a `finally`, never a trailing line: one failed assertion with the token
-// still unset leaves every later /api/ops and /api/worker spec answering 503.
-
-import { mkdir, readdir, rm, writeFile } from "node:fs/promises";
+// WHAT IS NOT HERE ANY MORE (release 19, A2b). Every test that exercised only a
+// route's request and response — the token gate's three answers (401, 503,
+// pass), unknown keys, traversing slugs, and each refusal answered before any
+// job — is a unit test now, beside its route: `app/api/ops/_door.test.ts` and
+// `app/api/ops/<route>/route.test.ts`, run by `pnpm --filter editor test`
+// against a temp copy of the same fixtures (`app/api/ops/_testCorpus.ts`).
+// What stays needs the running server: a write that revalidates a page, a job
+// that runs, a queue the server holds.
+
+import { mkdir, rm, writeFile } from "node:fs/promises";
import { test, expect, type APIRequestContext } from "@playwright/test";
import { baseUrl } from "./baseUrl";
import {
@@ -75,332 +65,6 @@ async function settings(): Promise<void> {
await writeSettings({ minFreeDiskGB: 0 });
}
-test("the token gate answers 401 for a missing and for a wrong bearer", async ({
- request,
-}) => {
- await resetData("empty");
- for (const headers of [undefined, { authorization: "Bearer wrong" }]) {
- const post = await request.post(`${baseUrl}/api/ops/refresh-report`, {
- ...(headers ? { headers } : {}),
- data: { all: true },
- });
- expect(post.status(), JSON.stringify(headers)).toBe(401);
- // The READ side is gated by the same token, not merely the write side.
- const get = await request.get(`${baseUrl}/api/ops/channel/anything`, {
- ...(headers ? { headers } : {}),
- });
- expect(get.status(), JSON.stringify(headers)).toBe(401);
- }
-});
-
-// UNSET IS OFF, on the ops door and on the worker door alike — one branch,
-// shared, and neither surface may decide for itself that "no token configured"
-// means "let them in".
-test("with no token configured every guarded route answers 503", async ({
- request,
-}) => {
- await resetData("empty");
- const toggle = async (query: string) => {
- const res = await request.get(`${baseUrl}/api/test/worker-token?${query}`);
- expect(res.status(), query).toBe(200);
- return (await res.json()) as { ok: boolean; enabled: boolean };
- };
-
- const off = await toggle("unset=1");
- expect(off.enabled).toBe(false);
- try {
- // The ops door, with the RIGHT token: it is the surface being off that
- // answers, not the credential being wrong.
- const tags = await request.get(`${baseUrl}/api/ops/tags`, {
- headers: AUTH,
- });
- expect(tags.status()).toBe(503);
- expect(((await tags.json()) as OpsResponse).error).toMatch(
- /set WORKER_TOKEN to enable/,
- );
-
- // And the LAN worker door, which shares the branch.
- const health = await request.get(`${baseUrl}/api/worker/health`, {
- headers: AUTH,
- });
- expect(health.status()).toBe(503);
- expect(((await health.json()) as OpsResponse).error).toMatch(
- /set WORKER_TOKEN to enable/,
- );
-
- // With no token configured a MISSING header is still 503, not 401: there is
- // nothing to be unauthorized against.
- const bare = await request.get(`${baseUrl}/api/ops/tags`);
- expect(bare.status()).toBe(503);
- } finally {
- // NOT a trailing line. See the header: the server outlives this spec.
- const on = await toggle(`set=${TOKEN}`);
- expect(on.enabled).toBe(true);
- }
-
- // Restored, and the same request now works — which is also the assertion
- // that the `finally` above did what it claims.
- const after = await request.get(`${baseUrl}/api/ops/tags`, { headers: AUTH });
- expect(after.status()).toBe(200);
-});
-
-test("an unknown body key is a 400 that names the accepted keys", async ({
- request,
-}) => {
- await resetData("empty");
- // A misspelled key would otherwise get a cheerful { ok: true } and a channel
- // that did not change.
- const { status, body } = await ops(request, "sync", {
- slug: "x",
- fullSweep: true,
- });
- expect(status).toBe(400);
- expect(body.ok).toBe(false);
- expect(body.error).toContain("unknown key(s): fullSweep");
- expect(body.error).toContain("full");
-
- // So is a nested one, on the route whose body carries an object.
- const patch = await ops(request, "channel-config", {
- slug: "x",
- patch: { downloadFilterExcluded: "rerun" },
- });
- expect(patch.status).toBe(400);
- expect(patch.body.error).toContain("downloadFilterExcluded");
-});
-
-test("a traversing slug is refused at the door, on every route that takes one", async ({
- request,
-}) => {
- await resetData("title-filter-channel");
- const channelsDir = resolvePath("test-transcripts/channels");
- const before = (await readdir(channelsDir)).sort();
- expect(before).toEqual(["test-filter"]);
-
- // EVERY SLUG BELOW REACHES A path.join UNDER channelsDir, and the readers
- // swallow their own errors — so an unchecked traversing segment would fail
- // SILENTLY (an empty config read as "channel not found") rather than loudly,
- // and any future writer on that path would land outside the corpus. reqSlug
- // is one check for all of them; this is the assertion that it is wired to
- // each.
- const cases: [string, Record<string, unknown>][] = [
- ["metadata-scan", { slug: "../../escape" }],
- ["refresh-metadata", { slug: "../../escape", id: "abc123" }],
- ["sync", { slug: "../../escape" }],
- ["download-missing", { slug: "../../escape" }],
- ["import-video", { slug: "../../escape", url: "https://example.com/v" }],
- ["retry-bucket", { slug: "../../escape", bucket: "noTranscript" }],
- ["refresh-report", { slug: "../../escape" }],
- ["channel-config", { slug: "../../escape", patch: { cookieMode: "always" } }],
- ["channel-priority", { slugs: ["../../escape"], tier: "paused" }],
- ["relocate", { slugs: ["../../escape"], root: "/tmp/ops-api-never" }],
- ["relocate-back", { slugs: ["../../escape"] }],
- ["fetch-posts", { slug: "../../escape", older: true }],
- ["capture-posts", { slug: "../../escape", ids: ["1"] }],
- ["persist-videos", { items: [{ slug: "../../escape", id: "abc123" }] }],
- ];
- for (const [action, data] of cases) {
- const { status, body } = await ops(request, action, data);
- expect(status, action).toBe(400);
- expect(body.error, action).toMatch(/is not a valid channel slug/);
- }
-
- // The read route takes its slug as a path SEGMENT rather than in a body, so
- // it spells the same check out. A slash-bearing value is not the case to
- // assert here — the router never matches one to a single dynamic segment, so
- // it 404s before the handler exists. What DOES reach the handler is a
- // one-segment name CHANNEL_SLUG_RE still refuses, and a leading dot is the
- // one that matters: it is how a dotfile beside the channels dir would be
- // named at.
- const read = await request.get(`${baseUrl}/api/ops/channel/.escape`, {
- headers: AUTH,
- });
- expect(read.status()).toBe(400);
- expect(((await read.json()) as OpsResponse).error).toMatch(
- /is not a valid channel slug/,
- );
-
- // NOTHING WAS TOUCHED: the corpus still holds exactly the fixture channel,
- // and the fixture's own config is byte-identical.
- expect((await readdir(channelsDir)).sort()).toEqual(before);
- expect(
- await readJson<Record<string, unknown>>(
- "test-transcripts/channels/test-filter/config.json",
- ),
- ).toEqual({
- handling: "youtube",
- name: "Test Title Filter",
- url: "https://www.youtube.com/@example/videos",
- downloadFilter: { include: "guest" },
- });
-});
-
-test("retry-bucket and transcribe-bucket ids must be in the bucket: a stray is refused, named, and nothing runs", async ({
- request,
-}) => {
- await resetData("title-filter-channel");
- await settings();
- const SLUG = "test-filter";
- // A report by hand, so the bucket is known without a regen. The stray ids are
- // in NO bucket of ANY report, so a debounced regen landing between this write
- // and the call cannot change the answer.
- await writeFile(
- resolvePath(`test-transcripts/channels/${SLUG}/snapshot.json`),
- JSON.stringify({ buckets: { downloadedNoTranscript: ["in-bucket-1"] } }),
- );
- const before = await listJobIds();
- const stray = await ops(request, "retry-bucket", {
- slug: SLUG,
- bucket: "downloadedNoTranscript",
- ids: ["stray-aaa", "stray-bbb"],
- });
- expect(stray.status).toBe(400);
- expect(stray.body.error).toMatch(/2 of "ids" not in .*: stray-aaa, stray-bbb/);
- // transcribe-bucket narrows the same bucket by the same rule.
- const strayT = await ops(request, "transcribe-bucket", {
- slug: SLUG,
- ids: ["stray-ccc"],
- });
- expect(strayT.status).toBe(400);
- expect(strayT.body.error).toMatch(/1 of "ids" not in .*: stray-ccc/);
- // An empty list is not "the whole bucket".
- const empty = await ops(request, "retry-bucket", {
- slug: SLUG,
- bucket: "downloadedNoTranscript",
- ids: [],
- });
- expect(empty.status).toBe(400);
- expect(await listJobIds()).toEqual(before);
-});
-
-test("fetch-posts refuses a channel that is not social, full with older, and an older fetch its fetcher cannot do — before any job", async ({
- request,
-}) => {
- await resetData("title-filter-channel");
- await settings();
- // A social channel whose fetcher has no older-posts walk. Nothing below
- // reaches a fetch: every case is refused before a job exists.
- await writeChannelConfig("example-bsky", {
- handling: "transcribe",
- sourceKind: "social",
- platform: "bluesky",
- postFetcher: "bluesky-atproto",
- socialHandle: "example.bsky.social",
- name: "Example (Bluesky)",
- url: "https://bsky.app/profile/example.bsky.social",
- });
- const before = await listJobIds();
-
- // The sentence the action gives for a video channel.
- const video = await ops(request, "fetch-posts", { slug: "test-filter" });
- expect(video.status).toBe(400);
- expect(video.body.error).toBe("test-filter is not a social channel.");
- const videoOlder = await ops(request, "fetch-posts", { slug: "test-filter", older: true });
- expect(videoOlder.status).toBe(400);
- expect(videoOlder.body.error).toBe("test-filter is not a social channel.");
-
- // Two walks at once, on any channel.
- const both = await ops(request, "fetch-posts", {
- slug: "example-bsky",
- full: true,
- older: true,
- });
- expect(both.status).toBe(400);
- expect(both.body.error).toMatch(/different walks — run one at a time/);
- const bothVideo = await ops(request, "fetch-posts", {
- slug: "test-filter",
- full: true,
- older: true,
- });
- expect(bothVideo.status).toBe(400);
- expect(bothVideo.body.error).toMatch(/different walks/);
-
- // A fetcher with no older walk, a floor that is not a date, a floor
- // without older, and a limit that is not a count.
- const bsky = await ops(request, "fetch-posts", { slug: "example-bsky", older: true });
- expect(bsky.status).toBe(400);
- expect(bsky.body.error).toMatch(/cannot fetch older posts/);
- const floor = await ops(request, "fetch-posts", {
- slug: "example-bsky",
- older: true,
- floor: "2020-13-01",
- });
- expect(floor.status).toBe(400);
- expect(floor.body.error).toMatch(/is not a date/);
- const floorAlone = await ops(request, "fetch-posts", {
- slug: "example-bsky",
- floor: "2020-01-01",
- });
- expect(floorAlone.status).toBe(400);
- expect(floorAlone.body.error).toMatch(/only to an older-posts fetch/);
- const limit = await ops(request, "fetch-posts", { slug: "example-bsky", limit: 0 });
- expect(limit.status).toBe(400);
- expect(limit.body.error).toMatch(/"limit" must be a whole number above zero/);
-
- expect(await listJobIds()).toEqual(before);
-});
-
-test("capture-posts refuses what it cannot capture, and an id not in the archive — before any job", async ({
- request,
-}) => {
- await resetData("title-filter-channel");
- await settings();
- // Nothing below reaches X: every case is refused before a job exists.
- await writeChannelConfig("example-bsky", {
- handling: "transcribe",
- sourceKind: "social",
- platform: "bluesky",
- postFetcher: "bluesky-atproto",
- socialHandle: "example.bsky.social",
- name: "Example (Bluesky)",
- url: "https://bsky.app/profile/example.bsky.social",
- });
- await writeChannelConfig("example-x", {
- handling: "transcribe",
- sourceKind: "social",
- platform: "twitter",
- postFetcher: "x-gallery-dl",
- socialHandle: "example_user",
- name: "Example (X)",
- url: "https://x.com/example_user",
- });
- const before = await listJobIds();
-
- const video = await ops(request, "capture-posts", { slug: "test-filter", ids: ["1"] });
- expect(video.status).toBe(400);
- expect(video.body.error).toBe("test-filter is not a social channel.");
-
- const bsky = await ops(request, "capture-posts", { slug: "example-bsky", ids: ["1"] });
- expect(bsky.status).toBe(400);
- expect(bsky.body.error).toMatch(/cannot capture posts/);
-
- const neither = await ops(request, "capture-posts", {
- slug: "example-x",
- ids: ["1"],
- shots: false,
- media: false,
- });
- expect(neither.status).toBe(400);
- expect(neither.body.error).toMatch(/both the screenshot and the media are turned off/);
-
- // The channel's posts archive is empty: every id is a stray, named.
- const stray = await ops(request, "capture-posts", { slug: "example-x", ids: ["111", "222"] });
- expect(stray.status).toBe(400);
- expect(stray.body.error).toBe("2 id(s) not in example-x's posts archive: 111, 222");
-
- // The body's shape.
- const noIds = await ops(request, "capture-posts", { slug: "example-x" });
- expect(noIds.status).toBe(400);
- expect(noIds.body.error).toMatch(/"ids" is required/);
- const badFlag = await ops(request, "capture-posts", { slug: "example-x", ids: ["1"], shots: "yes" });
- expect(badFlag.status).toBe(400);
- expect(badFlag.body.error).toMatch(/"shots" must be a boolean/);
- const unknown = await ops(request, "capture-posts", { slug: "example-x", ids: ["1"], limit: 5 });
- expect(unknown.status).toBe(400);
- expect(unknown.body.error).toMatch(/unknown key\(s\): limit/);
-
- expect(await listJobIds()).toEqual(before);
-});
-
// THE CHANNEL'S LIFECYCLE, OVER HTTP. create-channel is the New form,
// rename-channel and delete-channel the Danger zone's two forms, channel-config
// grows the Configure form's Sites section and the rack's two exclusion
@@ -911,6 +575,67 @@ test("lane flips a hold, and /api/auto-queue/status agrees", async ({
expect(bogus.body.error).toMatch(/transcription, download, digest, backfill/);
});
+// --- jobs (release 19, A2) --------------------------------------------------
+//
+// The /jobs row buttons over HTTP. Two fabricated jobs on one queue (the
+// /api/test/stuck-job harness: a running head that never finishes and one
+// queued behind it) are what a real backlog looks like, without the race of a
+// real one. The success paths live here because each verb is a server action
+// that revalidates /jobs, which only a running Next server can do.
+test("jobs over ops: the live head, one job's place, and the row buttons", async ({
+ request,
+}) => {
+ const queue = `ops-jobs-${Date.now()}`;
+ const stuck = async () =>
+ (await (await request.get(`${baseUrl}/api/test/stuck-job?queue=${queue}`)).json()) as {
+ id: string;
+ };
+ const head = await stuck();
+ const waiting = await stuck();
+ try {
+ const active = await request.get(`${baseUrl}/api/ops/jobs?active=1`, { headers: AUTH });
+ expect(active.status()).toBe(200);
+ const rows = ((await active.json()) as {
+ jobs: { id: string; status: string; queue?: { key: string; position: number } }[];
+ }).jobs.filter((j) => j.queue?.key === queue);
+ expect(rows.map((j) => [j.id, j.status, j.queue?.position])).toEqual([
+ [head.id, "running", 0],
+ [waiting.id, "queued", 1],
+ ]);
+
+ const one = await request.get(`${baseUrl}/api/ops/job/${waiting.id}`, { headers: AUTH });
+ const job = ((await one.json()) as {
+ job: { queue: { position: number; queued: number; head?: { id: string } } };
+ }).job;
+ expect(job.queue.position).toBe(1);
+ expect(job.queue.head?.id).toBe(head.id);
+ const tail = await request.get(`${baseUrl}/api/ops/job/${head.id}?tail=5`, { headers: AUTH });
+ expect(((await tail.json()) as { tail: string[] }).tail).toEqual([
+ "[fake] wedged job holding its slot",
+ ]);
+
+ // Already at the front of its queue: the row's own refusal, named.
+ const promote = await ops(request, "job", { verb: "promote", ids: [waiting.id] });
+ expect(promote.status).toBe(400);
+ expect(promote.body.error).toMatch(/1 of 1 not promoted: .*already at the front/);
+
+ // One unknown id does not stop the known one.
+ const cancel = await ops(request, "job", { verb: "cancel", ids: [waiting.id, "NOSUCHJOB"] });
+ expect(cancel.status).toBe(400);
+ expect(cancel.body.error).toMatch(/1 of 2 not cancelled: NOSUCHJOB/);
+ const cancelled = await request.get(`${baseUrl}/api/ops/job/${waiting.id}`, { headers: AUTH });
+ expect(((await cancelled.json()) as { job: { status: string } }).job.status).toBe("cancelled");
+ } finally {
+ // The head never ends on its own: free its slot whatever happened above.
+ const released = await ops(request, "job", { verb: "force-release", ids: [head.id] });
+ expect(released.body.ok).toBe(true);
+ }
+ const after = await request.get(`${baseUrl}/api/ops/job/${head.id}`, { headers: AUTH });
+ const ended = ((await after.json()) as { job: { status: string; queue?: unknown } }).job;
+ expect(ended.status).toBe("cancelled");
+ expect(ended.queue).toBeUndefined();
+});
+
// --- curated tags -----------------------------------------------------------
//
// The same two claims as everywhere else in this file: the route is an adapter
@@ -1025,39 +750,6 @@ test("tags defines a vocabulary, and tag-videos writes pins with their source",
expect(left.assignments[`${slug}/${id}`].suppressed).toEqual(["eva-collab"]);
});
-test("tag-videos refuses a traversing slug, a bad op and an unknown key", async ({
- request,
-}) => {
- await resetData("one-youtube-channel-with-data");
- const traversing = await ops(request, "tag-videos", {
- tag: "x",
- op: "add",
- videos: [{ slug: "../escape", id: "v" }],
- });
- expect(traversing.status).toBe(400);
- expect(traversing.body.error).toContain("not a valid channel slug");
-
- const badOp = await ops(request, "tag-videos", {
- tag: "x",
- op: "pin",
- videos: [{ slug: "test-youtube", id: "v" }],
- });
- expect(badOp.status).toBe(400);
- expect(badOp.body.error).toContain("add, remove, suppress, unsuppress");
-
- const unknown = await ops(request, "tag-videos", {
- tag: "x",
- op: "add",
- videos: [{ slug: "test-youtube", id: "v" }],
- provenance: "me",
- });
- expect(unknown.status).toBe(400);
- expect(unknown.body.error).toContain("unknown key(s): provenance");
-
- // Nothing was written by any of the three.
- expect(await pathExists("test-transcripts/tags.json")).toBe(false);
-});
-
// --- keep-videos: the bulk do-not-clean toggle --------------------------------
//
// The marker is the video page's own do-not-clean.json; what this pins is the
@@ -1131,294 +823,17 @@ test("keep-videos marks the matching video only, after a dry run that writes not
expect(badRegex.body.error).toContain("invalid pattern");
});
-test("persist-videos buckets a list in a dry run, and refuses a bad body and a low disk — before any job", async ({
- request,
-}) => {
- await resetData("one-youtube-channel-with-data");
- await settings();
- const slug = "test-youtube";
- const video = "20240101_test1234567";
- const before = await listJobIds();
-
- type Bucket = { count: number; items: { slug: string; id: string }[] };
- type PersistBody = OpsResponse & {
- dryRun?: boolean;
- plan?: Record<"saved" | "wrongHeight" | "toFetch" | "noUrl" | "unknown", Bucket> & {
- willFetch: number;
- };
- };
- const items = [
- { slug, id: video },
- { slug, id: "missing12345" },
- { slug: "no-such-channel", id: "abc123" },
- ];
- const dry = await ops(request, "persist-videos", {
- items,
- format: "video_720",
- dryRun: true,
- });
- expect(dry.status, JSON.stringify(dry.body)).toBe(200);
- const plan = (dry.body as PersistBody).plan!;
- expect((dry.body as PersistBody).dryRun).toBe(true);
- expect(plan.toFetch.items).toEqual([{ slug, id: video }]);
- expect(plan.unknown.count).toBe(2);
- expect(plan.saved.count).toBe(0);
- expect(plan.willFetch).toBe(1);
-
- // Nothing to fetch is an answer, not a job.
- const none = await ops(request, "persist-videos", {
- items: [{ slug: "no-such-channel", id: "abc123" }],
- });
- expect(none.status, JSON.stringify(none.body)).toBe(200);
- expect(none.body.jobIds).toEqual([]);
-
- // The body's shape.
- const cases: [Record<string, unknown>, RegExp][] = [
- [{}, /"items" is required/],
- [{ items: [] }, /"items" is required/],
- [{ items: [{ slug, id: video, height: 720 }] }, /unknown key\(s\): height/],
- [{ items: [{ slug, id: "../escape" }] }, /is not a video id/],
- [{ items, format: "1080p" }, /"format" must be one of original, video_720/],
- [{ items, replace: "always" }, /"replace" must be one of never, above-height/],
- [{ items, gapMs: -1 }, /"gapMs" must be a whole number, zero or above/],
- [{ items, minFreeMemMb: "4096" }, /"minFreeMemMb" must be a whole number/],
- [{ items, dryRun: "yes" }, /"dryRun" must be a boolean/],
- [{ items, ids: ["x"] }, /unknown key\(s\): ids/],
- ];
- for (const [data, error] of cases) {
- const res = await ops(request, "persist-videos", data);
- expect(res.status, JSON.stringify(data)).toBe(400);
- expect(res.body.error, JSON.stringify(data)).toMatch(error);
- }
-
- // A real run asks the disk floor before it queues anything.
- await writeSettings({ minFreeDiskGB: 1_000_000 });
- try {
- const low = await ops(request, "persist-videos", { items });
- expect(low.status, JSON.stringify(low.body)).toBe(400);
- expect(low.body.error).toMatch(/^Low disk space/);
- } finally {
- await settings();
- }
-
- expect(await listJobIds()).toEqual(before);
-});
-
-// --- the two build routes speak the same body ---------------------------------
-//
-// They did not. build-site took `siteIds` (a list) and build-deploy took
-// `siteId` (one), so the two halves of the same runbook sentence needed
-// different JSON and the only way to learn which was to get a 400. Both now
-// accept both keys; sending BOTH is still refused, because a caller with two
-// ideas about what to build should not have one picked for it.
-// Job sidecars on disk — the only way to say "and nothing was queued".
-async function listJobIds(): Promise<string[]> {
- return (await readdir(resolvePath("test-transcripts/.jobs")).catch(() => []))
- .filter((f) => f.endsWith(".meta.json"))
- .sort();
-}
-
// Hold the `publish` queue with a fabricated running job (/api/test/stuck-job;
// never released by the caller), so any build or deploy stage that a refusal
// under test failed to stop can only QUEUE, never run — every build and deploy
-// is a publish stage on that one queue since release 18. The fabricated holder
-// writes no .meta.json, so it is invisible to listJobIds; take `before` after
-// this. The next resetData cancels jobs newest first, so a queued stage is
-// removed before the holder's slot is freed and nothing is ever promoted.
+// is a publish stage on that one queue since release 18. The next resetData
+// cancels jobs newest first, so a queued stage is removed before the holder's
+// slot is freed and nothing is ever promoted.
async function holdPublishQueue(request: APIRequestContext) {
const hold = await request.get(`${baseUrl}/api/test/stuck-job?queue=publish`);
expect(hold.ok()).toBe(true);
}
-test("build-site and build-deploy each take siteId or siteIds, and refuse both or neither", async ({
- request,
-}) => {
- await resetData("title-filter-channel");
- await settings();
-
- for (const action of ["build-site", "build-deploy"]) {
- const both = await ops(request, action, {
- siteId: "a",
- siteIds: ["a"],
- });
- expect(both.status, action).toBe(400);
- expect(both.body.error, action).toContain("not both");
-
- const neither = await ops(request, action, {});
- expect(neither.status, action).toBe(400);
- expect(neither.body.error, action).toContain("siteId");
-
- // `all` is still exclusive of either spelling.
- const withAll = await ops(request, action, { all: true, siteId: "a" });
- expect(withAll.status, action).toBe(400);
- expect(withAll.body.error, action).toContain("not both");
-
- // A MALFORMED ID IS A 400 BEFORE ANY JOB STARTS, and that is the point of
- // validating in reqSiteIds rather than letting getSite throw: an id list of
- // ["good", "BAD"] used to queue the first build, throw on the second and
- // answer 500 with no ids at all — a real build running that nothing was
- // watching, and a --wait exiting 1 about it.
- const before = await listJobIds();
- const bad = await ops(request, action, { siteIds: ["buildsite", "BAD ID"] });
- expect(bad.status, action).toBe(400);
- expect(bad.body.error, action).toContain("not a valid site id");
- expect(bad.body.jobs, action).toBeUndefined();
- expect(await listJobIds(), action).toEqual(before);
- }
-});
-
-// --- preview deploys ----------------------------------------------------------
-//
-// NOTHING HERE STARTS A DEPLOY JOB. Wrangler has a fake since release 18
-// (fake-wrangler.mjs; publish.spec deploys through it), but the preview rules
-// are decided BEFORE a job exists — that is the whole of what they are — so
-// these assert that half, and prove it by watching the job sidecars not
-// appear.
-test("deploy-site and build-deploy refuse a preview name that is not one, before any job", async ({
- request,
-}) => {
- await resetData("title-filter-channel");
- await settings();
- await writeSite("previewsite", { cloudflareProject: "proj" });
-
- for (const action of ["deploy-site", "build-deploy"]) {
- // "main" is the production branch: deploying there is not a preview, it is
- // the live site, and the refusal says so rather than shipping it.
- const before = await listJobIds();
- const main = await ops(request, action, {
- siteId: "previewsite",
- preview: "main",
- });
- expect(main.status, action).toBe(400);
- expect(main.body.error, action).toContain("production branch");
- expect(await listJobIds(), action).toEqual(before);
-
- // Uppercase is refused rather than lowercased FOR the caller: silently
- // rewriting "Main" into "main" would deploy to production.
- const upper = await ops(request, action, {
- siteId: "previewsite",
- preview: "Main",
- });
- expect(upper.status, action).toBe(400);
- expect(upper.body.error, action).toContain("not a valid preview branch name");
- expect(upper.body.error, action).not.toContain("production branch");
-
- const spaced = await ops(request, action, {
- siteId: "previewsite",
- preview: "bad name",
- });
- expect(spaced.status, action).toBe(400);
- expect(spaced.body.error, action).toContain("not a valid preview branch name");
-
- const tooLong = await ops(request, action, {
- siteId: "previewsite",
- preview: "a".repeat(29),
- });
- expect(tooLong.status, action).toBe(400);
- expect(tooLong.body.error, action).toContain("at most 28 characters");
-
- expect(await listJobIds(), action).toEqual(before);
- }
-});
-
-test("a valid preview is accepted and reaches the action, on both deploy routes", async ({
- request,
-}) => {
- await resetData("title-filter-channel");
- await settings();
- // NO Cloudflare project, on purpose: the action's own refusal is the last
- // gate before wrangler, so a 400 saying THAT — rather than "unknown key:
- // preview" or a branch complaint — proves the name passed validation and the
- // route got all the way to the action, without a deploy ever starting.
- await writeSite("noproj", {});
-
- for (const action of ["deploy-site", "build-deploy"]) {
- const before = await listJobIds();
- const { status, body } = await ops(request, action, {
- siteId: "noproj",
- preview: "tags-exclude",
- });
- expect(status, action).toBe(400);
- expect(body.error, action).toContain("no Cloudflare Pages project");
- expect(body.error, action).not.toContain("unknown key");
- expect(await listJobIds(), action).toEqual(before);
- }
-});
-
-test("deploy-site refuses a site that was never built, in the deploy stage's words, before any job", async ({
- request,
-}) => {
- await resetData("title-filter-channel");
- await settings();
- await writeSite("previewsite", { cloudflareProject: "proj" });
-
- // A deploy ships the site's OWN bundle (export/.export-builds/<id>/out,
- // stamped built.json) — never "whatever export/out holds", which is how a
- // deploy-only used to ship one site's build to another's project. A site
- // with no build is refused before any job, with the stage's sentence.
- const before = await listJobIds();
- const { status, body } = await ops(request, "deploy-site", {
- siteId: "previewsite",
- });
- expect(status).toBe(400);
- expect(body.error).toMatch(
- /^no build of previewsite in .*\.export-builds\/previewsite — archilyzer publish build previewsite$/,
- );
- expect(await listJobIds()).toEqual(before);
-
- // A preview is refused for the same reason and just as early.
- const preview = await ops(request, "deploy-site", {
- siteId: "previewsite",
- preview: "tags-exclude",
- });
- expect(preview.status).toBe(400);
- expect(preview.body.error).toContain("no build of previewsite");
- expect(await listJobIds()).toEqual(before);
-});
-
-test("deploy-site takes siteId or siteIds, and has no all", async ({
- request,
-}) => {
- await resetData("title-filter-channel");
- await settings();
-
- const both = await ops(request, "deploy-site", { siteId: "a", siteIds: ["a"] });
- expect(both.status).toBe(400);
- expect(both.body.error).toContain("not both");
-
- const neither = await ops(request, "deploy-site", {});
- expect(neither.status).toBe(400);
- expect(neither.body.error).toContain("siteId");
-
- const before = await listJobIds();
- const bad = await ops(request, "deploy-site", { siteIds: ["deploysite", "BAD ID"] });
- expect(bad.status).toBe(400);
- expect(bad.body.error).toContain("not a valid site id");
- expect(await listJobIds()).toEqual(before);
-
- // Deploy-only has no all-sites form — build-deploy owns that — so `all` is
- // an unknown key here rather than a second spelling of it.
- const all = await ops(request, "deploy-site", { all: true });
- expect(all.status).toBe(400);
- expect(all.body.error).toContain("unknown key");
-});
-
-test("build-deploy refuses a preview alongside all rather than building everything", async ({
- request,
-}) => {
- await resetData("title-filter-channel");
- await settings();
-
- const before = await listJobIds();
- const { status, body } = await ops(request, "build-deploy", {
- all: true,
- preview: "tags-exclude",
- });
- expect(status).toBe(400);
- expect(body.error).toContain('"preview" is not supported with "all"');
- expect(await listJobIds()).toEqual(before);
-});
-
test("build-site with a bare siteId starts its build stage (the index first: there is none)", async ({
request,
}) => {
@@ -1450,108 +865,6 @@ test("build-site with a bare siteId starts its build stage (the index first: the
.toBe("publish-build-site");
});
-// The hub's deploy path (release 7). Every refusal here is answered BEFORE a
-// job exists, which is what lets a runbook's `pnpm ops deploy-hub --wait` fail
-// fast instead of queueing a deploy that can only fail. Both queues are held
-// (release 11, O4 review): a regression of a refusal below would otherwise
-// start a real hub build into export/out, or a real deploy.
-test("deploy-hub refuses no project, the homepage's project, and a bundle that is not the hub", async ({
- request,
-}) => {
- await resetData("title-filter-channel");
- await settings();
- await holdPublishQueue(request);
- const before = await listJobIds();
-
- // The fixture has no homepage.json at all: no project.
- const none = await ops(request, "deploy-hub", {});
- expect(none.status).toBe(400);
- expect(none.body.error).toBe(
- "The hub has no Cloudflare Pages project configured — set it on /sites under Hub.",
- );
-
- // The homepage's project is refused by name — homepage.json said
- // "archilyzer" before the hub could deploy, and a hub deployed there would
- // replace the software's own site.
- const hubFile = resolvePath("test-transcripts/sites/_homepage/homepage.json");
- await mkdir(resolvePath("test-transcripts/sites/_homepage"), { recursive: true });
- await writeFile(hubFile, JSON.stringify({ cloudflareProject: "archilyzer" }));
- const homepages = await ops(request, "deploy-hub", {});
- expect(homepages.status).toBe(400);
- expect(homepages.body.error).toContain('"archilyzer", which is the Archilyzer homepage\'s');
-
- // A real project, but no hub build: the hub's bundle is the FIXTURE's
- // (test-transcripts/.export-builds/_hub, stamped built.json) since release
- // 18, so this refusal fires in any checkout — export/out is not read.
- await writeFile(hubFile, JSON.stringify({ cloudflareProject: "archilyzer-hub" }));
- const unbuilt = await ops(request, "deploy-hub", { preview: "hub-check" });
- expect(unbuilt.status).toBe(400);
- expect(unbuilt.body.error).toMatch(/^no build of _hub in .*\.export-builds\/_hub — archilyzer publish hub$/);
-
- // build-hub: a preview without a deploy is a mistake, not a build.
- const previewOnly = await ops(request, "build-hub", { preview: "hub-check" });
- expect(previewOnly.status).toBe(400);
- expect(previewOnly.body.error).toContain('"preview" needs "deploy": true');
-
- expect(await listJobIds()).toEqual(before);
-});
-
-// The homepage's deploy path (release 11, slice O4). Every refusal is answered
-// BEFORE a job exists, like deploy-hub's.
-//
-// NOTHING HERE CAN REACH WRANGLER (the fake's, even), and three things make
-// sure of it:
-// 1. "built" is the homepage's stamp, `_homepage/built.json` in the fixture's
-// .export-builds (release 18) — the checkout's homepage/out is not asked
-// until a stamp says it was built, and the fixture has none;
-// 2. the missing-build call asks for a PREVIEW, never production;
-// 3. the whole test runs with the `publish` queue held by a fabricated job
-// (holdPublishQueue). A stage that a refusal failed to stop —
-// `build-homepage {preview}` is a build if the "preview needs deploy"
-// guard regresses — would only QUEUE; the listJobIds assertion then fails,
-// and the next resetData cancels jobs newest first, so the queued stage is
-// removed before the holder's slot is freed.
-test("deploy-homepage refuses a bad preview name and a missing build; build-homepage refuses a preview without a deploy", async ({
- request,
-}) => {
- await resetData("empty");
- await settings();
- await holdPublishQueue(request);
- const before = await listJobIds();
-
- // A bad preview name is judged before homepage/out is even looked at, so
- // these hold in any checkout.
- const production = await ops(request, "deploy-homepage", { preview: "main" });
- expect(production.status).toBe(400);
- expect(production.body.error).toBe(
- '"main" is the production branch; a preview needs another name.',
- );
- const shape = await ops(request, "deploy-homepage", { preview: "Not_Valid" });
- expect(shape.status).toBe(400);
- expect(shape.body.error).toMatch(/not a valid preview branch name/);
- // And on the build-then-deploy, before anything is built.
- const buildDeploy = await ops(request, "build-homepage", {
- deploy: true,
- preview: "main",
- });
- expect(buildDeploy.status).toBe(400);
- expect(buildDeploy.body.error).toContain("is the production branch");
-
- // build-homepage: a preview without a deploy is a mistake, not a build.
- const previewOnly = await ops(request, "build-homepage", { preview: "home-check" });
- expect(previewOnly.status).toBe(400);
- expect(previewOnly.body.error).toContain('"preview" needs "deploy": true');
-
- // No build: the homepage's stamp is the fixture's.
- const unbuilt = await ops(request, "deploy-homepage", { preview: "home-check" });
- expect(unbuilt.status).toBe(400);
- expect(unbuilt.body.error).toMatch(
- /^no build of _homepage in .*\.export-builds\/_homepage — archilyzer publish homepage$/,
- );
-
- expect(await listJobIds()).toEqual(before);
-});
-
// --- publish (release 18) ---------------------------------------------------
//
// ONE route for the stages, a verb in the body, and its GET the status. The
@@ -1609,65 +922,3 @@ test("publish: each verb enqueues its run of stages on the publish queue; a dupl
expect(stale.ok).toBe(true);
expect(stale.jobs?.find((j) => j.kind === "build-site")?.existing).toBe(true);
});
-
-test("publish: refusals before any job — a never-built deploy, a preview that is not one, a key the verb does not take", async ({
- request,
-}) => {
- await resetData("title-filter-channel");
- await settings();
- await writeSite("pubsite", { cloudflareProject: "pubproj" });
- await holdPublishQueue(request);
- const before = await listJobIds();
-
- const unbuilt = await ops(request, "publish", { verb: "deploy", siteId: "pubsite", preview: "r18" });
- expect(unbuilt.status).toBe(400);
- expect(unbuilt.body.error).toMatch(/^no build of pubsite in .* — archilyzer publish build pubsite$/);
-
- const main = await ops(request, "publish", { verb: "deploy", siteId: "pubsite", preview: "main" });
- expect(main.status).toBe(400);
- expect(main.body.error).toContain("production branch");
-
- const verb = await ops(request, "publish", { verb: "launch" });
- expect(verb.status).toBe(400);
- expect(verb.body.error).toContain('"verb" must be one of index, build, deploy, hub, homepage, now, stale');
-
- const extra = await ops(request, "publish", { verb: "index", siteId: "pubsite" });
- expect(extra.status).toBe(400);
- expect(extra.body.error).toContain('verb "index" takes no other keys');
-
- const hubPreview = await ops(request, "publish", { verb: "hub", preview: "r18" });
- expect(hubPreview.status).toBe(400);
- expect(hubPreview.body.error).toContain('need "deploy": true');
-
- expect(await listJobIds()).toEqual(before);
-});
-
-test("GET publish is the publish status: the index, the lane, a row per target", async ({ request }) => {
- await resetData("title-filter-channel");
- await settings();
- await writeSite("pubsite", { cloudflareProject: "pubproj" });
-
- const res = await request.get(`${baseUrl}/api/ops/publish`, {
- headers: { authorization: "Bearer test-worker-token" },
- });
- expect(res.status()).toBe(200);
- const status = (await res.json()) as {
- ok: boolean;
- index: { chip: { text: string } };
- lane: { held: boolean; enabled: boolean };
- sites: { target: string; chips: { built: { text: string } } }[];
- hub: { target: string };
- homepage: { target: string };
- plan: { steps: { kind: string }[] };
- };
- expect(status.ok).toBe(true);
- expect(status.index.chip.text).toBe("no index yet");
- expect(status.lane).toMatchObject({ held: false, enabled: false });
- expect(status.sites.map((x) => x.target)).toEqual(["pubsite"]);
- expect(status.sites[0].chips.built.text).toBe("update the index first");
- expect([status.hub.target, status.homepage.target]).toEqual(["_hub", "_homepage"]);
- expect(status.plan.steps.map((x) => x.kind)).toEqual(["update-index"]);
-
- const anon = await request.get(`${baseUrl}/api/ops/publish`);
- expect(anon.status()).toBe(401);
-});
diff --git a/editor/e2e/pacing.spec.ts b/editor/e2e/pacing.spec.ts
@@ -117,6 +117,8 @@ async function control(
action: "start" | "stop",
) {
const res = await request.post(`${baseUrl}/api/auto-queue/control`, {
+ // Behind the ops token since release 19 (A3): it starts and stops lanes.
+ headers: { authorization: "Bearer test-worker-token" },
data: { kind: "download", action },
});
if (action === "start") expect(res.ok()).toBeTruthy();
diff --git a/editor/e2e/rate-limit.spec.ts b/editor/e2e/rate-limit.spec.ts
@@ -146,6 +146,8 @@ async function control(
action: "start" | "stop",
) {
const res = await request.post(`${baseUrl}/api/auto-queue/control`, {
+ // Behind the ops token since release 19 (A3): it starts and stops lanes.
+ headers: { authorization: "Bearer test-worker-token" },
data: { kind: "download", action },
});
if (action === "start") expect(res.ok()).toBeTruthy();
diff --git a/plans/release-19.md b/plans/release-19.md
@@ -117,6 +117,160 @@ now: C1, C3 ──► C2 ──► C4, C5; C2 regenerated after A's actions la
### Track A
+Branch `worktree-agent-a103bef2ee4055362` (worktree `.claude/worktrees/agent-a103bef2ee4055362`, ports editor 4101,
+test 4111, export 4110) off `4cffda3f`, one Opus implementer, A1–A5 in order, one commit each; `r19/integration`
+merged in before this record. Scratch files `a-*` in the job's `tmp`.
+
+| commit | slice | one line |
+|---|---|---|
+| `988b282c` | A1 | `pnpm ops` reads `WORKER_TOKEN` / `ARCHILYZER_EDITOR_URL` from `editor/.env(.local)`; `--wait` asks `GET /api/ops/job/<id>` |
+| `5f86b10a` | A2 | `GET /api/ops/jobs`, `POST /api/ops/job {verb, ids}`; `get job`, `get jobs`, `job <verb>`, `job wait` |
+| `e065bd2d` | A2b | 19 of `ops-api.spec`'s request/response tests become route unit tests |
+| `30111658` | A3 | the read side (`get settings|storage|sites|workers|auto-queue|scheduler|cleanup`); `/api/auto-queue/control` behind the token |
+| `5f442885` | A4 | archival writes (settings patch, the publish lane, platform holds, workers, one video, cleanup, relocate dryRun); `archilyzer storage report` |
+| `c261b64a` | A5 | clip windows on `clips:<platform>`; an in-flight window answers its job |
+| `e53f3169` | — | `r19/integration` merged in |
+| `cde0e759` | A2b | the test corpus helper's paths opt out of tracing (`next-build-trace.test.mjs`) |
+
+#### Slice A1, as shipped — `pnpm ops` finds its token
+
+- `scripts/archilyzer-ops.mjs` `loadEditorEnv()` fills `WORKER_TOKEN` and `ARCHILYZER_EDITOR_URL` when unset — a
+ variable already set, even to `""`, wins — from `editor/.env.local` then `editor/.env` of the checkout the SCRIPT is
+ in (`REPO_ROOT` from `import.meta.url`, never the cwd), then of the main worktree when it runs in a linked one
+ (`mainWorktreeOf`: `.git` file → `gitdir` → `commondir`, read off disk, no git process). Only those two keys are
+ taken; the files also hold deploy credentials. A 401/503 prints `tokenHint`: where the token came from, never its
+ value. `parseDotenv` reads `KEY=value`, `export`, quotes, comments.
+- **What `--wait` polled and why it was a rewrite.** After three failed log polls it asked `/api/jobs/active`, which
+ since one-core phase 3 slice 2 is a `next.config` rewrite onto `/api/view/activeJobs` — kept at its old path for the
+ pages and pinned widgets that poll it — and that view builds the whole live payload (channel stats, the disk gate,
+ every runner's status) to answer "is job X still there". It now asks `GET /api/ops/job/<id>` (behind the token: one
+ registry record or `.meta.json` sidecar). An `archived` log status (an id the registry forgot) resolves to the
+ sidecar's terminal status instead of exiting 1 for a job that ended `done`; a 404 `{ok:false}` is "gone". While a
+ job waits, `followJob` prints `queued — position N on <queue>` once per change.
+- `GET /api/ops/job/<id>[?tail=N]` (`editor/app/api/ops/_jobs.ts`): the `getJobEntry` row without `logPath`, `queue:
+ {key, position, queued, head}` while queued or running, the runner's `progress`; `tail` (≤ 500) is
+ `common/jobs/listJobs.ts` `readLogTail` — the last N lines read from the end (256 KiB window, a partial first line
+ dropped).
+
+#### Slice A2, as shipped — jobs over ops
+
+- `GET /api/ops/jobs` — `?active` (registry queued/running, queue order, each with its place), `?failed`, `?kind`,
+ `?slug`, `?limit` (50, ≤ 500); a filtered history list reads the newest 2000 jobs (`JOBS_SCAN`); `active` with
+ `failed` is a 400; unknown query keys are a 400.
+- `POST /api/ops/job {verb, ids}` — `cancel | drain | promote | force-release | retry` per id through
+ `editor/app/jobs/actions.ts` (the row buttons' own actions), each answered in `results`; one refused id makes the
+ answer `ok: false` (400, naming it) without stopping the others. `retry-failed` (no ids) is Retry all;
+ `retryAllFailedAction` now also returns `jobIds` and cancels the new jobs' unread streams. `jobIds` = the jobs a
+ retry started, which `--wait` follows.
+- CLI: `get job <id> [--tail [N]]`, `get jobs [--active|--failed] [--kind] [--slug] [--limit]`, `job <verb> <id…>`,
+ `job wait <id…>` (no request; `followAll`, the `--wait` loop factored out). A read's flag on the wrong noun or on an
+ action is refused, not ignored.
+- `ops-api.spec` gains "jobs over ops": two `/api/test/stuck-job` holders on one queue — the live list, one job's
+ place and head, `?tail`, a promote refused ("already at the front"), a cancel of one known and one unknown id,
+ `force-release` of the head. The success paths revalidate `/jobs`, which only a running server can.
+
+#### Slice A2b, as shipped — the ops-api spec's request/response tests are unit tests
+
+- 19 of `ops-api.spec`'s 33 tests (32 at base + A2's) exercised only a route's request and response; they moved,
+ assertion for assertion, to `editor/app/api/ops/_door.test.ts` (the 401, the 503, unknown keys, the traversing slug
+ on fourteen routes) and route tests for retry-bucket, tag-videos, persist-videos, build-site (+ build-deploy),
+ deploy-site (+ build-deploy), deploy-hub, deploy-homepage and publish; fetch-posts and capture-posts gain one each.
+ **ops-api.spec: 33 → 14 tests.** What stays needs the running server: writes that revalidate, jobs that run, the
+ publish queue's runs, the jobs verbs.
+- `editor/app/api/ops/_testCorpus.ts` (imported only by tests): copies the named e2e fixture into a temp dir, sets the
+ e2e test server's environment (fake yt-dlp, gallery-dl, ffmpeg/ffprobe, whisper, chough, parakeet, diarize, claude,
+ findmnt, udisksctl, wrangler, next; the fake Cloudflare token; `ARCHILYZER_BRANCH=main`; `E2E_LIVE_CHECK=skip`) and
+ `holdPublishQueue()`s as e2e did, so a refusal that regresses can only queue a fake.
+- The 503 no longer needs `/api/test/worker-token`: the unit process unsets its own `WORKER_TOKEN`. That test route now
+ has no e2e caller (left in place).
+
+#### Slice A3, as shipped — the read side; the lane control route needs the token
+
+- `GET /api/ops/settings[?key]` (getSettings, migrated and defaulted), `storage` (`buildStorage`), `sites`
+ (`listSites`: id, title, `siteUrl`, Pages project, audience, `isListedSite`, search, publish policy, channels),
+ `workers` (`buildWorkersPayload`), `auto-queue` (`buildAutoQueueStatusPayload`), `scheduler`
+ (`buildSchedulerStatusPayload`), `cleanup/<slug>` (new `loadCleanupRow` — `loadCleanupSummary`'s row for one channel,
+ id lists as counts). One door, `editor/app/api/ops/_read.ts` `readRoute` (token, unknown query keys refused) and
+ `redactSecrets` (any string under a key naming a token/secret/password/api key/credential → `"<redacted>"`; a remote
+ worker's `token` is the one that exists today). CLI `get settings [<key>]`, `get storage|sites|workers|auto-queue|
+ scheduler`, `get cleanup <slug>`.
+- **`/api/auto-queue/control` answers `opsAuth`.** Its one UI caller — `RunnerOperationView`'s Start/Drain/Stop on
+ /operations — calls `startAutoQueueAction` / `stopAutoQueueAction` / new `drainAutoQueueAction` instead: `opsAuth`
+ has no exemption for the UI (the UI never called an ops route; it uses server actions, which carry Next's origin
+ check), and a page holds no token. The six e2e specs that post to the route (auto-queue, auto-subs-replace,
+ disk-space, lane-runner, pacing, rate-limit) send the test token.
+
+#### Slice A4, as shipped — archival writes
+
+- `POST /api/ops/settings {patch}` → `saveSettings` (one-level merge). Refused before anything is written
+ (`settings/patch.ts`): an unknown top-level key; **a key another writer owns** — `channelPriority` (compiled with the
+ lanes' trees in one write), `autoQueue`, `workers`, `storage` — naming that writer; and a value the schema would not
+ keep as sent (it coerces, never throws), every leaf of the patch compared with the schema's reading of the merged
+ settings and each difference named ("minFreeDiskGB: sent -5, would be saved as <the clamp>"). Answers `{changed, value}` read
+ back, redacted.
+- `lane` takes `publish`: `enabled` through `savePublishSettingsAction` with the stored fields, `held` through the
+ pause gate, Start/Drain/Stop the publish page's actions; the ingest lanes' drain goes through
+ `drainAutoQueueAction`.
+- `clear-platform-hold {platform}` (`clearPlatformHoldAction`, answers its sentence); `POST /api/ops/workers {op:
+ enable|disable|drain, ids}` (the live pool, per id); `transcribe-one {slug, id, file?}` (`transcribeOneAction` with a
+ file, else the video page's `whisperVideoAction`, which fetches audio through the managed path when there is none;
+ the file name is checked as one path segment); `delete-file {slug, id, file}` (`deleteVideoFileAction` —
+ `removeMediaFile`, refused while the media drive is not answering); `do-not-clean {slug, id, keep?}` (set, not
+ toggled); `POST /api/ops/cleanup {slug, sweep: transcribed|extra-formats|wrong-format|failed-transcriptions}` (the
+ channel page's four actions, as jobs); `relocate {dryRun}` (`previewRelocationAction` per slug — which, as on the
+ panel, tiers a never-tiered channel in place first; a slug with no channel is the bulk path's "channel not found").
+- `archilyzer storage report [<slug>…] [--json]` (`common/bin/storage-report.ts`): per channel the tierable media,
+ text and clip bytes off its report (`totalMediaBytes` is `isTierable`'s share) and where its media is — corpus disk,
+ a location (`config.mediaDir` under its root, longest match), a custom path, or legacy — with totals per place; a
+ missing figure is unknown (`—` / `null`), never 0; social channels left out; no drive touched, no editor needed.
+
+#### Slice A5, as shipped — fetch queue hygiene
+
+- **Its own queue.** `common/lib/queueKeys.ts` `clipWindowQueueKey`: a window job runs on `clips:<platform>` (from the
+ platform queue it would have used, so `platform:vimeo.com` → `clips:vimeo.com`) — single (`fetchWindowAction`) and
+ batch (`fetchWindowsAction`). A queue runs one job at a time and a tier orders only the waiting ones, so priority
+ alone could not get a window past a running multi-hour persist; the queue key can. **Ruling recorded:** one window
+ may now run beside one download on the same platform; windows stay one at a time per platform with the batches'
+ `clip-window:<platform>` gap, and the platform's hold, cooldown and 429 backoff are shared with the downloads.
+- **Deduped.** `common/jobs/windowJobs.ts` `windowInFlight`: 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 — read off the job's
+ replay spec. `/api/media/fetch-window` answers that job (202, `existing: true`, its window and file); a batch lists
+ such windows in `inFlight` and queues nothing for them, and the ops route adds their jobs to `jobIds` so `--wait`
+ follows the windows asked for. Retry (`stream: true`) is not deduplicated.
+- **Not done: `fetch-via-editor.mjs`.** It is `umtool/report-to-video/fetch-via-editor.mjs` — under the directory no
+ track touches — so its 10-minute poll timeout and its wait output are unchanged. The editor half covers its retry
+ loop (a re-run after its timeout gets the job already fetching), and `pnpm ops job wait <id>` waits with no timeout,
+ printing the queue position.
+
+#### Track A gates (A1–A5, on the merged tip)
+
+- `pnpm -r --no-bail --workspace-concurrency=1 exec tsc --noEmit` — clean before every commit and on the merged tip
+ (the A2/A2b/A3 commits under the 12:35 throttle ran the editor package's tsc alone; A1, A4, A5 and the tip ran the
+ whole workspace).
+- **editor unit** 205/205 (159 at base: +5 A1, +4 A2, +19 A2b, +7 A3, +9 A4, +2 A5). **test:scripts** 693: 689
+ pass, 3 skip, 1 fail — `next-build-trace.test.mjs`'s cwd-path check on the new test helper, fixed by `cde0e759`
+ (the rerun's one failure is `queue-lock.test.mjs` "prints a banner naming the holder", run while another session's
+ e2e held the machine lock; it passed in the gate run). `archilyzer-ops.test.mjs` 46 → 59. **common** 3473: 3469
+ pass, 4 fail — `fetchPosts.test.ts` and `relocateDir.test.ts` pass alone (load); `autoRunner.test.ts`'s two
+ ("computeLeafPending takes its channel listing from `shared`", "a channel with a move marker projects no work…")
+ fail the same way on an extracted `4cffda3f` tree: not this branch's. +8 here (listJobs 1, windowJobs 3,
+ storage-report 4). **mcp** 292/292.
+- **`pnpm --filter editor exec next build`** (5 GB scope) — ok.
+- **e2e** (detached, queued, memory-gated; `export/public` linked to the composed fixture as `r18-integration`'s is —
+ the primary's no longer has `summaries/`, so the export server never answered and the first launch timed out
+ waiting for it): `ops-api jobs fetch-window auto-queue lane-runner disk-space pacing rate-limit auto-subs-replace` —
+ **65 passed, 5 failed, 19.0 min**, at load average ~26–32. ops-api (14/14), fetch-window (9/9), lane-runner,
+ pacing and auto-subs-replace passed. The five failures are 30-s test timeouts before any code this branch changed ran: auto-queue's "UI:
+ build a policy…" (the Start click its own comment calls a race against the poll — it lost, the button went disabled
+ under it), "Drain completes when … parked" (timed out polling for the manual whisper job, before the Drain click),
+ disk-space's "prevents a download…" (the channel page's Download click), jobs' "kicks off the index update" (the
+ index stage still running at 30 s), rate-limit's "an overdue hold…" (`ECONNRESET` from the dev server on a control
+ POST; the same POST with the token passed in the other specs). **Recheck of the five alone** (load ~22, 2.0 min):
+ **4 passed, 1 failed** — the UI Start and the UI Drain (the two buttons A3 moved to server actions), disk-space and
+ rate-limit pass; jobs' index test hit its 30-s test timeout again while its own page snapshot shows the stage's log
+ ending `[stage] update-index _index: Done` — the index stage, which this branch does not touch, finishing late
+ under load.
+
### Track B
Branch `worktree-agent-a8b654c51bf472562` off `4cffda3f` (r18/integration with main merged), one Opus implementer,
diff --git a/scripts/archilyzer-ops.mjs b/scripts/archilyzer-ops.mjs
@@ -13,6 +13,10 @@
// pnpm ops get channel <slug> [--counts]
// pnpm ops get channels
// pnpm ops get tags [<tagId>]
+// pnpm ops get job <id> [--tail [<lines>]]
+// pnpm ops get jobs [--active | --failed] [--kind <kind>] [--slug <slug>] [--limit <n>]
+// pnpm ops job <cancel|drain|promote|force-release|retry> <id>... [--wait]
+// pnpm ops job retry-failed | job wait <id>...
// pnpm ops list
//
// ARCHILYZER_EDITOR_URL editor base URL (default http://localhost:3001)
@@ -23,6 +27,12 @@
// Unset on the SERVER => every route 503s; unset here
// => every route 401s.
//
+// Either variable, when unset, is read from the editor's own env files —
+// `editor/.env.local`, then `editor/.env` — of THIS checkout (found from the
+// script's path, never the cwd), then of the main worktree when this is a
+// linked one (a worktree has no `editor/.env`; the editor it talks to by
+// default is the primary's). See loadEditorEnv.
+//
// EXAMPLES
//
// pnpm ops sync --json '{"slug":"the-quartering"}' --wait
@@ -41,6 +51,13 @@
// pnpm ops lane --json '{"lane":"download","held":true}'
// pnpm ops refresh-report --json '{"all":true}'
// pnpm ops relocate --json '{"slugs":["x"],"locationId":"platter"}'
+// pnpm ops relocate --json '{"slugs":["x"],"locationId":"platter","dryRun":true}'
+// pnpm ops settings --json '{"patch":{"minFreeDiskGB":20}}'
+// pnpm ops lane --json '{"lane":"publish","action":"drain"}'
+// pnpm ops cleanup --json '{"slug":"x","sweep":"transcribed"}' --wait
+// pnpm ops workers --json '{"op":"disable","ids":["parakeet-cpu"]}'
+// pnpm ops get job <jobId> --tail 20
+// pnpm ops get jobs --failed --slug x
// pnpm ops build-site --json '{"siteId":"anilyzer"}' --wait
// pnpm ops build-deploy --json '{"siteIds":["anilyzer","jeralyzer"]}' --wait
// pnpm ops deploy-site --json '{"siteId":"anilyzer","preview":"tags-exclude"}' --wait
@@ -79,12 +96,14 @@
// The response JSON is printed verbatim on stdout (log lines from --wait go to
// stderr), so `pnpm ops … | jq` works.
+import { readFileSync } from "node:fs";
import { readFile } from "node:fs/promises";
-import { pathToFileURL } from "node:url";
+import path from "node:path";
+import { fileURLToPath, pathToFileURL } from "node:url";
const DEFAULT_URL = "http://localhost:3001";
-// Consecutive polls where NEITHER the job's log NOR the active list answered,
+// Consecutive polls where NEITHER the job's log NOR its ops record answered,
// after which --wait gives up. At the 30 s backoff ceiling that is ~5 minutes
// of an editor saying nothing at all, which is not a busy server — it is a
// server that is gone.
@@ -96,7 +115,7 @@ const MAX_PROBE_FAILURES = 10;
const GETTERS = {
// --counts adds the LIVE on-disk counts, which walk every video directory —
// opt-in for the same reason the route makes it opt-in.
- channel: (slug, counts) =>
+ channel: (slug, { counts }) =>
`/api/ops/channel/${encodeURIComponent(slug)}${counts ? "?counts=1" : ""}`,
// Every channel, one line each: slug, name, kind, platform and the sites
// that carry it ([] = private to this editor).
@@ -108,11 +127,75 @@ const GETTERS = {
// The publish status (release 18): the index, the lane, a row per site, the
// hub and the homepage with their chips, and the plan Publish now would run.
publish: () => "/api/ops/publish",
+ // One job: its record or sidecar, where it waits, and with --tail [N] the
+ // last N lines of its log (default 40).
+ job: (id, { tail }) =>
+ `/api/ops/job/${encodeURIComponent(id)}${tail ? `?tail=${tail}` : ""}`,
+ // The /jobs list: --active (the live head, in queue order), --failed,
+ // --kind <kind>, --slug <slug>, --limit <n>.
+ jobs: (_arg, { jobs }) => {
+ const q = new URLSearchParams();
+ if (jobs.active) q.set("active", "1");
+ if (jobs.failed) q.set("failed", "1");
+ if (jobs.kind) q.set("kind", jobs.kind);
+ if (jobs.slug) q.set("slug", jobs.slug);
+ if (jobs.limit) q.set("limit", String(jobs.limit));
+ const s = q.toString();
+ return `/api/ops/jobs${s ? `?${s}` : ""}`;
+ },
+ // THE READ SIDE (release 19, A3): what each page draws, behind the token.
+ // settings.json as the editor reads it (migrated, defaulted, secrets
+ // redacted), or one top-level block of it.
+ settings: (key) =>
+ key ? `/api/ops/settings?key=${encodeURIComponent(key)}` : "/api/ops/settings",
+ // The /storage page: each location, mounted or not, its free space and tiers.
+ storage: () => "/api/ops/storage",
+ // Every site: id, title, siteUrl, Pages project, audience, policy, channels.
+ sites: () => "/api/ops/sites",
+ // The transcription workers and what each is running.
+ workers: () => "/api/ops/workers",
+ // The four lanes: runner, policy, cooldowns, picks, pending.
+ "auto-queue": () => "/api/ops/auto-queue",
+ // The sync scheduler: each channel's schedule, the tick log.
+ scheduler: () => "/api/ops/scheduler",
+ // One channel's cleanup row: what each sweep would reclaim, what holds the rest.
+ cleanup: (slug) => `/api/ops/cleanup/${encodeURIComponent(slug)}`,
};
// Nouns whose read takes no argument. `get channel` without a slug is a
// mistake; `get tags` without one is the whole vocabulary.
-const GET_ARG_OPTIONAL = new Set(["tags", "channels", "publish"]);
+const GET_ARG_OPTIONAL = new Set([
+ "tags",
+ "channels",
+ "publish",
+ "jobs",
+ "settings",
+ "storage",
+ "sites",
+ "workers",
+ "auto-queue",
+ "scheduler",
+]);
+
+// The flags each noun takes beyond the shared ones; any other is refused.
+const GET_FLAGS = {
+ channel: ["counts"],
+ job: ["tail"],
+ jobs: ["active", "failed", "kind", "slug", "limit"],
+};
+
+// The row buttons of /jobs (`POST /api/ops/job`), as `job <verb> <id…>`.
+// `retry-failed` takes no ids; `wait` is not a request at all — it follows the
+// ids as --wait would.
+export const JOB_VERBS = [
+ "cancel",
+ "drain",
+ "promote",
+ "force-release",
+ "retry",
+ "retry-failed",
+ "wait",
+];
const ACTIONS = [
"channel-priority",
@@ -197,6 +280,19 @@ const ACTIONS = [
// ({path, start?, end?, workerId?, out?}) — the corpus's own engine and
// model, as a job. With --wait the result JSON is what stdout carries.
"transcribe",
+ // ARCHIVAL WRITES (release 19, A4). settings.json through the editor's one
+ // writer ({patch}); a platform's rate-limit hold cleared ({platform});
+ // workers switched on, off or drained ({op, ids}); one video transcribed
+ // ({slug, id, file?}), one of its files deleted ({slug, id, file}), its
+ // do-not-clean marker set ({slug, id, keep?}); a channel's cleanup sweep
+ // ({slug, sweep}).
+ "settings",
+ "clear-platform-hold",
+ "workers",
+ "transcribe-one",
+ "delete-file",
+ "do-not-clean",
+ "cleanup",
];
// The log line a transcribe job ends with: this marker, then the result as
@@ -216,6 +312,11 @@ export function parseArgs(argv) {
let wait = false;
let quiet = false;
let counts = false;
+ let tail = 0;
+ const jobs = {};
+ // The noun-specific flags seen, checked against GET_FLAGS once the noun is
+ // known: a flag a read does not take is refused, never ignored.
+ const used = new Set();
let waitTimeout = null;
const readTimeout = (raw) => {
const n = Number(raw);
@@ -246,6 +347,39 @@ export function parseArgs(argv) {
quiet = true;
} else if (arg === "--counts") {
counts = true;
+ used.add("counts");
+ } else if (arg === "--active" || arg === "--failed") {
+ jobs[arg.slice(2)] = true;
+ used.add(arg.slice(2));
+ } else if (/^--(kind|slug|limit)(=|$)/.test(arg)) {
+ const [, key, eq] = /^--(kind|slug|limit)(=?)/.exec(arg);
+ const raw = eq ? arg.slice(key.length + 3) : argv[++i];
+ if (raw === undefined || raw === "") {
+ return { error: `--${key} needs a value` };
+ }
+ if (key === "limit") {
+ const n = Number(raw);
+ if (!Number.isInteger(n) || n <= 0) {
+ return { error: "--limit needs a whole number above zero" };
+ }
+ jobs.limit = n;
+ } else {
+ jobs[key] = raw;
+ }
+ used.add(key);
+ } else if (arg === "--tail" || arg.startsWith("--tail=")) {
+ // An optional count: `--tail` alone is the last 40 lines.
+ const raw = arg.startsWith("--tail=")
+ ? arg.slice("--tail=".length)
+ : /^\d+$/.test(argv[i + 1] ?? "")
+ ? argv[++i]
+ : "40";
+ const n = Number(raw);
+ if (!Number.isInteger(n) || n <= 0) {
+ return { error: "--tail takes a whole number of lines above zero" };
+ }
+ tail = n;
+ used.add("tail");
} else if (arg === "--json") {
json = argv[++i];
if (json === undefined) {
@@ -296,14 +430,55 @@ export function parseArgs(argv) {
if (!positional[2] && !GET_ARG_OPTIONAL.has(noun)) {
return { error: `get ${noun}: needs an argument` };
}
+ const allowed = GET_FLAGS[noun] ?? [];
+ const stray = [...used].filter((f) => !allowed.includes(f));
+ if (stray.length) {
+ return {
+ error: `get ${noun} does not take ${stray.map((f) => `--${f}`).join(", ")}`,
+ };
+ }
return {
method: "GET",
- path: GETTERS[noun](positional[2], counts),
+ path: GETTERS[noun](positional[2], { counts, tail, jobs }),
wait: false,
quiet,
waitTimeout,
};
}
+ if (used.size) {
+ return {
+ error: `${[...used].map((f) => `--${f}`).join(", ")} belongs to a read (get …), not to "${positional[0]}"`,
+ };
+ }
+ if (positional[0] === "job") {
+ const verb = positional[1];
+ if (!JOB_VERBS.includes(verb)) {
+ return {
+ error: `job: unknown verb "${verb ?? ""}" — known: ${JOB_VERBS.join(", ")}`,
+ };
+ }
+ const ids = positional.slice(2);
+ if (verb === "retry-failed") {
+ if (ids.length) return { error: "job retry-failed takes no ids" };
+ } else if (ids.length === 0) {
+ return { error: `job ${verb}: needs one or more job ids` };
+ }
+ if (json !== null || file !== null) {
+ return { error: `job ${verb} takes its ids as arguments, not a body` };
+ }
+ if (verb === "wait") {
+ // Nothing to send: follow each id to its end, as --wait would.
+ return { waitFor: ids, wait: true, quiet, waitTimeout };
+ }
+ return {
+ method: "POST",
+ path: "/api/ops/job",
+ body: verb === "retry-failed" ? { verb } : { verb, ids },
+ wait,
+ quiet,
+ waitTimeout,
+ };
+ }
const action = positional[0];
if (!ACTIONS.includes(action)) {
return {
@@ -402,16 +577,53 @@ export function usage() {
" pnpm ops get channels",
" pnpm ops get tags [<tagId>]",
" pnpm ops get publish",
+ " pnpm ops get job <id> [--tail [<lines>]]",
+ " pnpm ops get jobs [--active | --failed] [--kind <kind>] [--slug <slug>]",
+ " [--limit <n>]",
+ " pnpm ops job cancel|drain|promote|force-release|retry <id>... [--wait]",
+ " pnpm ops job retry-failed [--wait]",
+ " pnpm ops job wait <id>... [--wait-timeout <seconds>]",
+ " pnpm ops get settings [<key>]",
+ " pnpm ops get storage | sites | workers | auto-queue | scheduler",
+ " pnpm ops get cleanup <slug>",
" pnpm ops list",
"",
`Actions: ${ACTIONS.join(", ")}`,
"",
"--wait follows the job's log and survives a poll that fails (a busy",
" in-process build starves the server): it backs off and, after three",
- " failures, asks /api/jobs/active whether the job is still there.",
+ " failures, asks /api/ops/job/<id> whether the job is still there and how",
+ " it ended. While the job waits it prints its queue position on stderr.",
"--wait-timeout <seconds> gives up and exits 1 instead of waiting forever.",
" Default: no timeout — the queue may legitimately hold a job for hours.",
"",
+ "get job <id> is one job: kind, channel, status, times, exit code, and while",
+ " it waits its queue (key, position, how many wait, the job at the head);",
+ " --tail [<lines>] adds the last lines of its log (40, at most 500).",
+ "get jobs lists jobs newest first (50; --limit, at most 500): --active is",
+ " every queued and running job in queue order, --failed the failed ones;",
+ " --kind and --slug narrow either. A filtered list looks through the",
+ " newest 2000 jobs.",
+ "job <verb> <id>... is a /jobs row's button: cancel (a running job is",
+ " stopped), drain (finish the sub-operations in flight, start no more),",
+ " promote (a queued job to the front of its queue), force-release (free a",
+ " wedged queue slot), retry (re-run from its replay descriptor, ahead of",
+ " the queue). retry-failed retries every failed job the editor still holds.",
+ " Each id is answered in \"results\"; one that could not be acted on makes",
+ " the answer ok: false without stopping the rest. --wait follows the jobs a",
+ " retry started. job wait <id>... follows jobs already running.",
+ "",
+ "get settings [<key>] is settings.json as the editor reads it (migrated and",
+ " defaulted; a secret is \"<redacted>\"), or one top-level key of it.",
+ "get storage is /storage: each location, mounted or not, free space, tiers.",
+ "get sites lists every site: id, title, siteUrl, Pages project, audience,",
+ " listed, search, publish policy, channels. (`get publish` is the status.)",
+ "get workers, get auto-queue and get scheduler are /workers, the four lanes",
+ " of /operations, and /operations/sync — the payloads those pages poll.",
+ "get cleanup <slug> is one channel's /cleanup row: what each sweep would",
+ " reclaim (they overlap — never add them), what holds the rest, and the",
+ " failed-transcriptions count. measured: false means unknown, not zero.",
+ "",
'publish runs publish stages on the editor\'s publish queue, one at a time,',
' under one run id: {"verb": …}. "index" updates the index; "build" builds',
' "siteId"/"siteIds" (forced; the index first when stale; "runner":',
@@ -440,6 +652,48 @@ export function usage() {
' report.md and evidence-pack.zip for the site\'s build to publish:',
' {"siteId", "reportId"?, "formats"?: ["html","pdf","md","zip"]}.',
"",
+ 'settings patches settings.json through the editor\'s one writer:',
+ ' {"patch": {"<top-level key>": value, …}} (SETTINGS.md). A block\'s keys',
+ ' merge one level; an array, a scalar or an object nested in a block',
+ ' replaces. Refused before anything is written: an unknown key; a key',
+ ' another writer owns (channelPriority, autoQueue, workers, storage — the',
+ ' refusal names it); a value the schema would not keep as sent (clamped,',
+ ' dropped), named with what would be saved. Answers the saved values.',
+ "",
+ 'lane takes {"lane": "transcription"|"download"|"digest"|"backfill"|',
+ ' "publish", "enabled"?, "held"?, "action"?: "start"|"stop"|"drain"}.',
+ "",
+ 'clear-platform-hold clears a platform\'s rate-limit hold, backoff and raised',
+ ' pace, as the lane strip\'s Clear hold: {"platform"}. Answers the sentence.',
+ "",
+ 'workers switches transcription workers as /workers does: {"op": "enable" |',
+ ' "disable" | "drain", "ids": [...]} — the live pool; a restart restores',
+ ' the launch default. An unknown id is named; the others still apply.',
+ "",
+ 'transcribe-one transcribes one video: {"slug", "id", "file"?}. With "file"',
+ ' (an audio file in its dir) that file; without, the video page\'s',
+ ' Transcribe — the audio on disk, or downloaded first through the channel\'s',
+ ' managed path. A job; --wait follows it.',
+ "",
+ 'delete-file deletes one file of one video, as its Files list does:',
+ ' {"slug", "id", "file"}. A tiered media file goes with its bytes on the',
+ ' media drive (never a bare link); refused while that drive is not',
+ ' answering. No undo.',
+ "",
+ 'do-not-clean sets one video\'s "Do not clean" marker: {"slug", "id",',
+ ' "keep"?: true}; "keep": false removes it. (keep-videos is the bulk form.)',
+ "",
+ 'cleanup runs one of a channel\'s cleanup sweeps as a job: {"slug", "sweep":',
+ ' "transcribed" (audio of transcribed videos — the /cleanup total) |',
+ ' "extra-formats" | "wrong-format" | "failed-transcriptions" (clears the',
+ ' list; deletes nothing)}. Do-not-clean videos are skipped; `get cleanup',
+ ' <slug>` says what each would reclaim.',
+ "",
+ 'relocate moves channels\' media to a location: {"slugs", "locationId" |',
+ ' "root"}. "dryRun": true answers each channel\'s preview (bytes to copy,',
+ ' free space both sides) and moves nothing — though, as on the Storage',
+ ' panel, a channel never tiered is first tiered in place on the corpus disk.',
+ "",
'channel-config changes a channel as its Configure form does: {"slug"} and',
' any of "patch" (form field names; "" clears one), "sites" (the WHOLE',
' membership set: [{"siteId", "groupId"? | "newGroupName"?}], [] = on no',
@@ -521,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",
@@ -561,12 +819,103 @@ export function usage() {
" stderr), so `pnpm ops transcribe ... --wait | jq -r .text` works.",
"",
"Env: ARCHILYZER_EDITOR_URL (default http://localhost:3001), WORKER_TOKEN,",
- " ARCHILYZER_AGENT (provenance of a tag write; default \"cli\")",
+ " ARCHILYZER_AGENT (provenance of a tag write; default \"cli\").",
+ " WORKER_TOKEN and ARCHILYZER_EDITOR_URL, when unset, are read from",
+ " editor/.env.local and editor/.env of this checkout (found from the",
+ " script, not the cwd), then of the main worktree.",
].join("\n");
}
+// THE REPO ROOT, FROM THIS FILE — never from the cwd. `pnpm ops` runs with the
+// repo root as its cwd, but `node scripts/archilyzer-ops.mjs` from a subdir,
+// or a tool that shells out from its own project directory, does not; every
+// session used to write a wrapper that sourced editor/.env first.
+const REPO_ROOT = path.resolve(path.dirname(fileURLToPath(import.meta.url)), "..");
+
+// The two variables this client reads from the editor's env files. Nothing
+// else is taken from them: they also hold deploy credentials, which have no
+// business in this process.
+const EDITOR_ENV_KEYS = ["WORKER_TOKEN", "ARCHILYZER_EDITOR_URL"];
+
+// `KEY=value` lines, as Next's own loader reads a .env: `#` comments, an
+// optional `export `, single or double quotes stripped. No interpolation.
+export function parseDotenv(text) {
+ const out = {};
+ for (const raw of text.split(/\r?\n/)) {
+ const line = raw.trim();
+ if (!line || line.startsWith("#")) continue;
+ const m = /^(?:export\s+)?([A-Za-z_][A-Za-z0-9_]*)\s*=\s*(.*)$/.exec(line);
+ if (!m) continue;
+ let value = m[2].trim();
+ const q = value[0];
+ if ((q === '"' || q === "'") && value.length >= 2 && value.endsWith(q)) {
+ value = value.slice(1, -1);
+ } else {
+ // An unquoted value ends at an inline comment.
+ value = value.replace(/\s+#.*$/, "");
+ }
+ out[m[1]] = value;
+ }
+ return out;
+}
+
+// The main worktree of a LINKED worktree, or null. A linked worktree's `.git`
+// is a file `gitdir: <main>/.git/worktrees/<name>`, and that directory's
+// `commondir` names the main `.git`. Read off disk — no git process.
+export function mainWorktreeOf(root, read = (p) => readFileSync(p, "utf8")) {
+ try {
+ const m = /^gitdir:\s*(.+)$/m.exec(read(path.join(root, ".git")));
+ if (!m) return null;
+ const gitdir = path.resolve(root, m[1].trim());
+ const common = path.resolve(gitdir, read(path.join(gitdir, "commondir")).trim());
+ return path.dirname(common);
+ } catch {
+ return null;
+ }
+}
+
+// The env files consulted, in order: this checkout's, then the main worktree's.
+export function editorEnvFiles(root = REPO_ROOT, read) {
+ const roots = [root];
+ const main = mainWorktreeOf(root, read);
+ if (main && path.resolve(main) !== path.resolve(root)) roots.push(main);
+ return roots.flatMap((r) => [
+ path.join(r, "editor", ".env.local"),
+ path.join(r, "editor", ".env"),
+ ]);
+}
+
+// FILL WHAT THE ENVIRONMENT LEFT UNSET from the editor's env files. A variable
+// already set — even to "" — wins: that is an explicit choice. The first file
+// that names a key supplies it. Returns where each came from, for the 401/503
+// hint (never the value).
+export function loadEditorEnv(
+ env = process.env,
+ files = editorEnvFiles(),
+ read = (p) => readFileSync(p, "utf8"),
+) {
+ const sources = {};
+ for (const key of EDITOR_ENV_KEYS) {
+ if (env[key] !== undefined) sources[key] = "environment";
+ }
+ for (const file of files) {
+ let parsed;
+ try {
+ parsed = parseDotenv(read(file));
+ } catch {
+ continue;
+ }
+ for (const key of EDITOR_ENV_KEYS) {
+ if (sources[key] || parsed[key] === undefined) continue;
+ env[key] = parsed[key];
+ sources[key] = file;
+ }
+ }
+ return sources;
+}
+
function baseUrl() {
- return (process.env.ARCHILYZER_EDITOR_URL ?? DEFAULT_URL).replace(/\/+$/, "");
+ return (process.env.ARCHILYZER_EDITOR_URL || DEFAULT_URL).replace(/\/+$/, "");
}
function authHeaders() {
@@ -574,15 +923,30 @@ function authHeaders() {
return token ? { authorization: `Bearer ${token}` } : {};
}
+// What a 401 or a 503 means HERE, with where the token came from — never what
+// it is. A 503 is the server's own switch (it runs without WORKER_TOKEN); a 401
+// is this side sending none, or one the server does not hold.
+export function tokenHint(status, sources, url = baseUrl()) {
+ const from = sources.WORKER_TOKEN;
+ if (status === 503) {
+ return `hint: the editor at ${url} runs without WORKER_TOKEN, so /api/ops is off — set it in its editor/.env and restart it`;
+ }
+ if (!from) {
+ return `hint: no WORKER_TOKEN in the environment or in ${editorEnvFiles().join(", ")}`;
+ }
+ return `hint: the WORKER_TOKEN from ${from} is not the one the editor at ${url} runs with`;
+}
+
// Follow a job's log to its terminal state. Returns the status string.
// Deliberately polls the SAME endpoint the editor's own log panel does, so a
// job started here and a job started by a click are observed identically.
-// NO AUTH HEADER, and that is not an omission. `/api/jobs/<id>/log` is
-// deliberately ungated (see its route for what it does and does not leak, and
-// why gating it would switch every run panel's log off on a default install) —
-// the browser polls it same-origin with no token. Sending one here implied a
-// gate that does not exist, which is worse than sending nothing: the next
-// person to read this would conclude the endpoint was protected.
+// NO AUTH HEADER on the log poll, and that is not an omission.
+// `/api/jobs/<id>/log` is deliberately ungated (see its route for what it does
+// and does not leak, and why gating it would switch every run panel's log off
+// on a default install) — the browser polls it same-origin with no token.
+// Sending one here implied a gate that does not exist, which is worse than
+// sending nothing: the next person to read this would conclude the endpoint
+// was protected.
//
// A POLL FAILURE IS NOT A JOB FAILURE, and this used to treat them as the same
// thing. The editor is single-process: a busy in-process build-index starves
@@ -594,10 +958,23 @@ function authHeaders() {
//
// After three consecutive failures the endpoint is no longer trusted to answer
// at all, and the question becomes a different one — IS THE JOB STILL THERE?
-// `/api/jobs/active` is cheap and is what the editor's own head polls. Listed
-// ⇒ keep waiting, however long that takes. Absent ⇒ it ended while we could not
-// see it, so one last log poll reads the terminal status; if even that fails,
-// the follow gives up rather than claiming an outcome it never read.
+// `GET /api/ops/job/<id>` answers it (behind the token): one registry record
+// and one sidecar. It used to be `/api/jobs/active`, the UI's live view — a
+// next.config REWRITE onto /api/view/activeJobs, kept at its old path for the
+// pages and pinned widgets that poll it, which builds the whole live payload
+// (channel stats, the disk gate, every runner's status) to answer a yes/no, on
+// a server already too busy to answer the log poll. Queued or running ⇒ keep
+// waiting, however long that takes. Ended ⇒ the status it ended with, after
+// one last log poll for the lines we missed. Unknown (404) ⇒ the follow gives
+// up rather than claiming an outcome it never read.
+//
+// "archived" IS NOT AN OUTCOME EITHER. The log route says it for any id the
+// registry no longer holds (evicted past 100, or the editor restarted); the
+// job's sidecar still says how it ended, and the job route reads it.
+//
+// QUEUE POSITION is printed (stderr, unless --quiet) whenever it changes while
+// the job waits: "queued — position 3 on youtube". A job queued behind hours of
+// other work otherwise looks exactly like a hung command.
//
// `--wait-timeout` bounds the whole thing for a caller that cannot hang (CI,
// an agent). Default none, because the honest default for a queue that may hold
@@ -606,15 +983,26 @@ export async function followJob(jobId, quiet, opts = {}) {
const doFetch = opts.fetch ?? fetch;
const sleep = opts.sleep ?? ((ms) => new Promise((r) => setTimeout(r, ms)));
const now = opts.now ?? (() => Date.now());
+ const say = opts.say ?? ((line) => process.stderr.write(`${line}\n`));
const deadline =
opts.timeoutSeconds > 0 ? now() + opts.timeoutSeconds * 1000 : null;
const base = opts.baseUrl ?? baseUrl();
+ const headers = opts.headers ?? authHeaders();
const id = encodeURIComponent(jobId);
let from = 0;
let failures = 0;
let probeFailures = 0;
let backoff = 1000;
+ let lastPlace = null;
+
+ const notePlace = (status, queueKey, position) => {
+ if (quiet || status !== "queued" || !(position > 0)) return;
+ const place = `${position}@${queueKey ?? ""}`;
+ if (place === lastPlace) return;
+ lastPlace = place;
+ say(`[${jobId}] queued — position ${position}${queueKey ? ` on ${queueKey}` : ""}`);
+ };
// One log poll. Returns the status, or null when the poll itself failed —
// never throws, so a transient fetch rejection cannot end the follow.
@@ -633,19 +1021,24 @@ export async function followJob(jobId, quiet, opts = {}) {
// Only advance once the chunk is in hand: a poll that failed halfway
// must re-ask for the same offset.
from = payload.nextOffset ?? from;
+ notePlace(payload.status, payload.queueKey, payload.queuePosition);
return payload.status ?? null;
} catch {
return null;
}
};
- const stillListed = async () => {
+ // The job route: { status } when it answered, "gone" when the editor does
+ // not know the id at all, null when the probe itself failed (or the editor
+ // predates the route — an HTML 404 is not an answer).
+ const probe = async () => {
try {
- const res = await doFetch(`${base}/api/jobs/active`);
- if (!res.ok) return null;
- const payload = await res.json();
- const jobs = Array.isArray(payload.jobs) ? payload.jobs : [];
- return jobs.some((j) => j && j.id === jobId);
+ const res = await doFetch(`${base}/api/ops/job/${id}`, { headers });
+ const payload = await res.json().catch(() => null);
+ if (!payload || typeof payload !== "object") return null;
+ if (res.status === 404 && payload.ok === false) return "gone";
+ if (!res.ok || !payload.job) return null;
+ return payload.job.status ?? null;
} catch {
return null;
}
@@ -654,6 +1047,7 @@ export async function followJob(jobId, quiet, opts = {}) {
// "queued" and "running" are the two NON-answers. Everything else is the job
// having ended, which is the only thing worth returning.
const terminal = (s) => s !== null && s !== "queued" && s !== "running";
+ const ended = (s) => s === "done" || s === "failed" || s === "cancelled";
for (;;) {
const status = await pollLog();
@@ -661,30 +1055,28 @@ export async function followJob(jobId, quiet, opts = {}) {
failures = 0;
probeFailures = 0;
backoff = 1000;
+ if (status === "archived") {
+ // Not in the registry: the sidecar knows how it ended.
+ const real = await probe();
+ return ended(real) ? real : status;
+ }
if (terminal(status)) return status;
} else {
failures++;
if (failures >= 3) {
- const listed = await stillListed();
- if (listed === false) {
- // Gone from the active list: it went terminal while the log endpoint
- // was unreachable. One more try at the status it ended with — and
- // ONLY a terminal one is an outcome. A log endpoint that came back
- // answering "running" means the active list was stale, not that the
- // job finished; returning that printed "running" as the result and
- // exited 1 for a job that was fine.
- const final = await pollLog();
- if (terminal(final)) return final;
- if (final === null) {
- throw new Error(
- `lost contact with job ${jobId}: it is no longer active and its log could not be read`,
- );
- }
- // Readable again and still going: back to waiting, from scratch.
- failures = 0;
- probeFailures = 0;
- backoff = 1000;
- } else if (listed === true) {
+ const real = await probe();
+ if (ended(real)) {
+ // It ended while the log endpoint was unreachable. One more try for
+ // the lines we missed — its status does not override the record's.
+ await pollLog();
+ return real;
+ }
+ if (real === "gone") {
+ throw new Error(
+ `lost contact with job ${jobId}: the editor no longer knows it and its log could not be read`,
+ );
+ }
+ if (real === "queued" || real === "running") {
// A job we can still see is a job to wait for.
failures = 0;
probeFailures = 0;
@@ -697,7 +1089,7 @@ export async function followJob(jobId, quiet, opts = {}) {
probeFailures++;
if (probeFailures >= MAX_PROBE_FAILURES) {
throw new Error(
- `lost contact with the editor at ${base}: ${MAX_PROBE_FAILURES} consecutive failed polls of job ${jobId} and of /api/jobs/active`,
+ `lost contact with the editor at ${base}: ${MAX_PROBE_FAILURES} consecutive failed polls of job ${jobId} and of /api/ops/job/${jobId}`,
);
}
}
@@ -713,6 +1105,34 @@ export async function followJob(jobId, quiet, opts = {}) {
}
}
+// Follow every id to its end, one after another, and say how each ended.
+// 0 only when every one finished `done`. A job still queued says where it waits
+// (followJob), so a list of ids behind one long job reads as a queue, not a
+// hang. In result mode (transcribe) each job's RESULT goes to stdout.
+async function followAll(jobIds, parsed, resultMode = false) {
+ let worst = 0;
+ for (const jobId of jobIds) {
+ const capture = resultMode ? makeResultCapture(parsed.resultMarker) : null;
+ const status = await followJob(jobId, parsed.quiet, {
+ timeoutSeconds: parsed.waitTimeout ?? 0,
+ ...(capture ? { onContent: capture.feed } : {}),
+ });
+ if (capture) {
+ const { echo, result } = capture.finish();
+ if (echo && !parsed.quiet) process.stderr.write(echo);
+ if (result !== null) {
+ console.log(JSON.stringify(result, null, 2));
+ } else if (status === "done") {
+ console.error(`[${jobId}] finished but its log carries no result`);
+ worst = 1;
+ }
+ }
+ console.error(`[${jobId}] ${status}`);
+ if (status !== "done") worst = 1;
+ }
+ return worst;
+}
+
async function main() {
const parsed = parseArgs(process.argv.slice(2));
if (parsed.help) {
@@ -729,6 +1149,11 @@ async function main() {
console.log(ACTIONS.join("\n"));
return 0;
}
+ if (parsed.waitFor) {
+ // `job wait <id…>`: no request of its own, only the follow.
+ loadEditorEnv();
+ return followAll(parsed.waitFor, parsed);
+ }
if (parsed.bodyFile) {
let raw;
try {
@@ -755,6 +1180,7 @@ async function main() {
if (parsed.defaultSource && parsed.body && parsed.body.source === undefined) {
parsed.body = { ...parsed.body, source: parsed.defaultSource };
}
+ const envSources = loadEditorEnv();
const url = `${baseUrl()}${parsed.path}`;
const res = await fetch(url, {
method: parsed.method,
@@ -772,6 +1198,9 @@ async function main() {
console.error(`HTTP ${res.status}: ${text.slice(0, 500)}`);
return 1;
}
+ if (res.status === 401 || res.status === 503) {
+ console.error(tokenHint(res.status, envSources));
+ }
// A result-carrying action under --wait keeps stdout for the RESULT: the
// response goes to stderr with the log.
const resultMode = Boolean(parsed.wait && parsed.resultMarker);
@@ -798,26 +1227,7 @@ async function main() {
printPreviewUrls(payload);
return 0;
}
- let worst = 0;
- for (const jobId of jobIds) {
- const capture = resultMode ? makeResultCapture(parsed.resultMarker) : null;
- const status = await followJob(jobId, parsed.quiet, {
- timeoutSeconds: parsed.waitTimeout ?? 0,
- ...(capture ? { onContent: capture.feed } : {}),
- });
- if (capture) {
- const { echo, result } = capture.finish();
- if (echo && !parsed.quiet) process.stderr.write(echo);
- if (result !== null) {
- console.log(JSON.stringify(result, null, 2));
- } else if (status === "done") {
- console.error(`[${jobId}] finished but its log carries no result`);
- worst = 1;
- }
- }
- console.error(`[${jobId}] ${status}`);
- if (status !== "done") worst = 1;
- }
+ const worst = await followAll(jobIds, parsed, resultMode);
// LAST, after the logs: with --wait the response scrolled off minutes ago,
// and the alias is the one thing the operator came for.
printPreviewUrls(payload);
diff --git a/scripts/archilyzer-ops.test.mjs b/scripts/archilyzer-ops.test.mjs
@@ -7,10 +7,15 @@ import assert from "node:assert/strict";
import test from "node:test";
import {
TRANSCRIBE_RESULT_MARKER,
+ editorEnvFiles,
followJob,
+ loadEditorEnv,
+ mainWorktreeOf,
makeResultCapture,
parseArgs,
+ parseDotenv,
previewUrlsIn,
+ tokenHint,
usage,
} from "./archilyzer-ops.mjs";
@@ -170,16 +175,24 @@ test("--wait-timeout is parsed, and refuses a non-number", () => {
});
// A fake editor. `log` is the sequence of log-poll outcomes — an Error is
-// thrown at the caller the way a starved server makes `fetch` reject.
-function fakeEditor({ log = [], active = [] } = {}) {
- const calls = { log: 0, active: 0 };
- const reply = (body) => ({ ok: true, json: async () => body });
- const doFetch = async (url) => {
- if (url.includes("/api/jobs/active")) {
- const next = active[Math.min(calls.active, active.length - 1)];
- calls.active++;
+// thrown at the caller the way a starved server makes `fetch` reject. `job` is
+// the sequence of /api/ops/job/<id> answers: a status string, "gone" (the
+// editor does not know the id: a JSON 404) or an Error.
+function fakeEditor({ log = [], job = [] } = {}) {
+ const calls = { log: 0, job: 0, jobHeaders: [] };
+ const reply = (body, status = 200) => ({
+ ok: status >= 200 && status < 300,
+ status,
+ json: async () => body,
+ });
+ const doFetch = async (url, init) => {
+ if (url.includes("/api/ops/job/")) {
+ const next = job[Math.min(calls.job, job.length - 1)];
+ calls.job++;
+ calls.jobHeaders.push(init?.headers ?? {});
if (next instanceof Error) throw next;
- return reply({ jobs: next });
+ if (next === "gone") return reply({ ok: false, error: "no job" }, 404);
+ return reply({ ok: true, job: { id: "j1", status: next } });
}
const next = log[Math.min(calls.log, log.length - 1)];
calls.log++;
@@ -194,6 +207,7 @@ const follow = (jobId, editor, opts = {}) =>
fetch: editor.doFetch,
sleep: async () => {},
baseUrl: "http://editor",
+ headers: { authorization: "Bearer t" },
...opts,
});
@@ -210,35 +224,78 @@ test("a poll that rejects does not end the follow", async () => {
});
assert.equal(await follow("j1", editor), "done");
// Two failures is below the probe threshold, so it never asked.
- assert.equal(editor.calls.active, 0);
+ assert.equal(editor.calls.job, 0);
});
-test("after three failures it asks whether the job is still there", async () => {
+test("after three failures it asks the job route whether the job is still there", async () => {
const boom = new Error("fetch failed");
- // Still listed ⇒ keep waiting, and the follow ends on the status it finally
+ // Still running ⇒ keep waiting, and the follow ends on the status it finally
// reads rather than on a guess.
const waiting = fakeEditor({
log: [boom, boom, boom, { content: "", nextOffset: 0, status: "done" }],
- active: [[{ id: "j1" }]],
+ job: ["running"],
});
assert.equal(await follow("j1", waiting), "done");
- assert.equal(waiting.calls.active, 1);
+ assert.equal(waiting.calls.job, 1);
+ // The job route is behind the token, and the probe carries it.
+ assert.deepEqual(waiting.calls.jobHeaders[0], { authorization: "Bearer t" });
- // Absent from the active list ⇒ it ended while the log was unreachable, so
- // one final poll reads the terminal status. `failed` is reported, not hidden.
+ // Ended while the log was unreachable ⇒ the record's status is the answer,
+ // even when the last log poll still fails. `failed` is reported, not hidden.
const gone = fakeEditor({
- log: [boom, boom, boom, { content: "", nextOffset: 0, status: "failed" }],
- active: [[]],
+ log: [boom, boom, boom, boom],
+ job: ["failed"],
});
assert.equal(await follow("j1", gone), "failed");
});
-test("a job that is gone AND unreadable refuses to claim an outcome", async () => {
+test("a job the editor no longer knows, with an unreadable log, refuses to claim an outcome", async () => {
const boom = new Error("fetch failed");
- const editor = fakeEditor({ log: [boom], active: [[]] });
+ const editor = fakeEditor({ log: [boom], job: ["gone"] });
await assert.rejects(follow("j1", editor), /lost contact with job j1/);
});
+test("'archived' is not an outcome: the sidecar's status is", async () => {
+ // The log route says "archived" for any id the registry forgot (evicted, or
+ // the editor restarted); --wait used to print it and exit 1 for a job that
+ // finished done.
+ const done = fakeEditor({
+ log: [{ content: "", nextOffset: 0, status: "archived" }],
+ job: ["done"],
+ });
+ assert.equal(await follow("j1", done), "done");
+ // A sidecar that says nothing terminal leaves "archived" standing.
+ const unknown = fakeEditor({
+ log: [{ content: "", nextOffset: 0, status: "archived" }],
+ job: ["archived"],
+ });
+ assert.equal(await follow("j1", unknown), "archived");
+});
+
+test("a waiting job's queue position is said once per change", async () => {
+ const lines = [];
+ const editor = fakeEditor({
+ log: [
+ { content: "", nextOffset: 0, status: "queued", queueKey: "youtube", queuePosition: 3 },
+ { content: "", nextOffset: 0, status: "queued", queueKey: "youtube", queuePosition: 3 },
+ { content: "", nextOffset: 0, status: "queued", queueKey: "youtube", queuePosition: 1 },
+ { content: "", nextOffset: 0, status: "running", queueKey: "youtube", queuePosition: 0 },
+ { content: "", nextOffset: 0, status: "done" },
+ ],
+ });
+ const status = await followJob("j1", false, {
+ fetch: editor.doFetch,
+ sleep: async () => {},
+ baseUrl: "http://editor",
+ say: (l) => lines.push(l),
+ });
+ assert.equal(status, "done");
+ assert.deepEqual(lines, [
+ "[j1] queued — position 3 on youtube",
+ "[j1] queued — position 1 on youtube",
+ ]);
+});
+
test("--wait-timeout gives up on a job that never ends", async () => {
const editor = fakeEditor({
log: [{ content: "", nextOffset: 0, status: "running" }],
@@ -314,9 +371,10 @@ test("previewUrlsIn is empty for a production deploy", () => {
assert.deepEqual(previewUrlsIn(null), []);
});
-test("a recovered log endpoint answering 'running' is not an outcome", async () => {
- // THE BUG: the absent-from-active branch returned whatever the final poll
- // said. A stale active list plus a recovered log endpoint therefore printed
+
+test("a job route answering 'running' keeps the follow going", async () => {
+ // THE OLD BUG, kept pinned: the absent-from-active branch returned whatever
+ // the final poll said, so a stale list plus a recovered log endpoint printed
// "running" as the job's result and exited 1 for a job that was fine.
const boom = new Error("fetch failed");
const editor = fakeEditor({
@@ -327,20 +385,99 @@ test("a recovered log endpoint answering 'running' is not an outcome", async ()
{ content: "", nextOffset: 0, status: "running" },
{ content: "", nextOffset: 0, status: "done" },
],
- active: [[]],
+ job: ["running"],
});
assert.equal(await follow("j1", editor), "done");
});
test("an editor that answers nothing at all is given up on, not waited on", async () => {
// Without --wait-timeout the null-probe path used to wait forever: the log
- // never answers AND the active list never answers, so nothing ever resets
- // the failure count and nothing ever concludes.
+ // never answers AND the job route never answers, so nothing ever resets the
+ // failure count and nothing ever concludes.
const boom = new Error("fetch failed");
- const editor = fakeEditor({ log: [boom], active: [boom] });
+ const editor = fakeEditor({ log: [boom], job: [boom] });
await assert.rejects(follow("j1", editor), /lost contact with the editor/);
});
+// THE TOKEN FROM THE EDITOR'S ENV FILES (A1). Every session used to write a
+// wrapper that sourced editor/.env before `pnpm ops`.
+test("parseDotenv reads KEY=value lines, quotes, export and comments", () => {
+ assert.deepEqual(
+ parseDotenv(
+ [
+ "# a comment",
+ "",
+ "WORKER_TOKEN=abc123",
+ "export ARCHILYZER_EDITOR_URL='http://localhost:3101'",
+ 'QUOTED="has # hash"',
+ "INLINE=value # trailing comment",
+ "not a line",
+ ].join("\n"),
+ ),
+ {
+ WORKER_TOKEN: "abc123",
+ ARCHILYZER_EDITOR_URL: "http://localhost:3101",
+ QUOTED: "has # hash",
+ INLINE: "value",
+ },
+ );
+});
+
+test("loadEditorEnv fills only what is unset, first file wins, and takes only its two keys", () => {
+ const files = {
+ "/a/editor/.env.local": "WORKER_TOKEN=local\n",
+ "/a/editor/.env": "WORKER_TOKEN=plain\nARCHILYZER_EDITOR_URL=http://x:1\nCLOUDFLARE_API_TOKEN=secret\n",
+ };
+ const read = (p) => {
+ if (!(p in files)) throw new Error("ENOENT");
+ return files[p];
+ };
+ const env = {};
+ const sources = loadEditorEnv(env, ["/missing/.env", ...Object.keys(files)], read);
+ assert.deepEqual(env, { WORKER_TOKEN: "local", ARCHILYZER_EDITOR_URL: "http://x:1" });
+ assert.deepEqual(sources, {
+ WORKER_TOKEN: "/a/editor/.env.local",
+ ARCHILYZER_EDITOR_URL: "/a/editor/.env",
+ });
+ // A variable already set — even to "" — is an explicit choice and wins.
+ const set = { WORKER_TOKEN: "" };
+ const s2 = loadEditorEnv(set, Object.keys(files), read);
+ assert.equal(set.WORKER_TOKEN, "");
+ assert.equal(s2.WORKER_TOKEN, "environment");
+});
+
+test("a linked worktree also reads the main worktree's env files, found off disk", () => {
+ const disk = {
+ "/repo/.claude/worktrees/w/.git": "gitdir: /repo/.git/worktrees/w\n",
+ "/repo/.git/worktrees/w/commondir": "../..\n",
+ };
+ const read = (p) => {
+ if (!(p in disk)) throw new Error("ENOENT");
+ return disk[p];
+ };
+ assert.equal(mainWorktreeOf("/repo/.claude/worktrees/w", read), "/repo");
+ assert.deepEqual(editorEnvFiles("/repo/.claude/worktrees/w", read), [
+ "/repo/.claude/worktrees/w/editor/.env.local",
+ "/repo/.claude/worktrees/w/editor/.env",
+ "/repo/editor/.env.local",
+ "/repo/editor/.env",
+ ]);
+ // The main checkout: `.git` is a directory, so reading it as a file fails.
+ assert.equal(mainWorktreeOf("/repo", read), null);
+ assert.deepEqual(editorEnvFiles("/repo", read), [
+ "/repo/editor/.env.local",
+ "/repo/editor/.env",
+ ]);
+});
+
+test("the token hint names where the token came from, never its value", () => {
+ assert.match(tokenHint(503, {}, "http://e"), /runs without WORKER_TOKEN/);
+ assert.match(tokenHint(401, {}, "http://e"), /no WORKER_TOKEN in the environment or in /);
+ const h = tokenHint(401, { WORKER_TOKEN: "/r/editor/.env" }, "http://e");
+ assert.match(h, /from \/r\/editor\/\.env is not the one the editor at http:\/\/e runs with/);
+});
+
+
test("--wait-timeout implies --wait", () => {
// A timeout on a wait nobody asked for is a typo with no effect, not a
// preference to honour silently.
@@ -555,3 +692,94 @@ test("followJob hands each chunk to onContent", async () => {
assert.equal(status, "done");
assert.deepEqual(seen, ["a\n", "b\n"]);
});
+
+// JOBS OVER OPS (A2): the /jobs list and its row buttons, without a browser.
+test("get job reads one job, --tail adds its log's last lines", () => {
+ assert.equal(parseArgs(["get", "job", "01ABC"]).path, "/api/ops/job/01ABC");
+ assert.equal(parseArgs(["get", "job", "01ABC", "--tail"]).path, "/api/ops/job/01ABC?tail=40");
+ assert.equal(parseArgs(["get", "job", "01ABC", "--tail", "5"]).path, "/api/ops/job/01ABC?tail=5");
+ assert.equal(parseArgs(["get", "job", "--tail=7", "01ABC"]).path, "/api/ops/job/01ABC?tail=7");
+ assert.match(parseArgs(["get", "job"]).error, /needs an argument/);
+ assert.match(parseArgs(["get", "job", "x", "--tail=0"]).error, /--tail takes/);
+});
+
+test("get jobs carries its filters as a query, and refuses one it does not take", () => {
+ assert.equal(parseArgs(["get", "jobs"]).path, "/api/ops/jobs");
+ assert.equal(
+ parseArgs(["get", "jobs", "--failed", "--kind", "sync", "--slug=demo", "--limit", "5"]).path,
+ "/api/ops/jobs?failed=1&kind=sync&slug=demo&limit=5",
+ );
+ assert.equal(parseArgs(["get", "jobs", "--active"]).path, "/api/ops/jobs?active=1");
+ assert.match(parseArgs(["get", "jobs", "--tail"]).error, /get jobs does not take --tail/);
+ assert.match(parseArgs(["get", "channel", "x", "--active"]).error, /does not take --active/);
+ assert.match(parseArgs(["get", "jobs", "--limit", "0"]).error, /--limit needs a whole number/);
+ assert.match(parseArgs(["get", "jobs", "--kind"]).error, /--kind needs a value/);
+ // A read's flag on an action is a mistake, not a silent no-op.
+ assert.match(parseArgs(["sync", "--failed"]).error, /belongs to a read/);
+});
+
+test("job <verb> <id...> posts the verb and the ids to /api/ops/job", () => {
+ const p = parseArgs(["job", "cancel", "A1", "B2", "--wait"]);
+ assert.equal(p.method, "POST");
+ assert.equal(p.path, "/api/ops/job");
+ assert.deepEqual(p.body, { verb: "cancel", ids: ["A1", "B2"] });
+ assert.equal(p.wait, true);
+ assert.deepEqual(parseArgs(["job", "retry-failed"]).body, { verb: "retry-failed" });
+ for (const verb of ["drain", "promote", "force-release", "retry"]) {
+ assert.deepEqual(parseArgs(["job", verb, "A1"]).body, { verb, ids: ["A1"] });
+ }
+ assert.match(parseArgs(["job", "kill", "A1"]).error, /unknown verb "kill"/);
+ assert.match(parseArgs(["job", "cancel"]).error, /needs one or more job ids/);
+ assert.match(parseArgs(["job", "retry-failed", "A1"]).error, /takes no ids/);
+ assert.match(parseArgs(["job", "cancel", "A1", "--json", "{}"]).error, /not a body/);
+});
+
+test("job wait <id...> sends nothing and follows the ids", () => {
+ const p = parseArgs(["job", "wait", "A1", "B2", "--wait-timeout", "60"]);
+ assert.deepEqual(p.waitFor, ["A1", "B2"]);
+ assert.equal(p.method, undefined);
+ assert.equal(p.waitTimeout, 60);
+});
+
+test("usage names the job reads and verbs", () => {
+ assert.match(usage(), /pnpm ops get job <id> \[--tail/);
+ assert.match(usage(), /pnpm ops get jobs \[--active \| --failed\]/);
+ assert.match(usage(), /job <verb> <id>\.\.\. is a \/jobs row's button/);
+});
+
+// THE READ SIDE (A3).
+test("the read-side nouns are GETs on their routes, named in the usage", () => {
+ assert.equal(parseArgs(["get", "settings"]).path, "/api/ops/settings");
+ assert.equal(parseArgs(["get", "settings", "autoQueue"]).path, "/api/ops/settings?key=autoQueue");
+ for (const noun of ["storage", "sites", "workers", "auto-queue", "scheduler"]) {
+ const p = parseArgs(["get", noun]);
+ assert.equal(p.method, "GET", noun);
+ assert.equal(p.path, `/api/ops/${noun}`, noun);
+ }
+ assert.equal(parseArgs(["get", "cleanup", "the-quartering"]).path, "/api/ops/cleanup/the-quartering");
+ assert.match(parseArgs(["get", "cleanup"]).error, /needs an argument/);
+ assert.match(usage(), /pnpm ops get settings \[<key>\]/);
+ assert.match(usage(), /get cleanup <slug> is one channel's \/cleanup row/);
+});
+
+// ARCHIVAL WRITES (A4).
+test("the archival writes are POSTs to their routes, each named in the usage", () => {
+ for (const action of [
+ "settings",
+ "clear-platform-hold",
+ "workers",
+ "transcribe-one",
+ "delete-file",
+ "do-not-clean",
+ "cleanup",
+ ]) {
+ const p = parseArgs([action, "--json", "{}"]);
+ assert.equal(p.method, "POST", action);
+ assert.equal(p.path, `/api/ops/${action}`, action);
+ assert.match(usage(), new RegExp(`^${action} `, "m"), action);
+ }
+ // `get settings` and `settings` are two different requests on one route.
+ assert.equal(parseArgs(["get", "settings"]).method, "GET");
+ assert.match(usage(), /"dryRun": true answers each channel's preview/);
+ assert.match(usage(), /"publish", "enabled"\?, "held"\?, "action"\?/);
+});