// `archilyzer run [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. // // - it does not see the editor's lanes either: a `run` over a channel the // editor's lane is working on at the same moment does the same videos // twice (the writes are atomic, so it is wasted work, not damage). Run it // when that lane is paused or elsewhere. // // Ctrl-C cancels the job the way the editor's Cancel does. Exit: 0 done; // 1 failed or refused (an operation this process will not run, one switched // off, a paused lane, unreachable media, no given id on disk); 2 usage (an // unknown operation or channel, a stray flag); 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 { 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 1; } 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 === ids.length) { out.error( `run: none of the ${ids.length} id(s) has a data// on disk (${absent.join(", ")}) — nothing to run`, ); return 1; } if (absent.length > 0) { out.error( `run: ${absent.length} of ${ids.length} id(s) have no data// 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; }