commit e7272a7246f174bdf21d506ad885049998668741
parent 093efc10a5ae7f19ebbc7a3737f4bd5346207734
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Fri, 9 Oct 2026 11:42:39 -0400
ops: the /jobs list and its row buttons over pnpm ops
Release 19 slice A2. GET /api/ops/jobs lists jobs (?active = the live head in
queue order with each job's place; ?failed, ?kind, ?slug, ?limit; a filtered
list reads the newest 2000). POST /api/ops/job {verb, ids} presses a /jobs
row's button per id — cancel, drain, promote, force-release, retry — through
editor/app/jobs/actions.ts, the same server actions the page calls, and
answers each id in `results`; `retry-failed` is Retry all. retryAllFailedAction
now also returns the new jobs' ids (and cancels the streams nobody reads).
The CLI gains `get job <id> [--tail [N]]`, `get jobs [--active|--failed]
[--kind] [--slug] [--limit]`, `job <verb> <id...>` and `job wait <id...>`,
which follows ids through followAll (the --wait loop, factored out). A read's
flag given to the wrong noun or to an action is refused, not ignored.
ops-api.spec gains one test: the success paths, which revalidate /jobs and so
need a running server.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Diffstat:
8 files changed, 536 insertions(+), 8 deletions(-)
diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md
@@ -2,6 +2,7 @@
## [Unreleased]
- **`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.
- **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/ops/job/route.test.ts b/editor/app/api/ops/job/route.test.ts
@@ -20,6 +20,28 @@ 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> };
@@ -127,3 +149,64 @@ test("a queued job says where it waits and who holds the head", async () => {
{ 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/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/e2e/ops-api.spec.ts b/editor/e2e/ops-api.spec.ts
@@ -911,6 +911,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
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)
@@ -104,7 +108,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).
@@ -116,11 +120,47 @@ 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}` : ""}`;
+ },
};
// 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"]);
+
+// 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",
@@ -224,6 +264,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);
@@ -254,6 +299,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) {
@@ -304,14 +382,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 {
@@ -410,6 +529,12 @@ 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 list",
"",
`Actions: ${ACTIONS.join(", ")}`,
@@ -421,6 +546,22 @@ export function usage() {
"--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.",
+ "",
'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":',
@@ -900,6 +1041,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 {
diff --git a/scripts/archilyzer-ops.test.mjs b/scripts/archilyzer-ops.test.mjs
@@ -692,3 +692,57 @@ 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/);
+});