Archilyzer · Source

archilyzer

Archilyzer
git clone https://archilyzer.pages.dev/source/archilyzer.git
Log | Files | Refs | README | LICENSE

commit 035730758d8877e4e93e452ba03c2cdd819bc2b1
parent 5221cddf6656fd323c88f4a84e94ebe224799a98
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Sun, 20 Sep 2026 19:43:51 -0400

Merge branch 'feat/ops-api'

Diffstat:
MRUNNING_IN_DOCKER.md | 49+++++++++++++++++++++++++++++++++++++++++++++++++
MSETUP.md | 2+-
Mcommon/jobs/jobSpec.ts | 7+++++++
Aeditor/app/api/ops/_lib.ts | 209+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/api/ops/build-deploy/route.ts | 33+++++++++++++++++++++++++++++++++
Aeditor/app/api/ops/build-index/route.ts | 11+++++++++++
Aeditor/app/api/ops/build-site/route.ts | 59+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/api/ops/channel-config/route.ts | 47+++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/api/ops/channel-priority/route.ts | 91+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/api/ops/channel/[slug]/route.ts | 100+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/api/ops/download-missing/route.ts | 21+++++++++++++++++++++
Aeditor/app/api/ops/import-video/route.ts | 17+++++++++++++++++
Aeditor/app/api/ops/lane/route.ts | 80+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/api/ops/metadata-scan/route.ts | 16++++++++++++++++
Aeditor/app/api/ops/refresh-report/route.ts | 43+++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/api/ops/relocate-back/route.ts | 20++++++++++++++++++++
Aeditor/app/api/ops/relocate/route.ts | 39+++++++++++++++++++++++++++++++++++++++
Aeditor/app/api/ops/retry-bucket/route.ts | 72++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/api/ops/sync/route.ts | 20++++++++++++++++++++
Aeditor/app/channels/components/channelConfigToForm.ts | 186+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/e2e/ops-api.spec.ts | 494+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mpackage.json | 3++-
Mplans/FACTS.md | 62++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Ascripts/archilyzer-ops.mjs | 245+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Ascripts/archilyzer-ops.test.mjs | 72++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
25 files changed, 1996 insertions(+), 2 deletions(-)

