Archilyzer · Source

archilyzer

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

commit a568fb57b7a0aafd712634db6e01dbb5094e2c97
parent bc891e659538c8459bb72d313d08077e19efed7a
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Mon, 28 Sep 2026 01:57:05 -0400

common: `archilyzer run`, `archilyzer mcp`, and every bin a subcommand

`run <operation> <channel> [ids…] [--lane local|remote]` runs one catalogued
operation over a channel offline through `runOperationChannelJob`, the body
the /channels buttons, stage cards and lane runners call: a real
`backfill-channel` / `digest-channel-*` job with its record and log under
`.jobs/`, the same summary line, the channel snapshot flushed at the end,
and `runManagedFunction`'s media guard. Diarization, both attribution
passes and digest run. Sync, the metadata scan and download are refused
(the per-platform download queue paces every request to one source inside
the editor), and transcription is refused (the editor's worker pool); the
sentence is derived from the descriptor. A switched-off operation, a paused
lane (the batch would idle-wait for a resume only the editor can give), an
unknown channel and --lane off digest are refused up front; ids with no
data dir are named. Ctrl-C cancels the job; exit 0/1/2/130.

`mcp [args…]` starts the MCP server as `pnpm --filter yt-dlp-transcript-mcp
exec tsx src/index.ts` does — mcp's tsx, cwd mcp/, env passed through, a
child rather than an import (mcp depends on common) — and prints nothing on
stdout, which is the JSON-RPC channel.

`_cli.ts` gains PASSTHROUGH commands: matched on the leading words only,
handed everything after their path verbatim. Every remaining bin is a
one-line row: `build stats|templates|archives` and `docs files` in-process;
`duplicates` (8 GB heap, as its script had), `posts fetch|check`,
`diarize backfill`, `digest plan|validate`, `reconcile video-dirs`,
`verify transcripts`, `transcribe` (transform.ts), `migrate
channel-priority` as children with their own flags (`_spawnBin.ts`);
`brand media` in-process with its argv.

Tests: run-operation.test.ts (6) over a temp corpus, mcp.test.ts (1: a real
initialize handshake through the CLI), _cli.test.ts +3 (passthrough, the
leading-words rule, every bin in common/bin reachable from a row).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

Diffstat:
Mcommon/bin/_cli.test.ts | 63++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcommon/bin/_cli.ts | 38+++++++++++++++++++++++++++++++++++++-
Acommon/bin/_spawnBin.ts | 81+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/bin/archilyzer.ts | 103+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Acommon/bin/mcp.test.ts | 54++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/bin/mcp.ts | 23+++++++++++++++++++++++
Acommon/bin/run-operation.test.ts | 187+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/bin/run-operation.ts | 175+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
8 files changed, 720 insertions(+), 4 deletions(-)

diff --git a/common/bin/_cli.test.ts b/common/bin/_cli.test.ts @@ -1,9 +1,13 @@ import { test } from "node:test"; import assert from "node:assert/strict"; +import { existsSync, readdirSync, readFileSync } from "node:fs"; +import path from "node:path"; +import { fileURLToPath } from "node:url"; import { parseArgv } from "./_parseFlags"; import { argumentProblem, booleanFlags, + passthroughCommand, resolveCommand, runCli, usage, @@ -195,7 +199,7 @@ test("parseCutArgs checks the changelog, the version and the date before anythin test("a refused release command exits 2 and prints why, touching nothing", async () => { const errors: string[] = []; const out = { log: () => {}, error: (s: string) => errors.push(s) }; - const ctx = (positionals: string[]) => ({ positionals, flags: {}, env: {} }); + const ctx = (positionals: string[]) => ({ positionals, flags: {}, env: {}, argv: [] }); assert.equal(await cutMain(ctx(["site", "next"]), out), 2); assert.equal(await showMain(ctx(["hub"]), out), 2); assert.equal(errors.length, 2); @@ -309,3 +313,60 @@ test("release show says the latest release, its date and what is pending, and wh assert.equal(formatAllLine([editor, exp]), "all: next 0.9.3, next-minor 0.10.0"); assert.equal(formatAllLine([editor]), null); }); + +// ── passthrough commands (one-core Phase 4 slice 3) ───────────────────────── + +test("a passthrough command gets every word after its path, verbatim and unchecked", async () => { + let seen: { argv: string[]; positionals: string[] } | null = null; + const table = [ + cmd(["posts", "fetch"], { + passthrough: true, + run: async (ctx) => { + seen = { argv: ctx.argv, positionals: ctx.positionals }; + return 7; + }, + }), + cmd(["posts"], { maxPositionals: 1 }), + ]; + const quiet = { log: () => {}, error: () => {} }; + assert.equal( + await runCli(table, ["posts", "fetch", "--slug", "x", "--full", "--", "--weird=1"], {}, quiet), + 7, + ); + assert.deepEqual(seen, { argv: ["--slug", "x", "--full", "--", "--weird=1"], positionals: [] }); + // `--help` as the first word after the path is the usage line, not the bin's. + const logged: string[] = []; + assert.equal( + await runCli(table, ["posts", "fetch", "--help"], {}, { log: (s) => logged.push(s), error: () => {} }), + 0, + ); + assert.match(logged.join("\n"), /archilyzer posts fetch/); +}); + +test("a passthrough path matches only the LEADING words, never a flag's value", () => { + const table = [cmd(["mcp"], { passthrough: true }), cmd(["run"], { maxPositionals: 9 })]; + assert.equal(passthroughCommand(table, ["run", "digest", "--lane", "mcp"]), null); + assert.deepEqual(passthroughCommand(table, ["mcp", "--local", "x"])?.path, ["mcp"]); + assert.equal(passthroughCommand(table, ["--local", "mcp"]), null); +}); + +test("every bin in common/bin is reachable as a subcommand", () => { + const here = path.dirname(fileURLToPath(import.meta.url)); + const table = readFileSync(path.join(here, "archilyzer.ts"), "utf8"); + const reached = new Set([ + ...[...table.matchAll(/import\("\.\/([\w-]+)"\)/g)].map((m) => m[1]), + ...[...table.matchAll(/script\(\[[^\]]*\],\s*"([\w-]+)\.ts"/g)].map((m) => m[1]), + ]); + const bins = readdirSync(here) + .filter((n) => n.endsWith(".ts") && !n.endsWith(".test.ts") && !n.startsWith("_")) + .map((n) => n.slice(0, -3)) + .filter((n) => n !== "archilyzer"); + assert.deepEqual(bins.filter((b) => !reached.has(b)), []); + // And every script row names a file that exists. + for (const m of table.matchAll(/script\(\[[^\]]*\],\s*"([\w-]+\.ts)"/g)) { + assert.ok(existsSync(path.join(here, m[1])), m[1]); + } + // No two rows share a path. + const paths = COMMANDS.map((c) => c.path.join(" ")); + assert.equal(new Set(paths).size, paths.length); +}); diff --git a/common/bin/_cli.ts b/common/bin/_cli.ts @@ -16,6 +16,9 @@ export type CommandContext = { positionals: string[]; flags: Record<string, FlagValue>; env: NodeJS.ProcessEnv; + // A passthrough command's words after its path, verbatim (every other + // command gets []). + argv: string[]; }; export type Command = { @@ -28,6 +31,12 @@ export type Command = { flags?: Record<string, FlagKind>; // At most this many positionals after the path (default 0). maxPositionals?: number; + // The command parses its own arguments (a bin with its own flags, the MCP + // server): runCli hands it everything after its path, untouched, as `argv`, + // and checks nothing. Matched on the LEADING words of argv only, so a flag + // value can never be mistaken for its path. `--help` as the first word after + // the path still prints the usage line. + passthrough?: boolean; // The exit code. run: (ctx: CommandContext) => Promise<number>; }; @@ -115,6 +124,15 @@ export async function runCli( env: NodeJS.ProcessEnv = process.env, out: { log: (s: string) => void; error: (s: string) => void } = console, ): Promise<number> { + const through = passthroughCommand(table, argv); + if (through) { + const rest = argv.slice(through.path.length); + if (rest[0] === "--help" || rest[0] === "-h") { + out.log(usage([through])); + return 0; + } + return through.run({ positionals: [], flags: {}, env, argv: rest }); + } const { flags, positionals } = parseArgv(argv, booleanFlags(table)); if (positionals.length === 0) { (flags.help ? out.log : out.error)(usage(table)); @@ -134,7 +152,25 @@ export async function runCli( out.error(`${problem}\n\n${usage([hit.command])}`); return 2; } - return hit.command.run({ positionals: hit.rest, flags, env }); + return hit.command.run({ positionals: hit.rest, flags, env, argv: [] }); +} + +/** + * The passthrough command whose path is the longest run of LEADING words of + * argv — or null. Only the leading words: `run digest --lane x` must never be + * read as some command named by a flag's value. + */ +export function passthroughCommand( + table: readonly Command[], + argv: readonly string[], +): Command | null { + let best: Command | null = null; + for (const c of table) { + if (!c.passthrough || c.path.length > argv.length) continue; + if (!c.path.every((w, i) => argv[i] === w)) continue; + if (!best || c.path.length > best.path.length) best = c; + } + return best; } /** diff --git a/common/bin/_spawnBin.ts b/common/bin/_spawnBin.ts @@ -0,0 +1,81 @@ +// Run a TypeScript entry point as a CHILD process under tsx, stdio inherited, +// and resolve with its exit code. For the CLI rows whose target parses its own +// argv at import time (the `_parseFlags` bins) or is another package's entry +// (the MCP server): importing either into the CLI's process would hand it the +// CLI's argv, or make common depend on a package that depends on common. +// +// Signals: a terminal's Ctrl-C reaches the child directly (same process +// group), so SIGINT is only held here — forwarding it too would deliver it +// twice, which a bin that treats a second Ctrl-C as "force" would honour. +// SIGTERM and SIGHUP come from a supervisor (an MCP client stopping its +// server, say) and are forwarded. + +import { spawn } from "node:child_process"; +import path from "node:path"; +import { fileURLToPath } from "node:url"; + +const COMMON = path.resolve(path.dirname(fileURLToPath(import.meta.url)), ".."); +export const REPO_ROOT = path.resolve(COMMON, ".."); + +/** The tsx that a package's own scripts run (its devDependency). */ +export function tsxFor(packageDir: string): string { + return path.join(packageDir, "node_modules", ".bin", "tsx"); +} + +export function runChild(opts: { + command: string; + args: string[]; + cwd?: string; + env?: NodeJS.ProcessEnv; +}): Promise<number> { + return new Promise((resolve) => { + const child = spawn(opts.command, opts.args, { + cwd: opts.cwd, + env: opts.env ?? process.env, + stdio: "inherit", + }); + const hold = () => {}; + const forward = (sig: NodeJS.Signals) => () => { + if (!child.killed) child.kill(sig); + }; + const onTerm = forward("SIGTERM"); + const onHup = forward("SIGHUP"); + process.on("SIGINT", hold); + process.on("SIGTERM", onTerm); + process.on("SIGHUP", onHup); + const done = (code: number) => { + process.off("SIGINT", hold); + process.off("SIGTERM", onTerm); + process.off("SIGHUP", onHup); + resolve(code); + }; + child.on("error", (err) => { + console.error(`${path.basename(opts.command)}: ${err.message}`); + done(1); + }); + child.on("exit", (code, signal) => { + done(code ?? (signal ? 128 + (signalNumber(signal) ?? 0) : 1)); + }); + }); +} + +function signalNumber(sig: NodeJS.Signals): number | undefined { + return ({ SIGHUP: 1, SIGINT: 2, SIGKILL: 9, SIGTERM: 15 } as Record<string, number>)[sig]; +} + +/** + * A common/bin script that reads process.argv itself, run with `argv` exactly + * as given. `heapMb` is the NODE_OPTIONS heap the old package.json script gave + * it (the corpus-wide passes need more than node's default). + */ +export function runBinScript(file: string, argv: string[], heapMb?: number): Promise<number> { + const env = { ...process.env }; + if (heapMb) { + env.NODE_OPTIONS = `${env.NODE_OPTIONS ?? ""} --max-old-space-size=${heapMb}`.trim(); + } + return runChild({ + command: tsxFor(COMMON), + args: [path.join(COMMON, "bin", file), ...argv], + env, + }); +} diff --git a/common/bin/archilyzer.ts b/common/bin/archilyzer.ts @@ -7,8 +7,6 @@ // Every row imports its implementation LAZILY, so `archilyzer sync tick` never // loads the AWS SDK and `archilyzer settings example` never opens LMDB. The // machinery (parser, lookup, usage) is `_cli.ts`; this file is only the table. -// -// Not here yet (one-core Phase 4 slice 3): `doctor`, `run <operation>`, `mcp`. import type { Command } from "./_cli"; import { runCli, runIfEntryPoint } from "./_cli"; @@ -23,6 +21,30 @@ export const COMMANDS: Command[] = [ }, }, { + path: ["build", "stats"], + usage: "rebuild the stats datasets (reads the index; the data phase's second step)", + run: async () => { + await (await import("./build-stats")).main(); + return 0; + }, + }, + { + path: ["build", "templates"], + usage: "bake each site's chart templates into its export staging dir (the data phase's third step)", + run: async () => { + await (await import("./build-chart-templates")).main(); + return 0; + }, + }, + { + path: ["build", "archives"], + usage: "warm the shared archive-zip cache for every enabled site's channels, once", + run: async () => { + await (await import("./build-archives")).main(); + return 0; + }, + }, + { path: ["compose", "site"], usage: "<id> compose one site's export/public (default: SITE_ID)", maxPositionals: 1, @@ -172,6 +194,65 @@ export const COMMANDS: Command[] = [ }, }, { + path: ["run"], + usage: + "<operation> <channel> [ids…] [--lane local|remote] run one catalogued operation over a channel offline, as the editor's job does (sync, downloads and transcription are refused: they run in the editor)", + flags: { lane: "string" }, + maxPositionals: Number.MAX_SAFE_INTEGER, + run: async ({ positionals, flags }) => { + const [operation, channel, ...ids] = positionals; + if (!operation || !channel) { + console.error("run: which operation, over which channel? `archilyzer run <operation> <channel> [ids…]`"); + return 2; + } + const lane = flags.lane; + if (lane !== undefined && lane !== "local" && lane !== "remote") { + console.error(`run: --lane is local or remote, not "${String(lane)}"`); + return 2; + } + const { runOperation } = await import("./run-operation"); + return runOperation( + { operation, channel, ids, lane: lane as "local" | "remote" | undefined }, + { signal: interrupted() }, + ); + }, + }, + // The bins that parse their own flags, run as children with their argv + // verbatim (_spawnBin.ts says why). Each row is the whole integration. + script(["duplicates"], "duplicate-shorts.ts", + "[--threshold N] [--all-durations] [--blocking title|duration|both] [--near F] [--tolerance N] … on-demand duplicate detection (after index + stats)", 8192), + script(["posts", "fetch"], "fetch-posts.ts", + "--slug <channel> [--full] [--limit N] fetch a social channel's posts into its posts corpus"), + script(["posts", "check"], "check-post-availability.ts", + "--slug <channel> [--mode stale|unchecked|all] [--limit N] which archived posts were deleted at the source"), + script(["diarize", "backfill"], "diarize-backfill.ts", + "[--dry-run] [--scope transcribed|channel:<slug>|video:<slug>/<id>] [--limit N] [--force] … diarize videos whose audio is still on disk"), + script(["digest", "plan"], "digest-plan.ts", + "[--lane local|remote] [--channels a,b] [--top N] [--json] [--census] … price the digest backfill; writes nothing"), + script(["digest", "validate"], "digest-validate.ts", + "<channel> [<channel> …] score digests already on disk"), + script(["reconcile", "video-dirs"], "reconcile-video-dirs.ts", + "[--channel <slug>] [--dry-run] [--verbose] rename video dirs to the canonical id layout"), + script(["verify", "transcripts"], "verify-transcripts.ts", + "--channel <slug> list duplicate and missing transcripts"), + script(["transcribe"], "transform.ts", + "--channel <slug> whisper over a channel's downloaded audio with THIS process's own worker pool — never beside a running editor"), + script(["migrate", "channel-priority"], "migrate-channel-priority.ts", + "[--dry-run] the one-shot channel-priority migration (plans/channel-priority.md, S5)"), + { + path: ["brand", "media"], + usage: "[--out <dir>] [--video-kit] render the Archilyzer Media channel's assets", + passthrough: true, + run: async ({ argv }) => (await import("./brand-media")).main(argv), + }, + { + path: ["mcp"], + usage: + "[--local <dir>|--remote <url>|--hub <url>] start the MCP server on stdio (as `pnpm --filter yt-dlp-transcript-mcp exec tsx src/index.ts`)", + passthrough: true, + run: async ({ argv }) => (await import("./mcp")).main(argv), + }, + { path: ["sync", "tick"], usage: "POST one scheduler tick to the editor (SYNC_TICK_URL, SYNC_TICK_TOKEN)", run: async () => (await import("./sync-tick")).tick(), @@ -184,6 +265,13 @@ export const COMMANDS: Command[] = [ (await import("./env-docs")).main({ check: flags.check === true }), }, { + path: ["docs", "files"], + usage: "[--check] write SITE.md + CHANNEL.md from the file schemas", + flags: { check: "boolean" }, + run: async ({ flags }) => + (await import("./file-schemas-docs")).main({ check: flags.check === true }), + }, + { path: ["settings", "example"], usage: "[--check] write settings.json.example + SETTINGS.md from the schema", flags: { check: "boolean" }, @@ -216,6 +304,17 @@ export const COMMANDS: Command[] = [ }, ]; +// A row for a bin that parses its own argv: a passthrough command that runs +// `common/bin/<file>` as a child with the words after its path. +function script(path: string[], file: string, usage: string, heapMb?: number): Command { + return { + path, + usage, + passthrough: true, + run: async ({ argv }) => (await import("./_spawnBin")).runBinScript(file, argv, heapMb), + }; +} + // The site a site command names: its argument, else SITE_ID (which is how // export's `build` / `deploy` scripts are called). Prints and returns null when // there is neither. diff --git a/common/bin/mcp.test.ts b/common/bin/mcp.test.ts @@ -0,0 +1,54 @@ +// `archilyzer mcp` starts the real MCP server, and stdout stays JSON-RPC. +// +// Run with: +// pnpm --filter yt-dlp-transcript-common test +// +// Spawns the CLI the way `claude mcp add … -- <command>` would, sends the 2025 +// initialize request, and reads the answer. The arguments after `mcp` reach the +// server untouched (the `--local` dir comes back in its stderr banner), and the +// FIRST line on stdout is the server's JSON-RPC answer: anything the CLI printed +// before it would break every MCP client. + +import { test, after } from "node:test"; +import assert from "node:assert/strict"; +import { spawn } from "node:child_process"; +import { mkdtempSync, rmSync } from "node:fs"; +import os from "node:os"; +import path from "node:path"; +import { fileURLToPath } from "node:url"; + +const HERE = path.dirname(fileURLToPath(import.meta.url)); +const TSX = path.join(HERE, "..", "node_modules", ".bin", "tsx"); +const DIR = mkdtempSync(path.join(os.tmpdir(), "archilyzer-mcp-")); +after(() => rmSync(DIR, { recursive: true, force: true })); + +test("archilyzer mcp answers initialize on stdout, with its argv passed through", async () => { + const child = spawn(TSX, [path.join(HERE, "archilyzer.ts"), "mcp", "--local", DIR], { + stdio: ["pipe", "pipe", "pipe"], + }); + let stdout = ""; + let stderr = ""; + child.stdout.on("data", (b) => (stdout += b)); + child.stderr.on("data", (b) => (stderr += b)); + child.stdin.write( + `${JSON.stringify({ + jsonrpc: "2.0", + id: 1, + method: "initialize", + params: { protocolVersion: "2025-06-18", capabilities: {}, clientInfo: { name: "test", version: "0" } }, + })}\n`, + ); + const deadline = Date.now() + 30_000; + while (!stdout.includes("\n") && Date.now() < deadline) { + await new Promise((r) => setTimeout(r, 50)); + } + child.stdin.end(); + const code = await new Promise<number | null>((r) => child.on("exit", (c) => r(c))); + const first = stdout.split("\n")[0]; + const msg = JSON.parse(first) as { id: number; result?: { serverInfo?: { name?: string } } }; + assert.equal(msg.id, 1); + assert.equal(msg.result?.serverInfo?.name, "yt-dlp-transcript-mcp"); + assert.match(stderr, new RegExp(`default corpus: local:${DIR.replace(/[.*+?^${}()|[\]\\]/g, "\\$&")}`)); + // stdin closed: the server ends, and the CLI ends with it. + assert.equal(code, 0); +}); diff --git a/common/bin/mcp.ts b/common/bin/mcp.ts @@ -0,0 +1,23 @@ +// `archilyzer mcp [args…]` — start the MCP server, exactly as +// `pnpm --filter yt-dlp-transcript-mcp exec tsx src/index.ts [args…]` does: +// mcp's own tsx, cwd mcp/, the environment passed through untouched (the +// server reads TRANSCRIPT_SITE_URL / TRANSCRIPT_HUB_URL / TRANSCRIPT_LOCAL_DIR, +// ARCHILYZER_EDITOR_URL and WORKER_TOKEN from it), stdio inherited. +// +// A CHILD, never an import: the mcp package depends on common, so common +// importing mcp would make the dependency a cycle. And NOTHING may be printed +// to stdout on the way in — stdout is the JSON-RPC channel. + +import { existsSync } from "node:fs"; +import path from "node:path"; +import { REPO_ROOT, runChild, tsxFor } from "./_spawnBin"; + +export async function main(argv: string[], root = REPO_ROOT): Promise<number> { + const dir = path.join(root, "mcp"); + const entry = path.join(dir, "src", "index.ts"); + if (!existsSync(entry)) { + console.error(`mcp: no MCP server at ${entry}`); + return 1; + } + return runChild({ command: tsxFor(dir), args: ["src/index.ts", ...argv], cwd: dir }); +} diff --git a/common/bin/run-operation.test.ts b/common/bin/run-operation.test.ts @@ -0,0 +1,187 @@ +// `archilyzer run <operation> <channel> [ids…]` over a temp corpus. +// +// Run with: +// pnpm --filter yt-dlp-transcript-common test +// +// Its own file for the SETTINGS SEAM, like operationBatchRelocation.test.ts: +// getPaths() memoizes its first answer, so TRANSCRIPTS_DIR and SETTINGS_FILE are +// set before anything imports it. Settings are re-read on every call, so each +// test writes the file it needs. +// +// The operation that RUNS here is diarization over videos that have a +// transcript and no audio: each classifies `missing-input` (with re-download +// off), which is counted and never dispatched — no engine, no model, no +// network — and still goes through the whole job: the record, the log, the +// batch, the summary line and the snapshot. + +import { mkdtempSync, readdirSync, readFileSync, writeFileSync } from "node:fs"; +import { mkdir, rm, symlink, writeFile } from "node:fs/promises"; +import os from "node:os"; +import path from "node:path"; +import { test, after } from "node:test"; +import assert from "node:assert/strict"; + +const ROOT = mkdtempSync(path.join(os.tmpdir(), "run-operation-")); +process.env.TRANSCRIPTS_DIR = ROOT; +const SETTINGS_FILE = path.join(ROOT, "settings.json"); +process.env.SETTINGS_FILE = SETTINGS_FILE; + +const { runOperation, offlineRefusal } = await import("./run-operation"); +const { getPaths } = await import("../lib/paths"); +const { operationCatalog } = await import("../lib/operations"); + +after(() => rm(ROOT, { recursive: true, force: true })); + +function settings(over: Record<string, unknown> = {}): void { + writeFileSync( + SETTINGS_FILE, + JSON.stringify({ + autoQueue: { backfill: { enabled: true, held: false } }, + backfill: { allowRedownload: false, concurrency: 1 }, + diarization: { enabled: true, segModel: "/models/seg.onnx", embModel: "/models/emb.onnx" }, + attribution: { enabled: false, diarizedEnabled: false, textOnlyEnabled: false }, + ...over, + }), + ); +} + +async function seed(slug: string, ids: string[]): Promise<void> { + const channelDir = path.join(getPaths().channelsDir, slug); + await mkdir(channelDir, { recursive: true }); + await writeFile( + path.join(channelDir, "config.json"), + JSON.stringify({ handling: "transcribe", url: "https://example.com/c" }), + ); + for (const id of ids) { + const videoDir = path.join(channelDir, "data", id); + await mkdir(videoDir, { recursive: true }); + await writeFile( + path.join(videoDir, "transcript.json"), + JSON.stringify({ transcription: [{ text: "hello", offsets: {} }] }), + ); + } +} + +function capture() { + const o = { stdout: "", stderr: "" }; + return { + o, + out: { + log: (s: string) => (o.stdout += `${s}\n`), + error: (s: string) => (o.stderr += `${s}\n`), + write: (s: string) => (o.stdout += s), + }, + }; +} + +const jobRecords = () => { + try { + return readdirSync(getPaths().jobsDir).filter((n) => n.endsWith(".json")); + } catch { + return []; + } +}; + +test("an unknown operation is refused with the list", async () => { + const { o, out } = capture(); + const code = await runOperation({ operation: "diarisation", channel: "x", ids: [] }, { out }); + assert.equal(code, 2); + assert.match(o.stderr, /no operation "diarisation"/); + assert.match(o.stderr, /Runs here: .*diarization.*digest/); + assert.match(o.stderr, /run only by the editor: .*sync.*transcription/); +}); + +test("sync, the scan, downloads and transcription are refused with a sentence, not half-run", async () => { + for (const id of ["sync", "metadata-scan", "download"]) { + const { o, out } = capture(); + assert.equal(await runOperation({ operation: id, channel: "x", ids: [] }, { out }), 2, id); + assert.match(o.stderr, /download queue inside the editor/, id); + } + const { o, out } = capture(); + assert.equal(await runOperation({ operation: "transcription", channel: "x", ids: [] }, { out }), 2); + assert.match(o.stderr, /worker pool/); + // Every catalogued operation has an answer, and exactly the registry's run. + const runs = operationCatalog().filter((op) => offlineRefusal(op) === null).map((op) => op.id); + assert.deepEqual(runs.sort(), ["attribution-diarized", "attribution-text", "diarization", "digest"]); +}); + +test("an unknown channel, a switched-off operation, a paused lane and a stray --lane are refused", async () => { + settings(); + await seed("known", ["v1"]); + let c = capture(); + assert.equal(await runOperation({ operation: "diarization", channel: "nope", ids: [] }, { out: c.out }), 2); + assert.match(c.o.stderr, /no channel "nope"/); + + c = capture(); + assert.equal(await runOperation({ operation: "attribution-text", channel: "known", ids: [] }, { out: c.out }), 1); + assert.match(c.o.stderr, /switched off in settings\.json/); + + settings({ autoQueue: { backfill: { enabled: true, held: true } } }); + c = capture(); + assert.equal(await runOperation({ operation: "diarization", channel: "known", ids: [] }, { out: c.out }), 1); + assert.match(c.o.stderr, /backfill lane is paused/); + + settings(); + c = capture(); + assert.equal( + await runOperation({ operation: "diarization", channel: "known", ids: [], lane: "local" }, { out: c.out }), + 2, + ); + assert.match(c.o.stderr, /--lane is the digest engine lane/); +}); + +test("the media guard refuses an unmounted channel before any job exists", async () => { + settings(); + const slug = "moved"; + const channelDir = path.join(getPaths().channelsDir, slug); + await mkdir(channelDir, { recursive: true }); + const target = path.join(ROOT, "platter", slug, "data"); // never created + await writeFile( + path.join(channelDir, "config.json"), + JSON.stringify({ handling: "transcribe", url: "https://example.com/m", dataDir: target }), + ); + await symlink(target, path.join(channelDir, "data")); + const before = jobRecords(); + const { o, out } = capture(); + const code = await runOperation({ operation: "diarization", channel: slug, ids: [] }, { out }); + assert.equal(code, 1); + assert.match(o.stderr, /moved/); + assert.match(o.stderr, /not mounted|unreachable|does not exist/); + assert.deepEqual(jobRecords(), before, "no job record for a refused run"); +}); + +test("diarization runs through the editor's job body: record, log, summary, snapshot", async () => { + settings(); + await seed("chan", ["v1", "v2"]); + const before = new Set(jobRecords()); + const { o, out } = capture(); + const code = await runOperation({ operation: "diarization", channel: "chan", ids: [] }, { out }); + assert.equal(code, 0, o.stderr + o.stdout); + assert.match(o.stdout, /Backfill chan: 0 done, 0 already current, 0 failed; 2 still need their media re-acquired/); + const added = jobRecords().filter((n) => !before.has(n)); + assert.equal(added.length >= 1, true, "a job record was written"); + const records = added.map((n) => JSON.parse(readFileSync(path.join(getPaths().jobsDir, n), "utf8"))); + const job = records.find((r) => r.kind === "backfill-channel"); + assert.ok(job, JSON.stringify(records)); + assert.equal(job.channelSlug, "chan"); + assert.equal(job.status, "done"); + assert.deepEqual(job.spec?.params?.kindIds, ["diarization"]); + // The report the editor would have refreshed at job end. + assert.ok( + readdirSync(path.join(getPaths().channelsDir, "chan")).some((n) => n.startsWith("snapshot")), + "the channel snapshot was regenerated", + ); +}); + +test("ids scope the run, and ids with no data dir are named", async () => { + settings(); + await seed("scoped", ["v1", "v2", "v3"]); + const { o, out } = capture(); + const code = await runOperation( + { operation: "diarization", channel: "scoped", ids: ["v2", "gone", "v2"] }, + { out }, + ); + assert.equal(code, 0, o.stderr); + assert.match(o.stderr, /1 of 2 id\(s\) have no data\/<id>\/ on disk and are skipped: gone/); + assert.match(o.stdout, /Backfill scoped: 0 done, 0 already current, 0 failed; 1 still need their media re-acquired/); +}); diff --git a/common/bin/run-operation.ts b/common/bin/run-operation.ts @@ -0,0 +1,175 @@ +// `archilyzer run <operation> <channel> [ids…]` — one catalogued operation over +// one channel, offline, through the body the editor's job runs. +// +// THE SAME BODY: `runOperationChannelJob` (controller/operationJobs.ts), which +// the /channels group buttons, the stage cards and the lane runners call. So a +// run here is a real job — kind `backfill-channel` or `digest-channel-*`, a +// record and a log under `transcripts/.jobs/`, the same summary line, and the +// media guard: `runManagedFunction` refuses a `needsMedia` kind for a channel +// whose media is unreachable, before any record is made. +// +// WHAT RUNS HERE AND WHAT DOES NOT, decided from the descriptor, not the id: +// - the registry operations (dispatch "backfill": diarization, both +// attribution passes, digest) run — their engines are subprocesses and +// HTTP endpoints this process can reach as well as the editor can; +// - an `external` operation is refused with a sentence. Sync, the metadata +// scan and downloads run on the per-platform download queue, which paces +// every request to one source INSIDE the editor; a second process would +// fetch beside it. Transcription is dispatched across the editor's worker +// pool (leases, tiers, drain), which does not exist out here. +// +// WHAT IS DIFFERENT FROM THE EDITOR'S RUN, deliberately and in the open: +// - it does not see the editor's jobs: the backfill lane's yield-to- +// transcription reads THIS process's transcription activity, which is none; +// - a paused lane is REFUSED up front rather than held: the batch would +// idle-wait for a resume that only the editor can give. +// +// Ctrl-C cancels the job the way the editor's Cancel does. Exit: 0 done, +// 1 failed or refused, 2 usage, 130 cancelled. + +import { existsSync } from "node:fs"; +import path from "node:path"; +import type { Paths } from "../lib/paths"; +import type { OperationDescriptor } from "../lib/operations"; + +export type RunOut = { log: (s: string) => void; error: (s: string) => void; write: (s: string) => void }; + +const defaultOut: RunOut = { + log: (s) => console.log(s), + error: (s) => console.error(s), + write: (s) => { + process.stdout.write(s); + }, +}; + +// Why an operation cannot run outside the editor, as one sentence — or null +// when it can. Read off the descriptor, so a new external entry gets an honest +// refusal without an edit here. +export function offlineRefusal(op: OperationDescriptor): string | null { + if (op.dispatch !== "external") return null; + if (op.lane.contendsFor === "network") { + return ( + `${op.label} runs on the per-platform download queue inside the editor, which paces ` + + `every request to one source; a second process would fetch beside it and defeat that ` + + `pacing. Run it from the editor, or through \`pnpm ops\`.` + ); + } + if (op.runner === "transcription") { + return ( + `${op.label} is dispatched across the editor's worker pool (leases, tiers, drain), ` + + `which does not exist outside the editor. Run it from the editor.` + ); + } + return `${op.label} is dispatched by the editor, not through the operation registry.`; +} + +export type RunOperationArgs = { + operation: string; + channel: string; + ids: string[]; + // Digest only: which engine lane. Absent = the configured one. + lane?: "local" | "remote"; +}; + +export async function runOperation( + args: RunOperationArgs, + deps: { paths?: Paths; out?: RunOut; signal?: AbortSignal } = {}, +): Promise<number> { + const out = deps.out ?? defaultOut; + const { operationCatalog, getOperation, DIGEST_OPERATION_ID } = await import("../lib/operations"); + const catalog = operationCatalog(); + const op = catalog.find((o) => o.id === args.operation); + if (!op) { + const runs = catalog.filter((o) => !offlineRefusal(o)).map((o) => o.id); + const refused = catalog.filter((o) => offlineRefusal(o)).map((o) => o.id); + out.error( + `run: no operation "${args.operation}". Runs here: ${runs.join(", ")}. ` + + `Catalogued but run only by the editor: ${refused.join(", ")}.`, + ); + return 2; + } + const refusal = offlineRefusal(op); + if (refusal) { + out.error(`run: ${refusal}`); + return 2; + } + if (args.lane && op.id !== DIGEST_OPERATION_ID) { + out.error(`run: --lane is the digest engine lane; ${op.label} has none`); + return 2; + } + + const paths = deps.paths ?? (await import("../lib/paths")).getPaths(); + const channelDir = path.join(paths.channelsDir, args.channel); + if (!existsSync(path.join(channelDir, "config.json"))) { + out.error(`run: no channel "${args.channel}" (no ${path.join(channelDir, "config.json")})`); + return 2; + } + + const { getSettings } = await import("../lib/settings"); + const settings = getSettings(); + const registered = getOperation(op.id); + if (!registered || !registered.enabled(settings)) { + out.error( + `run: ${op.label} is switched off in settings.json` + + (op.settingsBlock ? ` (its \`${op.settingsBlock}\` block)` : "") + + ` — the editor's lane would not run it either.`, + ); + return 1; + } + const { pauseLaneFor, isGateHeld } = await import("../lib/pauseGates"); + const gate = pauseLaneFor(op.id); + if (gate && isGateHeld(settings, gate)) { + out.error( + `run: the ${gate} lane is paused (autoQueue.${gate}.held) — the editor's run would hold ` + + `until it is resumed. Resume it on /operations/${op.id} first.`, + ); + return 1; + } + + const ids = [...new Set(args.ids.map((s) => s.trim()).filter(Boolean))]; + if (ids.length > 0) { + const absent = ids.filter((id) => !existsSync(path.join(channelDir, "data", id))); + if (absent.length > 0) { + out.error( + `run: ${absent.length} of ${ids.length} id(s) have no data/<id>/ on disk and are skipped: ${absent.join(", ")}`, + ); + } + } + + const { runOperationChannelJob } = await import("../controller/operationJobs"); + const result = await runOperationChannelJob({ + paths, + channelSlug: args.channel, + operation: op.id, + ids: ids.length > 0 ? ids : undefined, + ...(args.lane ? { digest: { lane: args.lane } } : {}), + }); + if (!result.ok) { + out.error(`run: ${result.error}`); + return 1; + } + out.error(`run: ${op.label} over ${args.channel}${ids.length ? ` (${ids.length} id(s))` : ""} — job ${result.jobId}`); + + const { getRegistry } = await import("../jobs/registry"); + const cancel = () => { + getRegistry().cancel(result.jobId); + }; + deps.signal?.addEventListener("abort", cancel, { once: true }); + const reader = result.stream.getReader(); + for (;;) { + const { value, done } = await reader.read(); + if (done) break; + out.write(value); + } + const { status } = await result.done; + deps.signal?.removeEventListener("abort", cancel); + + // The editor refreshes the channel's report after a job; a process that is + // about to exit has to ask for it now, or the debounce would never fire. + const { flushChannelSnapshotsNow } = await import("../jobs/snapshotScheduler"); + await flushChannelSnapshotsNow(); + + if (status === "done") return 0; + if (status === "cancelled") return 130; + return 1; +}