diff --git a/RUNNING_IN_DOCKER.md b/RUNNING_IN_DOCKER.md @@ -224,6 +224,55 @@ boot self-updates it. Or, once: docker compose exec editor yt-dlp -U ``` +### Driving the editor without a browser + +Every editor gesture is a server action, which is fine for a person and hostile +to a script: there is no URL to POST to. `/api/ops/*` is a thin layer over the +**same actions** — one route per gesture, no rule of its own — so a shell, a cron +job or an agent can run the archive without Playwright. + +It is gated by the **same `WORKER_TOKEN`** as `/api/worker/*`, deliberately: that +variable already means "this instance takes instructions from something that is +not the browser in front of it". Unset on the server and every route answers +**503** (the surface is off until you opt in); wrong or missing on the caller and +it answers **401**. + +```sh +export ARCHILYZER_EDITOR_URL=http://localhost:3001 +export WORKER_TOKEN=<the same secret the editor is running with> + +pnpm ops sync --json '{"slug":"the-quartering"}' --wait +pnpm ops metadata-scan --json '{"slug":"the-quartering"}' +pnpm ops channel-config --json '{"slug":"the-quartering","patch":{"downloadFilterExclude":"rerun"}}' +pnpm ops channel-priority --json '{"slugs":["the-quartering"],"operation":"download","tier":"paused"}' +pnpm ops lane --json '{"lane":"download","held":true}' +pnpm ops refresh-report --json '{"all":true}' +pnpm ops get channel the-quartering +pnpm ops list # every action name +``` + +Three things to know before you script against it: + +- **A job-starting action returns a `jobId` and does not stream.** The job may + sit in a platform queue behind other work for hours, so "started" is the + honest answer; `--wait` follows `/api/jobs/<id>/log` to the end and exits with + the job's status. +- **Unknown body keys are a 400.** A misspelled `downloadFilterExclude` would + otherwise save cleanly and leave a channel downloading everything. +- **`channel-config` patch keys are the CONFIGURE FORM's field names**, not + `config.json`'s — `downloadFilterInclude` / `downloadFilterExclude` rather than + a `downloadFilter` object. That is what routes them through the form's own + validators, so a bad regex is refused here with the sentence the form shows. + `""` clears a field, exactly as clearing the input does. + +The read side needs no new routes for jobs: `/api/jobs/active`, +`/api/jobs/<id>/log`, `/api/scheduler/status` and `/api/auto-queue/status` +already exist. `GET /api/ops/channel/<slug>` is the one addition — config, +report totals, bucket sizes, priority and, the part no directory listing can +tell you, whether the channel's media is actually **reachable**. Add `--counts` +(`?counts=1`) for the live on-disk counts; it is opt-in because it walks every +video directory, eleven thousand of them on the largest channel here. + ### Booting without resuming work `editor/instrumentation.ts` arms the sync heartbeat and every enabled auto-queue diff --git a/SETUP.md b/SETUP.md @@ -273,7 +273,7 @@ any of them via environment variables before launching: | `PARAKEET_CLI` / `PARAKEET_MODEL` / `PARAKEET_STITCH_BIN` | `parakeet-cli` / — / `scripts/parakeet-stitch.mjs` | parakeet.cpp CLI, model, and wrapper. | | `FFMPEG_BIN` / `FFPROBE_BIN` | `ffmpeg` / `ffprobe` (PATH) | Audio transcode + duration checks. | | `RSYNC_BIN` | `rsync` (PATH) | Saved-video backup. | -| `WORKER_TOKEN` | — | Bearer token for the remote-worker transcription API (set on both ends when used). | +| `WORKER_TOKEN` | — | Bearer token for the remote-worker transcription API (set on both ends when used), and for the `/api/ops/*` HTTP layer over the editor's actions — see [RUNNING_IN_DOCKER.md](RUNNING_IN_DOCKER.md#driving-the-editor-without-a-browser) and `pnpm ops`. Unset means both surfaces are off. | Feature-area docs cover their own env vars: [SCHEDULED_SYNC.md](SCHEDULED_SYNC.md) (`SYNC_HEARTBEAT_SECONDS`, `SYNC_TICK_URL`, `SYNC_TICK_TOKEN`) and diff --git a/common/jobs/jobSpec.ts b/common/jobs/jobSpec.ts @@ -53,6 +53,13 @@ const REPLAY_BUCKETS: ReadonlySet<string> = new Set<ReplayBucket>([ "supersededAutoSubs", ]); +// Is this bucket name one a replayed job can be re-derived from? Exported so a +// caller naming a bucket (the ops API's retry-bucket route) can decide whether +// the job is replayable without re-spelling the set. +export function isReplayBucket(value: unknown): value is ReplayBucket { + return typeof value === "string" && REPLAY_BUCKETS.has(value); +} + // Defensive parse for a spec read back from JSON (the <id>.meta.json sidecar). // Returns null on anything malformed so a hand-edited or stale file can't // crash a reader. Mirrors the tolerance of readJobMeta / readWorkerDefaults. diff --git a/editor/app/api/ops/_lib.ts b/editor/app/api/ops/_lib.ts @@ -0,0 +1,209 @@ +import { NextResponse } from "next/server"; +import { authorizeWorkerRequest } from "yt-dlp-transcript-common/lib/workerToken"; +import { isValidChannelSlug } from "yt-dlp-transcript-common/controller/channels"; +import type { StreamActionResult } from "yt-dlp-transcript-common/jobs/streamCommand"; +import type { QueueOutcome } from "../../channels/lib/queueForSlugs"; + +// THE OPS API IS ADAPTERS, AND NOTHING ELSE. +// +// Every route under /api/ops is ~5 lines that validate a JSON body and call ONE +// existing server action. No route may contain a rule the UI does not already +// enforce: the point of the layer is that an agent driving the editor over HTTP +// and an operator clicking the same button get the same refusal, with the same +// sentence, from the same code. A check written here would be a second opinion +// nobody maintains. +// +// THE TOKEN IS THE WORKER TOKEN, ON PURPOSE. `WORKER_TOKEN` already gates the +// LAN worker protocol and already means "this instance accepts instructions +// from something that is not the browser in front of it". A second secret would +// be a second thing to distribute, rotate and leave unset; the failure modes are +// identical, so the gate is. Unset => 503 (the surface is off, you opt in), +// wrong => 401. +// +// A JOB-STARTING ROUTE RETURNS A jobId AND NEVER STREAMS. runManagedFunction +// hands back a ReadableStream the browser consumes; an HTTP caller wants to +// disconnect and poll. So every adapter cancels the stream (which stops pushing +// into the controller and leaves the on-disk log running — see streamCommand's +// `cancel()` note) and returns the id. Follow it with /api/jobs/<id>/log. + +export type OpsBody = Record<string, unknown>; + +// Thrown by the field readers below; caught by `ops()` and rendered as a 400. +export class OpsInputError extends Error {} + +export function opsFail( + error: string, + status = 400, + extra?: Record<string, unknown>, +): NextResponse { + return NextResponse.json({ ok: false, error, ...extra }, { status }); +} + +// Auth + body parse + unknown-key rejection, wrapped around one handler. +// +// UNKNOWN KEYS ARE A 400, not a silent ignore. A caller that misspells +// `downloadFilterExclude` would otherwise get a cheerful `{ ok: true }` and a +// channel that still downloads everything. The allow-list IS the route's +// documented body shape. +export async function ops( + request: Request, + allowedKeys: readonly string[], + run: (body: OpsBody) => Promise<NextResponse>, +): Promise<NextResponse> { + const auth = authorizeWorkerRequest(request.headers.get("authorization")); + if (!auth.ok) return opsFail(auth.error, auth.status); + let body: unknown; + try { + body = await request.json(); + } catch { + return opsFail("malformed JSON body"); + } + if (typeof body !== "object" || body === null || Array.isArray(body)) { + return opsFail("body must be a JSON object"); + } + const unknown = Object.keys(body as OpsBody).filter( + (k) => !allowedKeys.includes(k), + ); + if (unknown.length) { + return opsFail( + `unknown key(s): ${unknown.join(", ")} — this route accepts ${ + allowedKeys.length ? allowedKeys.join(", ") : "no keys" + }`, + ); + } + try { + return await run(body as OpsBody); + } catch (e) { + if (e instanceof OpsInputError) return opsFail(e.message); + return opsFail((e as Error).message, 500); + } +} + +// A GET route's gate. Same token, no body. +export function opsAuth(request: Request): NextResponse | null { + const auth = authorizeWorkerRequest(request.headers.get("authorization")); + return auth.ok ? null : opsFail(auth.error, auth.status); +} + +// --- field readers ---------------------------------------------------------- + +export function reqString(body: OpsBody, key: string): string { + const v = body[key]; + if (typeof v !== "string" || !v.trim()) { + throw new OpsInputError(`"${key}" is required and must be a non-empty string`); + } + return v.trim(); +} + +// A CHANNEL SLUG, NOT MERELY A STRING. Every slug below reaches a `path.join` +// under `channelsDir`, and the readers swallow their own errors — so a +// traversing segment would fail SILENTLY (an empty config, an "empty channel") +// rather than loudly. `isValidChannelSlug` is CHANNEL_SLUG_RE, which forbids +// "/" and "..", and is what every other slug-taking surface in the app uses. +// +// One reader for every route rather than a check per route: a route added later +// gets this for free by calling reqSlug instead of reqString, and there is one +// place to be wrong. +export function reqSlug(body: OpsBody, key: string): string { + const v = reqString(body, key); + if (!isValidChannelSlug(v)) { + throw new OpsInputError( + `"${v}" is not a valid channel slug (letters, digits, ".", "_", "-"; must start with a letter or digit)`, + ); + } + return v; +} + +export function reqSlugs(body: OpsBody, key: string): string[] { + const values = reqStringArray(body, key); + for (const v of values) { + if (!isValidChannelSlug(v)) { + throw new OpsInputError( + `"${v}" is not a valid channel slug (letters, digits, ".", "_", "-"; must start with a letter or digit)`, + ); + } + } + return values; +} + +export function optString(body: OpsBody, key: string): string | undefined { + const v = body[key]; + if (v === undefined) return undefined; + if (typeof v !== "string") { + throw new OpsInputError(`"${key}" must be a string`); + } + return v; +} + +export function optBool(body: OpsBody, key: string): boolean | undefined { + const v = body[key]; + if (v === undefined) return undefined; + if (typeof v !== "boolean") { + throw new OpsInputError(`"${key}" must be a boolean`); + } + return v; +} + +export function reqStringArray(body: OpsBody, key: string): string[] { + const v = body[key]; + if ( + !Array.isArray(v) || + v.length === 0 || + v.some((s) => typeof s !== "string" || !s.trim()) + ) { + throw new OpsInputError( + `"${key}" is required and must be a non-empty array of strings`, + ); + } + return (v as string[]).map((s) => s.trim()); +} + +export function oneOf<T extends string>( + body: OpsBody, + key: string, + values: readonly T[], +): T { + const v = reqString(body, key); + if (!(values as readonly string[]).includes(v)) { + throw new OpsInputError(`"${key}" must be one of ${values.join(", ")}`); + } + return v as T; +} + +// --- result mapping --------------------------------------------------------- + +// StreamActionResult -> { ok: true, jobId } | { ok: false, error, info? }. +// The stream is cancelled, never returned: see the header. +export function jobResponse(result: StreamActionResult): NextResponse { + if (!result.ok) { + return opsFail(result.error, 400, result.info ? { info: true } : undefined); + } + void result.stream.cancel(); + return NextResponse.json({ ok: true, jobId: result.jobId }); +} + +// The `{ error } | undefined` shape every /channels action returns. +export function actionResponse( + result: { error: string } | undefined, +): NextResponse { + if (result?.error) return opsFail(result.error); + return NextResponse.json({ ok: true }); +} + +// The `{ ok } | { ok: false, error }` shape the lane/pause actions return. +export function okResponse( + result: { ok: boolean; error?: string }, +): NextResponse { + if (!result.ok) return opsFail(result.error ?? "action failed"); + return NextResponse.json({ ok: true }); +} + +// A bulk fan-out's { queued, skipped }. A skip is not a failure — the caller +// gets both lists and decides, exactly as the bulk bar in the UI does. +export function queueResponse(outcome: QueueOutcome): NextResponse { + return NextResponse.json({ + ok: true, + queued: outcome.queued, + skipped: outcome.skipped, + }); +} diff --git a/editor/app/api/ops/build-deploy/route.ts b/editor/app/api/ops/build-deploy/route.ts @@ -0,0 +1,33 @@ +import { NextResponse } from "next/server"; +import { + buildAndDeployAction, + buildAndDeployAllSitesAction, +} from "../../../sites/lib/buildAction"; +import { jobResponse, OpsInputError, ops, optBool, optString } from "../_lib"; + +export const dynamic = "force-dynamic"; + +// POST { siteId: string, skipArchives? } | { all: true, skipArchives? } +// -> { ok: true, jobId } +// +// One managed job either way (build then deploy, one log, one Cancel), so the +// caller polls /api/jobs/<jobId>/log exactly as for a single site. +export async function POST(request: Request) { + return ops(request, ["siteId", "all", "skipArchives"], async (body) => { + const skipArchives = optBool(body, "skipArchives"); + const all = optBool(body, "all"); + const siteId = optString(body, "siteId"); + if (all) { + if (siteId) throw new OpsInputError('send either "siteId" or "all", not both'); + return jobResponse(await buildAndDeployAllSitesAction(skipArchives)); + } + if (!siteId?.trim()) { + throw new OpsInputError('"siteId" is required (or send { "all": true })'); + } + return jobResponse(await buildAndDeployAction(siteId.trim(), skipArchives)); + }); +} + +export function GET() { + return NextResponse.json({ ok: false, error: "POST only" }, { status: 405 }); +} diff --git a/editor/app/api/ops/build-index/route.ts b/editor/app/api/ops/build-index/route.ts @@ -0,0 +1,11 @@ +import { buildIndexAction } from "../../../sites/lib/buildAction"; +import { jobResponse, ops, optString } from "../_lib"; + +export const dynamic = "force-dynamic"; + +// POST { queueKey? } -> { ok: true, jobId }. Rebuilds the LMDB corpus index. +export async function POST(request: Request) { + return ops(request, ["queueKey"], async (body) => + jobResponse(await buildIndexAction(optString(body, "queueKey"))), + ); +} diff --git a/editor/app/api/ops/build-site/route.ts b/editor/app/api/ops/build-site/route.ts @@ -0,0 +1,59 @@ +import { NextResponse } from "next/server"; +import { + buildAllSitesAction, + buildExportAction, +} from "../../../sites/lib/buildAction"; +import { jobResponse, OpsInputError, ops, optBool } from "../_lib"; + +export const dynamic = "force-dynamic"; + +// POST { siteIds: string[], skipData?, skipArchives? } | { all: true, skipArchives? } +// +// BUILD WITHOUT DEPLOYING. `siteIds` queues one build-export job per site on the +// shared build queue (they run one at a time, as they do from /sites) and +// returns `{ jobs: [{ siteId, jobId }] }` — a list, because there is a job per +// site and a caller waiting on them needs all the ids. `{ all: true }` is the +// single build-all job instead, which is the docker fan-out. +export async function POST(request: Request) { + return ops( + request, + ["siteIds", "all", "skipData", "skipArchives"], + async (body) => { + const skipArchives = optBool(body, "skipArchives"); + if (optBool(body, "all")) { + if (body.siteIds !== undefined) { + throw new OpsInputError('send either "siteIds" or "all", not both'); + } + return jobResponse(await buildAllSitesAction(skipArchives)); + } + const raw = body.siteIds; + if ( + !Array.isArray(raw) || + raw.length === 0 || + raw.some((s) => typeof s !== "string" || !s.trim()) + ) { + throw new OpsInputError( + '"siteIds" is required and must be a non-empty array of strings (or send { "all": true })', + ); + } + const skipData = optBool(body, "skipData"); + const jobs: { siteId: string; jobId: string }[] = []; + const skipped: { siteId: string; reason: string }[] = []; + for (const siteId of (raw as string[]).map((s) => s.trim())) { + const result = await buildExportAction( + siteId, + undefined, + skipData, + skipArchives, + ); + if (!result.ok) { + skipped.push({ siteId, reason: result.error }); + continue; + } + void result.stream.cancel(); + jobs.push({ siteId, jobId: result.jobId }); + } + return NextResponse.json({ ok: true, jobs, skipped }); + }, + ); +} diff --git a/editor/app/api/ops/channel-config/route.ts b/editor/app/api/ops/channel-config/route.ts @@ -0,0 +1,47 @@ +import { getPaths } from "yt-dlp-transcript-common/lib/paths"; +import { readChannelConfig } from "yt-dlp-transcript-common/controller/channels"; +import { + applyChannelFormPatch, + channelConfigToFormData, + validateChannelFormPatch, +} from "../../../channels/components/channelConfigToForm"; +import { updateChannelAction } from "../../../channels/actions"; +import { actionResponse, OpsInputError, ops, reqSlug } from "../_lib"; + +export const dynamic = "force-dynamic"; + +// POST { slug: string, patch: { <Configure-form field>: string|number|boolean|null } } +// +// The patch keys are the FORM's field names, not ChannelConfig's, because the +// form is what validates them: `downloadFilterInclude` / `downloadFilterExclude` +// rather than a `downloadFilter` object, `ytdlpExtraArgs` as a string or an +// array of lines. See channelConfigToForm.ts for why the patch is laid over the +// channel's current form representation instead of being written directly. +// +// `""` (or null) clears a field, exactly as clearing the input does. +export async function POST(request: Request) { + return ops(request, ["slug", "patch"], async (body) => { + const slug = reqSlug(body, "slug"); + const patch = body.patch; + if ( + typeof patch !== "object" || + patch === null || + Array.isArray(patch) + ) { + throw new OpsInputError('"patch" must be a JSON object'); + } + // KEYS FIRST, CHANNEL SECOND. A misspelled field is a fact about the + // request; reporting "channel not found" for it would hide the real error. + try { + validateChannelFormPatch(patch as Record<string, unknown>); + } catch (e) { + throw new OpsInputError((e as Error).message); + } + const paths = getPaths(); + const existing = await readChannelConfig(paths, slug); + if (!existing) throw new OpsInputError(`Channel "${slug}" not found`); + const fd = channelConfigToFormData(existing); + applyChannelFormPatch(fd, patch as Record<string, unknown>); + return actionResponse(await updateChannelAction(slug, undefined, fd)); + }); +} diff --git a/editor/app/api/ops/channel-priority/route.ts b/editor/app/api/ops/channel-priority/route.ts @@ -0,0 +1,91 @@ +import { NextResponse } from "next/server"; +import { + STORED_CHANNEL_TIERS, + PRIORITY_OPERATIONS, + type PriorityOperation, + type StoredChannelTier, +} from "yt-dlp-transcript-common/lib/channelPriority"; +import { + applyChannelPriorityPresetAction, + setChannelOperationTierAction, + setChannelTierAction, +} from "../../../channels/actions"; +import { + actionResponse, + OpsInputError, + ops, + optString, + reqSlugs, +} from "../_lib"; + +export const dynamic = "force-dynamic"; + +// POST { slugs: string[], tier?: "normal"|"low"|"paused"|null, operation?: one +// of sync|transcription|download|digest|backfill, preset?: "sync-only"|"clear" } +// +// Three named gestures, one route, because they are one gesture in the UI: the +// /channels deck's tier control. `operation` present pins that operation's +// override (null clears it back to the base tier); absent sets the BASE tier. +// `preset` is the deck's two shortcuts and takes no tier. +export async function POST(request: Request) { + return ops( + request, + ["slugs", "tier", "operation", "preset"], + async (body) => { + const slugs = reqSlugs(body, "slugs"); + const preset = optString(body, "preset"); + if (preset !== undefined) { + if (preset !== "sync-only" && preset !== "clear") { + throw new OpsInputError('"preset" must be "sync-only" or "clear"'); + } + return actionResponse( + await applyChannelPriorityPresetAction(slugs, preset), + ); + } + const rawTier = body.tier; + if (rawTier !== null && typeof rawTier !== "string") { + throw new OpsInputError( + `"tier" must be one of ${STORED_CHANNEL_TIERS.join(", ")} (or null with an "operation", to clear the override)`, + ); + } + if ( + rawTier !== null && + !(STORED_CHANNEL_TIERS as readonly string[]).includes(rawTier) + ) { + throw new OpsInputError( + `"tier" must be one of ${STORED_CHANNEL_TIERS.join(", ")}`, + ); + } + const operation = optString(body, "operation"); + if (operation !== undefined) { + if (!(PRIORITY_OPERATIONS as readonly string[]).includes(operation)) { + throw new OpsInputError( + `"operation" must be one of ${PRIORITY_OPERATIONS.join(", ")}`, + ); + } + return actionResponse( + await setChannelOperationTierAction( + slugs, + operation as PriorityOperation, + rawTier as StoredChannelTier | null, + ), + ); + } + if (rawTier === null) { + throw new OpsInputError( + 'a null "tier" clears an operation override — name an "operation", or send a stored tier', + ); + } + return actionResponse( + await setChannelTierAction(slugs, rawTier as StoredChannelTier), + ); + }, + ); +} + +export function GET() { + return NextResponse.json( + { ok: false, error: "POST only" }, + { status: 405 }, + ); +} diff --git a/editor/app/api/ops/channel/[slug]/route.ts b/editor/app/api/ops/channel/[slug]/route.ts @@ -0,0 +1,100 @@ +import { NextResponse } from "next/server"; +import { getPaths } from "yt-dlp-transcript-common/lib/paths"; +import { + isValidChannelSlug, + readChannelConfig, + readChannelSnapshot, + readChannelStat, +} from "yt-dlp-transcript-common/controller/channels"; +import { inspectChannelMedia } from "yt-dlp-transcript-common/lib/channelMedia"; +import { + overridesOf, + rankOf, + tierOf, +} from "yt-dlp-transcript-common/lib/channelPriority"; +import { getSettings } from "yt-dlp-transcript-common/lib/settings"; +import { opsAuth } from "../../_lib"; + +export const dynamic = "force-dynamic"; + +// GET /api/ops/channel/<slug> +// +// THE READ SIDE, so a caller never opens transcripts/ itself. Everything a +// script needs to decide what to do next about one channel: its stored config, +// its report totals and bucket SIZES (not the id lists — a bucket can hold tens +// of thousands of ids and a caller deciding "is there work" only needs the +// count; /api/ops/retry-bucket resolves the ids server-side anyway), its +// priority tier, and — the one that cannot be inferred from the filesystem — +// whether its media is reachable. +// +// `media.status` is `inspectChannelMedia`'s, which is the ONE module that can +// tell an unmounted drive from an empty channel. Every enumerator else swallows +// ENOENT on data/ as "no videos", so a script reading counts alone would read an +// unmounted platter as a channel that has downloaded nothing. +// +// `?counts=1` ADDS THE LIVE ON-DISK COUNTS, and it is opt-in because +// readChannelStat walks every video directory — eleven thousand readdirs on the +// largest channel here. The report's totals answer the same question from a +// file, and a poll loop asking for them every few seconds must not be the thing +// that hammers the platter. +export async function GET( + request: Request, + { params }: { params: Promise<{ slug: string }> }, +) { + const denied = opsAuth(request); + if (denied) return denied; + const { slug } = await params; + // The same shape check every POST route applies through reqSlug (_lib.ts); + // spelled out here only because the slug arrives as a route param, not in a + // body. readChannelConfig swallows its own errors, so a traversing segment + // would fail silently rather than loudly. + if (!isValidChannelSlug(slug)) { + return NextResponse.json( + { ok: false, error: `"${slug}" is not a valid channel slug` }, + { status: 400 }, + ); + } + const paths = getPaths(); + const config = await readChannelConfig(paths, slug); + if (!config) { + return NextResponse.json( + { ok: false, error: `Channel "${slug}" not found` }, + { status: 404 }, + ); + } + const priority = getSettings().channelPriority; + const wantCounts = new URL(request.url).searchParams.get("counts") === "1"; + const [snapshot, stat, media] = await Promise.all([ + readChannelSnapshot(paths, slug), + wantCounts ? readChannelStat(paths, slug) : Promise.resolve(null), + inspectChannelMedia(paths, slug, config), + ]); + const buckets = snapshot + ? Object.fromEntries( + Object.entries(snapshot.buckets).map(([k, v]) => [ + k, + Array.isArray(v) ? v.length : v, + ]), + ) + : null; + return NextResponse.json({ + ok: true, + slug, + config, + priority: { + tier: tierOf(priority, slug), + overrides: overridesOf(priority, slug), + rank: rankOf(priority, slug), + }, + // null unless ?counts=1 — see the header. + counts: stat, + media, + report: snapshot + ? { + generatedAt: snapshot.generatedAt, + totals: snapshot.totals, + buckets, + } + : null, + }); +} diff --git a/editor/app/api/ops/download-missing/route.ts b/editor/app/api/ops/download-missing/route.ts @@ -0,0 +1,21 @@ +import { downloadMissingAction } from "../../../channels/[slug]/pipelineActions"; +import { jobResponse, ops, optBool, optString, reqSlug } from "../_lib"; + +export const dynamic = "force-dynamic"; + +// POST { slug, queueKey?, ignoreArchive?, abortOnError? } -> { ok: true, jobId } +export async function POST(request: Request) { + return ops( + request, + ["slug", "queueKey", "ignoreArchive", "abortOnError"], + async (body) => + jobResponse( + await downloadMissingAction( + reqSlug(body, "slug"), + optString(body, "queueKey"), + optBool(body, "ignoreArchive"), + optBool(body, "abortOnError"), + ), + ), + ); +} diff --git a/editor/app/api/ops/import-video/route.ts b/editor/app/api/ops/import-video/route.ts @@ -0,0 +1,17 @@ +import { importVideoAction } from "../../../channels/[slug]/pipelineActions"; +import { jobResponse, ops, optString, reqSlug, reqString } from "../_lib"; + +export const dynamic = "force-dynamic"; + +// POST { slug: string, url: string, queueKey?: string } -> { ok: true, jobId } +export async function POST(request: Request) { + return ops(request, ["slug", "url", "queueKey"], async (body) => + jobResponse( + await importVideoAction( + reqSlug(body, "slug"), + reqString(body, "url"), + optString(body, "queueKey"), + ), + ), + ); +} diff --git a/editor/app/api/ops/lane/route.ts b/editor/app/api/ops/lane/route.ts @@ -0,0 +1,80 @@ +import { NextResponse } from "next/server"; +import { 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 { + pauseLaneAction, + resumeLaneAction, + saveAutoQueueAction, + startAutoQueueAction, + stopAutoQueueAction, +} from "../../../operations/actions"; +import { OpsInputError, okResponse, ops, oneOf, optBool } from "../_lib"; + +export const dynamic = "force-dynamic"; + +// POST { lane: transcription|download|digest|backfill, +// held?: boolean, enabled?: boolean, action?: "start"|"stop"|"drain" } +// +// One lane, up to three independent changes, applied in the order the operator +// would: enable the policy, set the gate, then bring the runner up or down. +// Each is the existing action and nothing more. +// +// `held` IS NOT `enabled`. A hold is a dispatch gate (the runner idle-waits and +// keeps its place); disabling is the policy's master switch. The pair is what +// /operations draws as two separate controls and this route keeps them separate. +// +// `action` mirrors /api/auto-queue/control's body — same three verbs, same +// meanings — so a caller that already drives that route needs nothing new. +export async function POST(request: Request) { + return ops(request, ["lane", "held", "enabled", "action"], async (body) => { + const lane = oneOf(body, "lane", LANES) as AutoQueueKind; + const enabled = optBool(body, "enabled"); + const held = optBool(body, "held"); + const action = body.action; + if ( + action !== undefined && + action !== "start" && + action !== "stop" && + action !== "drain" + ) { + throw new OpsInputError('"action" must be "start", "stop" or "drain"'); + } + if (enabled === undefined && held === undefined && action === undefined) { + throw new OpsInputError( + 'nothing to do — send at least one of "enabled", "held" or "action"', + ); + } + 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 (held !== undefined) { + const gated = held + ? await pauseLaneAction(lane) + : 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); + } + return NextResponse.json({ ok: true, lane }); + }); +} diff --git a/editor/app/api/ops/metadata-scan/route.ts b/editor/app/api/ops/metadata-scan/route.ts @@ -0,0 +1,16 @@ +import { runMetadataScanAction } from "../../../channels/[slug]/pipelineActions"; +import { jobResponse, ops, optString, reqSlug } from "../_lib"; + +export const dynamic = "force-dynamic"; + +// POST { slug: string, queueKey?: string } -> { ok: true, jobId } +export async function POST(request: Request) { + return ops(request, ["slug", "queueKey"], async (body) => + jobResponse( + await runMetadataScanAction( + reqSlug(body, "slug"), + optString(body, "queueKey"), + ), + ), + ); +} diff --git a/editor/app/api/ops/refresh-report/route.ts b/editor/app/api/ops/refresh-report/route.ts @@ -0,0 +1,43 @@ +import { NextResponse } from "next/server"; +import { + refreshAllChannelSnapshotsAction, + refreshChannelSnapshotAction, +} from "../../../channels/actions"; +import { + actionResponse, + OpsInputError, + ops, + optBool, + optString, + reqSlug, +} from "../_lib"; + +export const dynamic = "force-dynamic"; + +// POST { slug: string } | { all: true } +// +// The single-channel form REGENERATES SYNCHRONOUSLY (it is a filesystem scan, +// not a job) and returns `{ ok: true }` once snapshot.json is on disk. The +// `all` form queues one refresh-report job per channel and returns the bulk +// { queued, skipped } the /channels header button shows. +export async function POST(request: Request) { + return ops(request, ["slug", "all"], async (body) => { + const all = optBool(body, "all"); + const slug = optString(body, "slug"); + if (all) { + if (slug) { + throw new OpsInputError('send either "slug" or "all", not both'); + } + const result = await refreshAllChannelSnapshotsAction(); + return NextResponse.json({ ok: true, ...result }); + } + if (!slug?.trim()) { + throw new OpsInputError('"slug" is required (or send { "all": true })'); + } + // Re-read through reqSlug now that we know it is the single-channel form: + // the shape check belongs on the value that reaches a path.join. + return actionResponse( + await refreshChannelSnapshotAction(reqSlug(body, "slug")), + ); + }); +} diff --git a/editor/app/api/ops/relocate-back/route.ts b/editor/app/api/ops/relocate-back/route.ts @@ -0,0 +1,20 @@ +import { moveChannelMediaBackAction } from "../../../channels/[slug]/storageActions"; +import { ops, queueResponse, reqSlugs } from "../_lib"; +import { queueForSlugs } from "../../../channels/lib/queueForSlugs"; + +export const dynamic = "force-dynamic"; + +// POST { slugs: string[] } -> { ok: true, queued, skipped } +// +// One job per slug on the shared relocation queue, and the per-channel action's +// own busy guard decides each one — the same fan-out loop the bulk move uses, +// so a refusal reads identically whether it came from the panel or from here. +export async function POST(request: Request) { + return ops(request, ["slugs"], async (body) => + queueResponse( + await queueForSlugs(reqSlugs(body, "slugs"), { + run: (slug) => moveChannelMediaBackAction(slug), + }), + ), + ); +} diff --git a/editor/app/api/ops/relocate/route.ts b/editor/app/api/ops/relocate/route.ts @@ -0,0 +1,39 @@ +import { bulkRelocateChannelMediaAction } from "../../../channels/bulkStorageActions"; +import { + OpsInputError, + ops, + optString, + queueResponse, + reqSlugs, +} from "../_lib"; + +export const dynamic = "force-dynamic"; + +// POST { slugs: string[], locationId?: string, root?: string } +// -> { ok: true, queued, skipped } +// +// THE DESTINATION IS A LOCATION ID WHEREVER POSSIBLE — the root is resolved on +// the server from settings.storage.locations, so a caller holding a stale root +// 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. +export async function POST(request: Request) { + return ops(request, ["slugs", "locationId", "root"], async (body) => { + const slugs = reqSlugs(body, "slugs"); + const locationId = optString(body, "locationId"); + const root = optString(body, "root"); + 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 }, + ), + ); + }); +} diff --git a/editor/app/api/ops/retry-bucket/route.ts b/editor/app/api/ops/retry-bucket/route.ts @@ -0,0 +1,72 @@ +import { getPaths } from "yt-dlp-transcript-common/lib/paths"; +import { readChannelSnapshot } from "yt-dlp-transcript-common/controller/channels"; +import { isReplayBucket } from "yt-dlp-transcript-common/jobs/jobSpec"; +import { retryBucketAction } from "../../../channels/[slug]/pipelineActions"; +import { + jobResponse, + OpsInputError, + ops, + optBool, + optString, + reqSlug, + reqString, +} from "../_lib"; + +export const dynamic = "force-dynamic"; + +// POST { slug, bucket, queueKey?, abortOnError?, handlingOverride?, +// forceCookies?, replaceAutoSubs? } -> { ok: true, jobId } +// +// THE BUCKET IS RESOLVED FROM THE SNAPSHOT HERE, and that is not new logic: the +// bucket controls in the UI pass `snapshot.buckets[key]` to the same action, and +// the action's own argument is a list of ids. An HTTP caller naming a bucket +// rather than pasting ids is the same gesture — and `bucketKey` travelling with +// it is what makes the job replayable against the CURRENT bucket, exactly as a +// clicked one is. +export async function POST(request: Request) { + return ops( + request, + [ + "slug", + "bucket", + "queueKey", + "abortOnError", + "handlingOverride", + "forceCookies", + "replaceAutoSubs", + ], + async (body) => { + const slug = reqSlug(body, "slug"); + const bucket = reqString(body, "bucket"); + const snapshot = await readChannelSnapshot(getPaths(), slug); + if (!snapshot) { + throw new OpsInputError( + `Channel "${slug}" has no report yet — run /api/ops/refresh-report first.`, + ); + } + const ids = (snapshot.buckets as Record<string, unknown>)[bucket]; + if (!Array.isArray(ids)) { + throw new OpsInputError( + `"${bucket}" is not a bucket on this channel's report — known buckets: ${Object.keys( + snapshot.buckets, + ).join(", ")}`, + ); + } + return jobResponse( + await retryBucketAction( + slug, + ids as string[], + optString(body, "queueKey"), + optBool(body, "abortOnError"), + optString(body, "handlingOverride"), + // Replayable only for the buckets a replay can re-derive; an + // ad-hoc bucket name still runs, it just carries no spec — the + // same distinction the UI's named vs checkbox controls make. + isReplayBucket(bucket) ? bucket : undefined, + optBool(body, "forceCookies"), + optBool(body, "replaceAutoSubs"), + ), + ); + }, + ); +} diff --git a/editor/app/api/ops/sync/route.ts b/editor/app/api/ops/sync/route.ts @@ -0,0 +1,20 @@ +import { syncAction } from "../../../channels/[slug]/pipelineActions"; +import { jobResponse, ops, optBool, optString, reqSlug } from "../_lib"; + +export const dynamic = "force-dynamic"; + +// POST { slug: string, full?: boolean, queueKey?: string } -> { ok: true, jobId } +// +// `full` forces the periodic whole-listing sweep now; without it the job +// decides for itself from the configured cadence. +export async function POST(request: Request) { + return ops(request, ["slug", "full", "queueKey"], async (body) => + jobResponse( + await syncAction( + reqSlug(body, "slug"), + optString(body, "queueKey"), + optBool(body, "full"), + ), + ), + ); +} diff --git a/editor/app/channels/components/channelConfigToForm.ts b/editor/app/channels/components/channelConfigToForm.ts @@ -0,0 +1,186 @@ +import type { ChannelConfig } from "yt-dlp-transcript-common/lib/channelConfig"; + +// THE INVERSE OF parseChannelForm, and it exists for exactly one caller: +// /api/ops/channel-config. +// +// WHY A ROUND TRIP THROUGH FormData rather than a direct config write. The form +// parser is where a download-filter regex is refused, where a sync cadence is +// bounded and where a cleared field is turned into an absent key — and +// `updateChannelAction` DELETES every CHANNEL_FORM_FIELDS key from the stored +// config before layering the parse result on top, so that a cleared input +// actually clears. Feed it a FormData carrying only a patch and it would clear +// everything the patch did not name. So the patch is applied ON TOP of the +// channel's current form representation, which is what this builds, and the +// action then behaves EXACTLY as it does for the browser form. One validator, +// one writer, one set of error sentences. +// +// Every key below is a field name `parseChannelForm` reads. Checkbox fields +// follow the browser: present means checked, absent means unchecked — which is +// why they are emitted only when true. + +// The checkbox-shaped fields (presence, not value). +export const CHANNEL_FORM_FLAGS = [ + "keepSourceVideo", + "downloadFilterIncludeLivestreams", + "audioCheckEnabled", + "audioCheckResumeDuringProbe", +] as const; + +// The value-shaped fields. A patch may set any of these to "" to clear it. +export const CHANNEL_FORM_VALUES = [ + "name", + "handling", + "url", + "platform", + "sourceKind", + "postFetcher", + "socialHandle", + "audioFormat", + "downloadFormat", + "keepLatest", + "extractionMode", + "savedVideosDir", + "ytdlpExtraArgs", + "sleepBetweenDownloadsSeconds", + "syncIntervalMinutes", + "fullSweepIntervalMinutes", + "cookiesFromBrowser", + "cookieMode", + "downloadFilterInclude", + "downloadFilterExclude", + "audioCheckIntervalSeconds", + "audioCheckMaxRollbacks", + "audioCheckCopyTimeoutSeconds", +] as const; + +export type ChannelFormFlag = (typeof CHANNEL_FORM_FLAGS)[number]; +export type ChannelFormValue = (typeof CHANNEL_FORM_VALUES)[number]; +export type ChannelFormField = ChannelFormFlag | ChannelFormValue; + +export const CHANNEL_FORM_FIELD_NAMES: readonly ChannelFormField[] = [ + ...CHANNEL_FORM_VALUES, + ...CHANNEL_FORM_FLAGS, +]; + +function put(fd: FormData, key: string, value: string | undefined | null): void { + if (value === undefined || value === null || value === "") return; + fd.set(key, value); +} + +function putNum(fd: FormData, key: string, value: number | undefined): void { + if (value === undefined || value === null) return; + fd.set(key, String(value)); +} + +// The channel's stored config as the Configure form would post it. +export function channelConfigToFormData(config: ChannelConfig): FormData { + const fd = new FormData(); + // Both are REQUIRED by parseChannelForm and are not CHANNEL_FORM_FIELDS, so + // they always travel. + fd.set("name", config.name ?? ""); + fd.set("handling", config.handling ?? "transcribe"); + put(fd, "url", config.url); + put(fd, "platform", config.platform); + put(fd, "sourceKind", config.sourceKind); + put(fd, "postFetcher", config.postFetcher); + put(fd, "socialHandle", config.socialHandle); + put(fd, "audioFormat", config.audioFormat); + put(fd, "downloadFormat", config.downloadFormat); + if (config.keepSourceVideo) fd.set("keepSourceVideo", "on"); + putNum(fd, "keepLatest", config.keepLatest); + put(fd, "extractionMode", config.extractionMode); + put(fd, "savedVideosDir", config.savedVideosDir); + if (config.ytdlpExtraArgs?.length) { + fd.set("ytdlpExtraArgs", config.ytdlpExtraArgs.join("\n")); + } + putNum(fd, "sleepBetweenDownloadsSeconds", config.sleepBetweenDownloadsSeconds); + putNum(fd, "syncIntervalMinutes", config.syncIntervalMinutes); + putNum(fd, "fullSweepIntervalMinutes", config.fullSweepIntervalMinutes); + put(fd, "cookiesFromBrowser", config.cookiesFromBrowser); + put(fd, "cookieMode", config.cookieMode); + put(fd, "downloadFilterInclude", config.downloadFilter?.include); + put(fd, "downloadFilterExclude", config.downloadFilter?.exclude); + if (config.downloadFilter?.includeLivestreams) { + fd.set("downloadFilterIncludeLivestreams", "on"); + } + if (config.audioCheck?.enabled) { + fd.set("audioCheckEnabled", "on"); + putNum(fd, "audioCheckIntervalSeconds", config.audioCheck.intervalSeconds); + putNum(fd, "audioCheckMaxRollbacks", config.audioCheck.maxRollbacks); + putNum( + fd, + "audioCheckCopyTimeoutSeconds", + config.audioCheck.copyTimeoutSeconds, + ); + if (config.audioCheck.resumeDuringProbe) { + fd.set("audioCheckResumeDuringProbe", "on"); + } + } + return fd; +} + +export type ChannelConfigPatch = Record<string, unknown>; + +// Lay a JSON patch over a FormData built above. Values: a string (or number, +// stringified) sets the field, `""` clears it; a flag takes a boolean. +// `ytdlpExtraArgs` additionally accepts an array of lines, because that is the +// shape it has in config.json and a caller should not have to know the textarea +// is newline-separated. +// +// Throws on a key this form has no field for — the allow-list IS the documented +// body shape, and a misspelled key that silently did nothing would look like a +// successful save. +export function applyChannelFormPatch( + fd: FormData, + patch: ChannelConfigPatch, +): void { + validateChannelFormPatch(patch); + for (const [key, value] of Object.entries(patch)) { + if ((CHANNEL_FORM_FLAGS as readonly string[]).includes(key)) { + if (value) fd.set(key, "on"); + else fd.delete(key); + continue; + } + if (key === "ytdlpExtraArgs" && Array.isArray(value)) { + const joined = (value as string[]).join("\n"); + if (joined.trim()) fd.set(key, joined); + else fd.delete(key); + continue; + } + if (value === null || value === "") { + fd.delete(key); + continue; + } + fd.set(key, typeof value === "number" ? String(value) : (value as string)); + } +} + +// SHAPE ONLY, AND SEPARATE ON PURPOSE: a misspelled field name is a fact about +// the REQUEST, so the route checks it before it looks the channel up. Otherwise +// `{ slug: "typo", patch: { downloadFilterExcluded: … } }` reports the channel +// and never mentions the key that was actually wrong. +export function validateChannelFormPatch(patch: ChannelConfigPatch): void { + for (const [key, value] of Object.entries(patch)) { + if ((CHANNEL_FORM_FLAGS as readonly string[]).includes(key)) { + if (typeof value !== "boolean") { + throw new Error(`"${key}" must be a boolean`); + } + continue; + } + if (!(CHANNEL_FORM_VALUES as readonly string[]).includes(key)) { + throw new Error( + `"${key}" is not a channel config field — accepted: ${CHANNEL_FORM_FIELD_NAMES.join(", ")}`, + ); + } + if (key === "ytdlpExtraArgs" && Array.isArray(value)) { + if (value.some((v) => typeof v !== "string")) { + throw new Error(`"${key}" must be an array of strings`); + } + continue; + } + if (value === null || typeof value === "number" || typeof value === "string") { + continue; + } + throw new Error(`"${key}" must be a string, a number, or "" to clear it`); + } +} diff --git a/editor/e2e/ops-api.spec.ts b/editor/e2e/ops-api.spec.ts @@ -0,0 +1,494 @@ +// /api/ops — the HTTP door onto the editor's server actions. +// +// What this spec is really pinning is the claim the layer rests on: that an +// agent driving the editor over HTTP and an operator clicking the same button +// get the SAME answer from the SAME code. So the assertions are deliberately +// about the shared sentences — the download-filter message that +// 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 503 BRANCH IS NOT REACHABLE FROM HERE. The test server runs with +// WORKER_TOKEN=test-worker-token (editor/package.json, dev:test) and there is one +// server for the whole suite, so no spec can observe the endpoint disabled. 401 +// (missing and wrong) is covered below; the 503 is authorizeWorkerRequest's own +// first branch, shared with /api/worker/* and unit-tested by nothing else +// either. See plans/FACTS.md. + +import { readdir, rm } from "node:fs/promises"; +import { test, expect, type APIRequestContext } from "@playwright/test"; +import { baseUrl } from "./baseUrl"; +import { + generateReport, + pathExists, + readJson, + resetData, + resolvePath, + writeChannelConfig, + writeSettings, +} from "./helpers"; + +const TOKEN = "test-worker-token"; +const AUTH = { authorization: `Bearer ${TOKEN}` }; + +type OpsResponse = { + ok?: boolean; + error?: string; + jobId?: string; + queued?: string[]; + skipped?: { slug: string; reason: string }[]; +}; + +async function ops( + request: APIRequestContext, + action: string, + data: Record<string, unknown>, +): Promise<{ status: number; body: OpsResponse }> { + const res = await request.post(`${baseUrl}/api/ops/${action}`, { + headers: AUTH, + data, + }); + return { status: res.status(), body: (await res.json()) as OpsResponse }; +} + +// EVERY SPEC HERE NAMES THE DISK FLOOR. Without it the merged fixture default +// applies, and a spec that starts a download-shaped job on a nearly-full host +// would be refused by lowDiskError() with a returned { ok: false } no assertion +// reads. 0 is the fixture's own value; naming it makes that a decision. +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); + } +}); + +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" }], + ["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"] }], + ]; + 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("channel-config round-trips a download filter and refuses a bad regex", async ({ + page, + request, +}) => { + test.setTimeout(120_000); + await resetData("title-filter-channel"); + await settings(); + const SLUG = "test-filter"; + const CONFIG = `test-transcripts/channels/${SLUG}/config.json`; + // TWO FIELDS THE CLEAR-THEN-LAYER PATH ACTUALLY THREATENS. name/handling/url + // survive a patch trivially — they are re-posted by the serializer because + // parseChannelForm requires the first two and CHANNEL_FORM_FIELDS spares + // none of the rest. `cookieMode` and `keepLatest` are in CHANNEL_FORM_FIELDS, + // so updateChannelAction deletes them from the baseline before layering, and + // a patch that failed to re-post them would silently clear both. + // + // `keepLatest: 0` is the sharp one: 0 is the explicit "disabled" sentinel, and + // a serializer that tested the value for truthiness rather than for + // null/undefined would drop it and read as "inherit" on the next save. + await writeChannelConfig(SLUG, { + handling: "youtube", + name: "Test Title Filter", + url: "https://www.youtube.com/@example/videos", + downloadFilter: { include: "guest" }, + cookieMode: "always", + keepLatest: 0, + }); + const before = await readJson<Record<string, unknown>>(CONFIG); + expect(before.downloadFilter).toEqual({ include: "guest" }); + expect(before.cookieMode).toBe("always"); + expect(before.keepLatest).toBe(0); + + // THE SAME SENTENCE THE FORM SHOWS. title-filter.spec.ts reads this off the + // page after typing "elf(" into the exclude input; the route reaches it + // through the same parseChannelForm, which is the whole point of routing a + // patch through a FormData rather than writing config.json directly. + const bad = await ops(request, "channel-config", { + slug: SLUG, + patch: { downloadFilterExclude: "elf(" }, + }); + expect(bad.status).toBe(400); + expect(bad.body.error).toMatch( + /Download filter exclude is not a valid regular expression/, + ); + // And nothing was written. + expect( + (await readJson<Record<string, unknown>>(CONFIG)).downloadFilter, + ).toEqual({ include: "guest" }); + + const good = await ops(request, "channel-config", { + slug: SLUG, + patch: { downloadFilterInclude: "guest|special", downloadFilterExclude: "rerun" }, + }); + expect(good.body).toEqual({ ok: true }); + const after = await readJson<Record<string, unknown>>(CONFIG); + expect(after.downloadFilter).toEqual({ + include: "guest|special", + exclude: "rerun", + }); + // A PATCH IS A PATCH. updateChannelAction clears every form-managed key + // before layering the parse result on, so a route that posted only the patch + // would have silently dropped every one of these. + expect(after.name).toBe(before.name); + expect(after.handling).toBe(before.handling); + expect(after.url).toBe(before.url); + expect(after.cookieMode).toBe("always"); + expect(after.keepLatest).toBe(0); + + // "" clears a field, exactly as clearing the input does. + const cleared = await ops(request, "channel-config", { + slug: SLUG, + patch: { downloadFilterInclude: "", downloadFilterExclude: "" }, + }); + expect(cleared.body).toEqual({ ok: true }); + const emptied = await readJson<Record<string, unknown>>(CONFIG); + expect("downloadFilter" in emptied).toBe(false); + // Clearing one field clears ONLY that field. + expect(emptied.cookieMode).toBe("always"); + expect(emptied.keepLatest).toBe(0); + + // The read route sees the same config, and says whether the media is there. + await generateReport(page, SLUG); + const read = await request.get(`${baseUrl}/api/ops/channel/${SLUG}`, { + headers: AUTH, + }); + expect(read.status()).toBe(200); + const view = (await read.json()) as { + ok: boolean; + config: { name: string }; + media: { status: string }; + report: { totals: { videos: number } } | null; + }; + expect(view.ok).toBe(true); + expect(view.config.name).toBe(before.name); + expect(view.media.status).toBe("in-place"); + expect(view.report?.totals.videos).toBeGreaterThanOrEqual(0); + + const missing = await request.get(`${baseUrl}/api/ops/channel/nope`, { + headers: AUTH, + }); + expect(missing.status()).toBe(404); +}); + +test("channel-priority writes a per-operation override", async ({ request }) => { + await resetData("title-filter-channel"); + await settings(); + const SLUG = "test-filter"; + + const pinned = await ops(request, "channel-priority", { + slugs: [SLUG], + operation: "download", + tier: "paused", + }); + expect(pinned.body).toEqual({ ok: true }); + await expect + .poll(async () => { + const s = await readJson<{ + channelPriority?: { + channels?: Record<string, { overrides?: Record<string, string> }>; + }; + }>("test-settings.json"); + return s.channelPriority?.channels?.[SLUG]?.overrides?.download ?? null; + }) + .toBe("paused"); + + // null clears the override — and only with an operation named, because a bare + // null tier has no meaning for the BASE tier. + const bare = await ops(request, "channel-priority", { + slugs: [SLUG], + tier: null, + }); + expect(bare.status).toBe(400); + expect(bare.body.error).toMatch(/clears an operation override/); + + const cleared = await ops(request, "channel-priority", { + slugs: [SLUG], + operation: "download", + tier: null, + }); + expect(cleared.body).toEqual({ ok: true }); + await expect + .poll(async () => { + const s = await readJson<{ + channelPriority?: { + channels?: Record<string, { overrides?: Record<string, string> }>; + }; + }>("test-settings.json"); + return s.channelPriority?.channels?.[SLUG]?.overrides?.download ?? null; + }) + .toBe(null); + + const bogus = await ops(request, "channel-priority", { + slugs: [SLUG], + tier: "urgent", + }); + expect(bogus.status).toBe(400); + expect(bogus.body.error).toMatch(/normal, low, paused/); +}); + +test("metadata-scan starts a job, and the job says it is a metadata-scan", async ({ + request, +}) => { + test.setTimeout(120_000); + await resetData("title-filter-channel"); + await settings(); + const SLUG = "test-filter"; + + const { status, body } = await ops(request, "metadata-scan", { slug: SLUG }); + expect(status).toBe(200); + expect(body.ok).toBe(true); + expect(body.jobId).toBeTruthy(); + + // THE ROUTE DOES NOT STREAM, so the id is the whole contract: the caller + // follows the same log endpoint the editor's own panel polls. + // The sidecar is written with a `void` promise right after enqueue, so poll. + const metaPath = `test-transcripts/.jobs/${body.jobId}.meta.json`; + await expect + .poll(async () => { + const meta = await readJson<{ kind: string; channelSlug: string }>( + metaPath, + ).catch(() => null); + return meta ? `${meta.kind}/${meta.channelSlug}` : null; + }) + .toBe(`metadata-scan/${SLUG}`); + + const log = await request.get( + `${baseUrl}/api/jobs/${body.jobId}/log?from=0`, + ); + expect(log.status()).toBe(200); + const payload = (await log.json()) as { status: string }; + expect( + ["queued", "running", "done", "failed", "cancelled"].includes( + payload.status, + ), + ).toBe(true); + + const unknownChannel = await ops(request, "metadata-scan", { slug: "nope" }); + expect(unknownChannel.status).toBe(400); + expect(unknownChannel.body.error).toContain('Channel "nope" not found'); +}); + +test("refresh-report regenerates snapshot.json", async ({ request }) => { + test.setTimeout(120_000); + await resetData("title-filter-channel"); + await settings(); + const SLUG = "test-filter"; + const SNAP = `test-transcripts/channels/${SLUG}/snapshot.json`; + // Removed rather than assumed absent: a debounced regen armed by the previous + // spec can land between resetData's copy and here, and the point of the + // assertion below is the ROUTE's effect, not the scheduler's. + await rm(resolvePath(SNAP), { force: true }); + expect(await pathExists(SNAP)).toBe(false); + + // No page, no click: the route IS the refresh. It regenerates SYNCHRONOUSLY + // (a filesystem scan, not a job), so { ok: true } means the file is there. + const first = await ops(request, "refresh-report", { slug: SLUG }); + expect(first.body).toEqual({ ok: true }); + const snapshot = await readJson<{ generatedAt: string; totals: { videos: number } }>(SNAP); + expect(snapshot.generatedAt).toBeTruthy(); + expect(snapshot.totals.videos).toBeGreaterThanOrEqual(0); + + // Re-running REWRITES it. Polled through the action itself because two scans + // of a six-video fixture can land in the same millisecond. + await expect + .poll(async () => { + await ops(request, "refresh-report", { slug: SLUG }); + return (await readJson<{ generatedAt: string }>(SNAP)).generatedAt; + }) + .not.toBe(snapshot.generatedAt); + + // The bulk form queues a job per channel and reports both lists. + const all = await ops(request, "refresh-report", { all: true }); + expect(all.status).toBe(200); + expect(all.body.ok).toBe(true); + expect([...(all.body.queued ?? []), ...(all.body.skipped ?? []).map((s) => s.slug)]).toContain( + SLUG, + ); + + const both = await ops(request, "refresh-report", { slug: SLUG, all: true }); + expect(both.status).toBe(400); + expect(both.body.error).toMatch(/not both/); +}); + +test("relocate refuses a busy channel with the sentence the panel shows", async ({ + request, +}, testInfo) => { + test.setTimeout(120_000); + await resetData("slow-pipeline-channel"); + await settings(); + const SLUG = "slow-channel"; + + // --test-slow makes the fake yt-dlp sleep 30s, so the channel is genuinely + // busy for the length of this assertion rather than racily so. + const started = await ops(request, "sync", { slug: SLUG }); + expect(started.body.ok).toBe(true); + await expect + .poll(async () => { + const res = await request.get(`${baseUrl}/api/jobs/active`); + const body = (await res.json()) as + | { channelSlug?: string; status: string }[] + | { jobs?: { channelSlug?: string; status: string }[] }; + const jobs = Array.isArray(body) ? body : (body.jobs ?? []); + return jobs.filter( + (j) => + j.channelSlug === SLUG && + (j.status === "running" || j.status === "queued"), + ).length; + }) + .toBeGreaterThan(0); + + const root = testInfo.outputPath("media-root"); + const refused = await ops(request, "relocate", { slugs: [SLUG], root }); + expect(refused.status).toBe(200); + // A SKIP IS NOT A FAILURE — the bulk bar renders both numbers, and so does + // this. The reason is channelMediaBusyReason's, word for word. + expect(refused.body.queued).toEqual([]); + expect(refused.body.skipped?.[0]?.slug).toBe(SLUG); + expect(refused.body.skipped?.[0]?.reason).toMatch( + /running\/queued job\(s\) for this channel/, + ); + + // Exactly one destination, and it must be absolute. + const neither = await ops(request, "relocate", { slugs: [SLUG] }); + expect(neither.status).toBe(400); + expect(neither.body.error).toMatch(/exactly one of "locationId".*or "root"/); +}); + +test("lane flips a hold, and /api/auto-queue/status agrees", async ({ + request, +}) => { + await resetData("empty"); + await settings(); + + const held = async (): Promise<boolean> => { + const res = await request.get(`${baseUrl}/api/auto-queue/status`); + const body = (await res.json()) as Record<string, { held?: boolean }>; + return body.download?.held === true; + }; + expect(await held()).toBe(false); + + const hold = await ops(request, "lane", { lane: "download", held: true }); + expect(hold.body).toEqual({ ok: true, lane: "download" }); + await expect.poll(held).toBe(true); + + const release = await ops(request, "lane", { lane: "download", held: false }); + expect(release.body).toEqual({ ok: true, lane: "download" }); + await expect.poll(held).toBe(false); + + // `enabled` is the policy's master switch and is NOT the gate — two controls + // in the UI, two keys here. + const enable = await ops(request, "lane", { lane: "download", enabled: true }); + expect(enable.body.ok).toBe(true); + await expect + .poll(async () => { + const s = await readJson<{ + autoQueue?: Record<string, { enabled?: boolean }>; + }>("test-settings.json"); + return s.autoQueue?.download?.enabled ?? null; + }) + .toBe(true); + await ops(request, "lane", { lane: "download", enabled: false }); + + const nothing = await ops(request, "lane", { lane: "download" }); + expect(nothing.status).toBe(400); + expect(nothing.body.error).toMatch(/nothing to do/); + + const bogus = await ops(request, "lane", { lane: "transcode", held: true }); + expect(bogus.status).toBe(400); + expect(bogus.body.error).toMatch(/transcription, download, digest, backfill/); +}); diff --git a/package.json b/package.json @@ -22,7 +22,8 @@ "wt": "node scripts/worktree.mjs", "e2e:sharded": "node scripts/run-sharded-e2e.mjs", "test:scripts": "node --test scripts/*.test.mjs umtool/report-to-video/*.test.mjs", - "lint": "pnpm --filter export run lint" + "lint": "pnpm --filter export run lint", + "ops": "node scripts/archilyzer-ops.mjs" }, "devDependencies": { "tsx": "^4.21.0" diff --git a/plans/FACTS.md b/plans/FACTS.md @@ -4148,3 +4148,65 @@ never pulls `reader-fs.ts` into a client chunk. has no `runner`, so the operator presses Run. The obvious next step is for the download lane to run it for a channel whose filter has unscanned listed videos — which is exactly the backlog the snapshot already carries. + +--- + +## The ops API (`/api/ops/*`) — added 2026-09-20 + +**It is adapters, and nothing else.** Each route under +`editor/app/api/ops/<action>/route.ts` is ~5 lines: validate a JSON body, call +ONE existing server action, map its result. No route contains a rule the UI does +not already enforce — the point of the layer is that an agent over HTTP and an +operator clicking the same button get the same refusal, with the same sentence, +from the same code. A check written in a route would be a second opinion nobody +maintains. `editor/app/api/ops/_lib.ts` holds auth, body parsing and the three +result mappers (`jobResponse` / `actionResponse` / `queueResponse`). + +**`WORKER_TOKEN` is shared with `/api/worker/*` by design.** It already means +"this instance accepts instructions from something that is not the browser in +front of it", and the failure modes are identical, so the gate is: unset → 503 +(the surface is off until you opt in), wrong → 401. `authorizeWorkerRequest` +(`common/lib/workerToken.ts`) is the single implementation; nothing new was +written. A second secret would be a second thing to distribute, rotate and leave +unset. + +**A job-starting route returns `{ ok: true, jobId }` and NEVER streams.** +`runManagedFunction` hands back a `ReadableStream` the browser consumes; an HTTP +caller wants to hang up and poll. Every job adapter calls `result.stream.cancel()` +— which stops pushing into the controller and leaves the on-disk log running +(`streamCommand.ts`'s `cancel()` note) — and the caller follows +`/api/jobs/<id>/log`. Returning the id is also the honest answer: the queue may +hold the job behind other work for hours, so "started" is not "running". + +**Unknown body keys are a 400, never a silent ignore.** The allow-list passed to +`ops()` IS the route's documented body shape. A caller that misspells +`downloadFilterExclude` would otherwise get `{ ok: true }` and a channel that +still downloads everything. + +**`ops/channel-config`'s patch keys are the FORM's field names, not +`ChannelConfig`'s** — `downloadFilterInclude` / `downloadFilterExclude` rather +than a `downloadFilter` object, `ytdlpExtraArgs` as a string or an array of +lines. That is what routes them through `parseChannelForm`'s validators. The +patch is laid over the channel's CURRENT form representation +(`editor/app/channels/components/channelConfigToForm.ts`) rather than posted +alone, because `updateChannelAction` deletes every `CHANNEL_FORM_FIELDS` key from +the stored config before layering the parse result on — a FormData carrying only +a patch would clear everything the patch did not name. + +**`GET /api/ops/channel/<slug>`'s live counts are opt-in (`?counts=1`).** +`readChannelStat` walks every video directory — eleven thousand readdirs on the +largest channel here — so a poll loop asking for a channel's state would be the +thing hammering the platter. The report's totals answer the same question from a +file. The field is `counts` and is null without the flag. + +**No action needed a refactor to be callable from a route.** `revalidatePath` is +supported in Route Handlers (Next 16 — +`docs/01-app/03-api-reference/04-functions/revalidatePath.md:10`), no adapted +action calls `cookies()`, and the only actions that `redirect()` +(`createChannelAction`, `gotoVideoAction`) are deliberately not exposed. + +**The 503-when-unset branch is not covered by e2e.** The editor test server runs +with `WORKER_TOKEN=test-worker-token` in `editor/package.json`'s `dev:test`, one +server for the whole suite, so no spec can observe the disabled state. +`editor/e2e/ops-api.spec.ts` covers 401 (missing and wrong); the 503 is +`authorizeWorkerRequest`'s own first branch, shared with `/api/worker/*`. diff --git a/scripts/archilyzer-ops.mjs b/scripts/archilyzer-ops.mjs @@ -0,0 +1,245 @@ +#!/usr/bin/env node +// archilyzer-ops — drive a running editor over HTTP, without a browser. +// +// Every editor gesture used to be reachable only as a server action, which meant +// an agent that wanted to sync a channel or fix a download filter had to drive +// Playwright. /api/ops is a thin adapter layer over those same actions, and this +// is its client. +// +// USAGE +// +// pnpm ops <action> [--json '<body>'] [--wait] [--quiet] +// pnpm ops get channel <slug> [--counts] +// pnpm ops list +// +// ARCHILYZER_EDITOR_URL editor base URL (default http://localhost:3001) +// WORKER_TOKEN the shared secret the editor is running with. +// Unset on the SERVER => every route 503s; unset here +// => every route 401s. +// +// EXAMPLES +// +// pnpm ops sync --json '{"slug":"the-quartering"}' --wait +// pnpm ops metadata-scan --json '{"slug":"the-quartering"}' +// pnpm ops channel-config --json '{"slug":"x","patch":{"downloadFilterExclude":"rerun"}}' +// pnpm ops channel-priority --json '{"slugs":["x"],"operation":"download","tier":"paused"}' +// 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 get channel the-quartering +// +// --wait follows /api/jobs/<jobId>/log to the end for a job-starting action and +// exits 0 only if the job finished `done`. Without it the command returns as +// soon as the job is QUEUED, which is the honest answer: the queue may hold it +// behind other work for hours. +// +// The response JSON is printed verbatim on stdout (log lines from --wait go to +// stderr), so `pnpm ops … | jq` works. + +const DEFAULT_URL = "http://localhost:3001"; + +// The read-side routes, reachable as `get <noun> <arg>`. Kept tiny and explicit: +// an ops API that let a caller assemble arbitrary GET paths would be a proxy, +// not an adapter. +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) => + `/api/ops/channel/${encodeURIComponent(slug)}${counts ? "?counts=1" : ""}`, +}; + +const ACTIONS = [ + "channel-priority", + "channel-config", + "metadata-scan", + "import-video", + "refresh-report", + "sync", + "download-missing", + "retry-bucket", + "build-index", + "build-deploy", + "build-site", + "relocate", + "relocate-back", + "lane", +]; + +export function parseArgs(argv) { + const positional = []; + let json = null; + let wait = false; + let quiet = false; + let counts = false; + for (let i = 0; i < argv.length; i++) { + const arg = argv[i]; + if (arg === "--wait") { + wait = true; + } else if (arg === "--quiet") { + quiet = true; + } else if (arg === "--counts") { + counts = true; + } else if (arg === "--json") { + json = argv[++i]; + if (json === undefined) { + return { error: "--json needs a JSON object argument" }; + } + } else if (arg.startsWith("--json=")) { + json = arg.slice("--json=".length); + } else if (arg === "--help" || arg === "-h") { + return { help: true }; + } else if (arg.startsWith("-")) { + return { error: `unknown flag: ${arg}` }; + } else { + positional.push(arg); + } + } + if (positional.length === 0) return { help: true }; + let body = {}; + if (json !== null) { + try { + body = JSON.parse(json); + } catch (e) { + return { error: `--json is not valid JSON: ${e.message}` }; + } + if (typeof body !== "object" || body === null || Array.isArray(body)) { + return { error: "--json must be a JSON object" }; + } + } + if (positional[0] === "list") { + return { list: true }; + } + if (positional[0] === "get") { + const noun = positional[1]; + if (!noun || !GETTERS[noun]) { + return { + error: `get: unknown noun "${noun ?? ""}" — known: ${Object.keys(GETTERS).join(", ")}`, + }; + } + if (!positional[2]) return { error: `get ${noun}: needs an argument` }; + return { + method: "GET", + path: GETTERS[noun](positional[2], counts), + wait: false, + quiet, + }; + } + const action = positional[0]; + if (!ACTIONS.includes(action)) { + return { + error: `unknown action "${action}" — known: ${ACTIONS.join(", ")}`, + }; + } + if (positional.length > 1) { + return { + error: `"${action}" takes no positional arguments — pass its body with --json`, + }; + } + return { method: "POST", path: `/api/ops/${action}`, body, wait, quiet }; +} + +export function usage() { + return [ + "Usage: pnpm ops <action> [--json '<body>'] [--wait]", + " pnpm ops get channel <slug> [--counts]", + " pnpm ops list", + "", + `Actions: ${ACTIONS.join(", ")}`, + "", + "Env: ARCHILYZER_EDITOR_URL (default http://localhost:3001), WORKER_TOKEN", + ].join("\n"); +} + +function baseUrl() { + return (process.env.ARCHILYZER_EDITOR_URL ?? DEFAULT_URL).replace(/\/+$/, ""); +} + +function authHeaders() { + const token = process.env.WORKER_TOKEN ?? ""; + return token ? { authorization: `Bearer ${token}` } : {}; +} + +// 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. +async function followJob(jobId, quiet) { + let from = 0; + for (;;) { + const res = await fetch( + `${baseUrl()}/api/jobs/${encodeURIComponent(jobId)}/log?from=${from}`, + { headers: authHeaders() }, + ); + if (!res.ok) throw new Error(`log poll failed: HTTP ${res.status}`); + const payload = await res.json(); + if (payload.content && !quiet) process.stderr.write(payload.content); + from = payload.nextOffset ?? from; + const status = payload.status; + if (status !== "queued" && status !== "running") return status; + await new Promise((r) => setTimeout(r, 1000)); + } +} + +async function main() { + const parsed = parseArgs(process.argv.slice(2)); + if (parsed.help) { + console.log(usage()); + return 0; + } + if (parsed.error) { + console.error(parsed.error); + console.error(""); + console.error(usage()); + return 2; + } + if (parsed.list) { + console.log(ACTIONS.join("\n")); + return 0; + } + const url = `${baseUrl()}${parsed.path}`; + const res = await fetch(url, { + method: parsed.method, + headers: { + ...authHeaders(), + ...(parsed.method === "POST" ? { "content-type": "application/json" } : {}), + }, + ...(parsed.method === "POST" ? { body: JSON.stringify(parsed.body) } : {}), + }); + const text = await res.text(); + let payload; + try { + payload = JSON.parse(text); + } catch { + console.error(`HTTP ${res.status}: ${text.slice(0, 500)}`); + return 1; + } + console.log(JSON.stringify(payload, null, 2)); + if (!res.ok || payload.ok === false) return 1; + if (!parsed.wait) return 0; + const jobIds = payload.jobId + ? [payload.jobId] + : Array.isArray(payload.jobs) + ? payload.jobs.map((j) => j.jobId) + : []; + if (jobIds.length === 0) { + // Not a job-starting action (or it queued nothing). --wait is satisfied. + return 0; + } + let worst = 0; + for (const jobId of jobIds) { + const status = await followJob(jobId, parsed.quiet); + console.error(`[${jobId}] ${status}`); + if (status !== "done") worst = 1; + } + return worst; +} + +// Importable for the arg-parsing tests; only the CLI entry point runs main(). +if (process.argv[1] && import.meta.url === `file://${process.argv[1]}`) { + main().then( + (code) => process.exit(code), + (e) => { + console.error(e.message); + process.exit(1); + }, + ); +} diff --git a/scripts/archilyzer-ops.test.mjs b/scripts/archilyzer-ops.test.mjs @@ -0,0 +1,72 @@ +// Arg parsing for scripts/archilyzer-ops.mjs. No network: parseArgs is pure and +// returns the request it WOULD make, which is the whole surface worth pinning — +// the routes themselves are covered by editor/e2e/ops-api.spec.ts. +// +// Run with: pnpm test:scripts +import assert from "node:assert/strict"; +import test from "node:test"; +import { parseArgs, usage } from "./archilyzer-ops.mjs"; + +test("no arguments prints usage", () => { + assert.equal(parseArgs([]).help, true); + assert.match(usage(), /pnpm ops <action>/); +}); + +test("an action becomes a POST to its route", () => { + const p = parseArgs(["sync", "--json", '{"slug":"x","full":true}']); + assert.equal(p.method, "POST"); + assert.equal(p.path, "/api/ops/sync"); + assert.deepEqual(p.body, { slug: "x", full: true }); + assert.equal(p.wait, false); +}); + +test("--wait and --json= are both accepted", () => { + const p = parseArgs(["metadata-scan", '--json={"slug":"x"}', "--wait"]); + assert.equal(p.wait, true); + assert.deepEqual(p.body, { slug: "x" }); +}); + +test("an action with no body posts an empty object", () => { + const p = parseArgs(["build-index"]); + assert.deepEqual(p.body, {}); +}); + +test("an unknown action is refused by name, with the list", () => { + const p = parseArgs(["sinc"]); + assert.match(p.error, /unknown action "sinc"/); + assert.match(p.error, /metadata-scan/); +}); + +test("malformed --json is refused before any request", () => { + assert.match(parseArgs(["sync", "--json", "{"]).error, /not valid JSON/); + assert.match(parseArgs(["sync", "--json", "[1]"]).error, /must be a JSON object/); + assert.match(parseArgs(["sync", "--json"]).error, /needs a JSON object/); +}); + +test("positional arguments after an action are refused", () => { + // `pnpm ops sync the-quartering` reads naturally and would otherwise be a + // silent no-op body, so it is an error that names the fix. + assert.match(parseArgs(["sync", "the-quartering"]).error, /--json/); +}); + +test("get channel becomes a GET on the read route", () => { + const p = parseArgs(["get", "channel", "the quartering"]); + assert.equal(p.method, "GET"); + assert.equal(p.path, "/api/ops/channel/the%20quartering"); +}); + +test("get refuses an unknown noun and a missing argument", () => { + assert.match(parseArgs(["get", "site", "x"]).error, /unknown noun/); + assert.match(parseArgs(["get", "channel"]).error, /needs an argument/); +}); + +test("an unknown flag is refused", () => { + assert.match(parseArgs(["sync", "--force"]).error, /unknown flag/); +}); + +test("get channel --counts asks for the live on-disk counts", () => { + const p = parseArgs(["get", "channel", "x", "--counts"]); + assert.equal(p.path, "/api/ops/channel/x?counts=1"); + // Off by default: the counts walk every video directory. + assert.equal(parseArgs(["get", "channel", "x"]).path, "/api/ops/channel/x"); +});