Archilyzer · Source

archilyzer

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

commit 5e4d31824f4dbabbc19322d419bbd69d5c1f16b3
parent d6d99ff8c4fd42920ae526065ba3b8ae3ae52b78
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Mon, 24 Aug 2026 10:07:40 -0400

unit executor: any backfill kind can run on a corpus-less box

Every backfill kind now declares its unit contract — inputs(files) (ordered,
transcript.cues.json LAST so the executor's materialization passes the
isCuesJsonFresh mtime gate), outputs, and applyResult (the primary-side apply
through the guarded writers: attribution re-checks the downgrade rule against
current disk, digest lands section-wise through writeDigestSection, only
diarization.json copies verbatim). BackfillRunOptions carries the injected
identity (attributionSettings / diarizationSettings / digestContext) that
each kind's run() threads into the controllers' existing-but-dead settings
overrides; attributeOneVideo gains the context param digestVideo already had.

workerServer gains startWorkerUnit (scratch corpus under
.worker-scratch/<job>/channels, cues written last, "disabled"/"not-configured"
treated as an injection-dropped ERROR) and readWorkerUnitResult (the
generalisation of the hardcoded transcript.json read to kind.outputs), behind
new token-guarded /api/worker/unit routes cloned from the transcribe set —
which refuse anything that is not a backfill kind, keeping download politeness
and the transcription pool single-scheduler. Executor deployment documented in
RUNNING_IN_DOCKER.md: this app + ARCHILYZER_IDLE_BOOT=1 + WORKER_TOKEN + an
empty transcripts dir.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

Diffstat:
MRUNNING_IN_DOCKER.md | 34++++++++++++++++++++++++++++++++++
Mcommon/controller/attributeOne.ts | 11+++++++++--
Mcommon/controller/workerServer.ts | 186++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcommon/lib/backfillKinds.ts | 231+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Acommon/lib/backfillUnit.test.ts | 228+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/CHANGELOG.md | 1+
Aeditor/app/api/worker/unit/[id]/events/route.ts | 26++++++++++++++++++++++++++
Aeditor/app/api/worker/unit/[id]/result/route.ts | 31+++++++++++++++++++++++++++++++
Aeditor/app/api/worker/unit/[id]/route.ts | 30++++++++++++++++++++++++++++++
Aeditor/app/api/worker/unit/route.ts | 51+++++++++++++++++++++++++++++++++++++++++++++++++++
10 files changed, 822 insertions(+), 7 deletions(-)

diff --git a/RUNNING_IN_DOCKER.md b/RUNNING_IN_DOCKER.md @@ -240,6 +240,40 @@ Boots the server with all of that stopped. You can start any of it from the UI afterwards. The shutdown reaper and the persisted-pause restore stay armed either way — both only ever *stop* work. +### Running as a unit executor for another machine + +A second machine can take whole work *units* (diarization, speaker +attribution, CPU transcription) from a primary instance without holding any +corpus at all. Two shapes, by cost: + +- **LLM calls only** (digest / attribution — the bulk of any backlog): run + nothing but `ollama serve` with the primary's exact model tag pulled, and + register the box on the primary as an **LLM endpoint** worker (Settings → + Transcription workers). No repo, no container, no token. The primary + verifies the model tag before use and refuses an endpoint that lacks it — + freshness pins the model identity, so "almost the right model" would write + permanently-stale records. +- **A unit executor** (everything else): boot this app with + + ```sh + ARCHILYZER_IDLE_BOOT=1 # never arm runners or resume sweeps here + WORKER_TOKEN=<shared secret> # enables /api/worker/*; off without it + TRANSCRIPTS_DIR=/some/empty/dir + SETTINGS_FILE=/some/where/settings.json # pin it — the cwd fallback is a trap off-repo + ``` + + and register it on the primary as a **Remote** worker with matching token. + Units arrive with their inputs, run against a scratch corpus under + `TRANSCRIPTS_DIR/.worker-scratch/`, and the primary pulls the produced + sidecars back and applies them through its own guarded writers. The + executor's own settings are never consulted for the work's identity — the + primary injects its model/prompt configuration into every unit, and a unit + that would fall back to local defaults fails loudly instead. Tag the worker + (e.g. `cpu, diarization`) to say what it should take. + +The executor never opens the primary's LMDB, never downloads media, and never +arms a sweep — the primary stays the sole scheduler. + ### GPU transcription Two overlays, because two different engines get the GPU. Neither is required — diff --git a/common/controller/attributeOne.ts b/common/controller/attributeOne.ts @@ -64,7 +64,10 @@ import { } from "../lib/digestPrompt"; import { chunkCuesForContext } from "../lib/transcriptWindow"; import { transcriptToMarkdown } from "../lib/transcriptToMarkdown"; -import { readDigestContext } from "../lib/digestContext-server"; +import { + readDigestContext, + type DigestContext, +} from "../lib/digestContext-server"; import type { Cue } from "../lib/vtt"; import { isCuesJsonFresh, @@ -88,6 +91,9 @@ export type AttributeOneOptions = { // unit executor, which passes the primary's injected identity config so a // bare box's default settings can never leak into provenance. appConfig?: DigestAppConfig; + // Pre-read channel context (note + hash). A unit executor injects the + // primary's, since its scratch corpus has no digest-context.md to read. + context?: DigestContext; // Redo even when the recorded identity matches. Never overrides the DOWNGRADE // rule — forcing a regeneration is not the same as asking for a worse record, // and nothing in the UI should be able to request the second by accident. @@ -167,7 +173,8 @@ export async function attributeOneVideo( return "already-exists"; } - const context = await readDigestContext(opts.paths, opts.channelSlug); + const context = + opts.context ?? (await readDigestContext(opts.paths, opts.channelSlug)); const warnings: AttributionWarning[] = []; const started = Date.now(); let speakers: AttributionSpeaker[]; diff --git a/common/controller/workerServer.ts b/common/controller/workerServer.ts @@ -7,12 +7,23 @@ import path from "node:path"; import fs from "fs-extra"; -import { getPaths } from "../lib/paths"; +import { getPaths, type Paths } from "../lib/paths"; import { getRegistry, newJobId, type JobStatus } from "../jobs/registry"; import { getWorkerPool } from "../jobs/workerPool"; import { runManagedFunction } from "../jobs/streamCommand"; import { makeTaskTracker } from "../jobs/taskHooks"; import { transcribeWithWorker } from "./transcribeOne"; +import { + getBackfillKind, + type BackfillRunOutcome, +} from "../lib/backfillKinds"; +import { CUES_JSON_FILENAME } from "../lib/videoStatus"; +import type { + AttributionSettings, + DiarizationSettings, +} from "../lib/settings"; +import type { DigestAppConfig } from "../lib/digest"; +import type { DigestContext } from "../lib/digestContext-server"; type WorkerJob = { remoteJobId: string; @@ -25,6 +36,9 @@ type WorkerJob = { // Byte offset already streamed to the requester, so /events can return only // new log output on each poll. logOffset: number; + // Set for a UNIT job (a backfill kind run against a scratch corpus). The op + // decides which output files the result endpoint reads back. + unit?: { op: string; videoDir: string }; }; const ID_RE = /^[A-Za-z0-9_-]+$/; @@ -170,6 +184,176 @@ export async function startWorkerTranscriptionShared( return { remoteJobId }; } +// --------------------------------------------------------------------------- +// Unit jobs: run ONE backfill kind (attribution, diarization, …) against a +// scratch corpus materialized from the uploaded inputs. The generalisation of +// the transcription protocol above: same auth, same poll/pull/DELETE shape, +// but the payload is a backfill kind's declared inputs and the result its +// declared outputs. The primary applies those outputs through the kinds' +// guarded writers — the executor never touches a real corpus. +// --------------------------------------------------------------------------- + +const UNIT_OUTCOME_FILENAME = "unit-outcome.json"; + +export type StartWorkerUnitInput = { + // A backfill kind id. `download`/`transcription` are ExternalOperations, not + // kinds, so getBackfillKind() refusing them is the door guard that keeps + // download politeness (and the transcription pool) single-scheduler. + op: string; + channelSlug: string; + videoId: string; + // Input files by (base)name, base64-encoded. Written in arrival order with + // transcript.cues.json FORCED LAST: isCuesJsonFresh compares mtimes, and a + // cues file older than its metadata/raw transcript makes the run skip with + // stale-cues — a unit that silently does nothing. + files: Record<string, string>; + // The primary's resolved freshness target, verbatim. + target: unknown; + force?: boolean; + // The primary's IDENTITY config, verbatim — model, promptVersion, numCtx, + // context — with any baseUrl stripped so the executor localises its own + // endpoint. Mandatory in spirit: a kind that answers "disabled" or + // "not-configured" here is treated as an ERROR, because on a bare box that + // answer means the injection was dropped and default settings leaked in. + config?: { + attribution?: AttributionSettings; + diarization?: DiarizationSettings; + appConfig?: DigestAppConfig; + context?: DigestContext; + }; +}; + +export async function startWorkerUnit( + input: StartWorkerUnitInput, +): Promise<{ remoteJobId: string }> { + const kind = getBackfillKind(input.op); + if (!kind) { + throw new Error( + `operation "${input.op}" is not a backfill kind this executor can run`, + ); + } + if (!ID_RE.test(input.channelSlug) || !ID_RE.test(input.videoId)) { + throw new Error("invalid channelSlug or videoId"); + } + const paths = getPaths(); + const remoteJobId = newJobId(); + const scratchRoot = path.join(paths.workerScratchDir, remoteJobId); + const scratchChannels = path.join(scratchRoot, "channels"); + const videoDir = path.join( + scratchChannels, + input.channelSlug, + "data", + input.videoId, + ); + await fs.ensureDir(videoDir); + // Materialize the inputs, cues LAST (see StartWorkerUnitInput.files). The + // names cross the trust boundary — basename + allowlist, like uploads. + const names = Object.keys(input.files ?? {}); + const ordered = [ + ...names.filter((n) => n !== CUES_JSON_FILENAME), + ...names.filter((n) => n === CUES_JSON_FILENAME), + ]; + for (const name of ordered) { + const safe = sanitizeAudioName(name); + await fs.writeFile( + path.join(videoDir, safe), + Buffer.from(input.files[name], "base64"), + ); + } + // The kind runs against a SCRATCH corpus root — the real channelsDir is + // never touched, and a bare executor has none anyway. + const scratchPaths: Paths = { ...paths, channelsDir: scratchChannels }; + const res = await runManagedFunction({ + kind: "worker-unit", + queueKey: "", // immediate; the primary's scheduler already throttled this + paths, + // Deliberately NO channelSlug: the completion hook would request a channel + // snapshot for a channel this box does not have. + videoId: input.videoId, + fn: async (onLog, signal) => { + const outcome: BackfillRunOutcome = await kind.run({ + paths: scratchPaths, + videoDir, + videoId: input.videoId, + channelSlug: input.channelSlug, + target: input.target, + force: input.force, + attributionSettings: input.config?.attribution, + diarizationSettings: input.config?.diarization, + appConfig: input.config?.appConfig, + digestContext: input.config?.context, + onLog, + signal, + }); + if (outcome === "disabled" || outcome === "not-configured") { + throw new Error( + `unit ${input.op} reported "${outcome}" on the executor — the config ` + + `injection was dropped and this box's default settings leaked in`, + ); + } + if (outcome === "failed") { + throw new Error(`unit ${input.op} failed for ${input.videoId}`); + } + await fs.writeFile( + path.join(scratchRoot, UNIT_OUTCOME_FILENAME), + JSON.stringify({ outcome }), + ); + }, + }); + if (!res.ok) { + await fs.remove(scratchRoot).catch(() => {}); + throw new Error(res.error); + } + res.stream.cancel().catch(() => {}); + jobs().set(remoteJobId, { + remoteJobId, + workDir: scratchRoot, + inPlace: false, + managedJobId: res.jobId, + logOffset: 0, + unit: { op: input.op, videoDir }, + }); + return { remoteJobId }; +} + +export type WorkerUnitResult = { + // The kind's BackfillRunOutcome ("done", "already-present", "skipped", …). + outcome: string; + // The kind's declared output files, as UTF-8 JSON text by name. May be empty + // for a no-op outcome. + files: Record<string, string>; +}; + +// The generalisation of readWorkerResult: instead of one hardcoded +// transcript.json, the unit's kind declares its outputs. null until the run +// has written its outcome marker. +export async function readWorkerUnitResult( + remoteJobId: string, +): Promise<WorkerUnitResult | null> { + const job = jobs().get(remoteJobId); + if (!job?.unit) return null; + const outcomeRaw = await fs + .readFile(path.join(job.workDir, UNIT_OUTCOME_FILENAME), "utf8") + .catch(() => null); + if (outcomeRaw === null) return null; + let outcome = "done"; + try { + const parsed = JSON.parse(outcomeRaw) as { outcome?: string }; + if (typeof parsed.outcome === "string") outcome = parsed.outcome; + } catch { + // keep "done" — the marker only exists after a successful run + } + const kind = getBackfillKind(job.unit.op); + const files: Record<string, string> = {}; + for (const name of kind?.outputs ?? []) { + const p = path.join(job.unit.videoDir, name); + if (await fs.pathExists(p)) { + files[name] = await fs.readFile(p, "utf8"); + } + } + return { outcome, files }; +} + export type WorkerJobEvents = { status: JobStatus | "unknown"; fraction?: number; diff --git a/common/lib/backfillKinds.ts b/common/lib/backfillKinds.ts @@ -58,7 +58,11 @@ // into an existing pass and `lane` gets the concurrent queue and the share. import type { Paths } from "./paths"; -import type { SiteSettings } from "./settings"; +import type { + AttributionSettings, + DiarizationSettings, + SiteSettings, +} from "./settings"; import { BACKFILL_QUEUE, DIGEST_LOCAL_QUEUE, @@ -66,13 +70,20 @@ import { TRANSCRIPTION_QUEUE, } from "./queueKeys"; import { + DIGEST_FILENAME, + isDigestSectionKind, isSectionFresh, type DigestAppConfig, type DigestFreshnessTarget, + type DigestItem, type DigestLane, + type DigestProvenance, + type DigestRecord, type DigestSectionKind, + type DigestWarning, } from "./digest"; -import { loadDigest } from "./digest-server"; +import { loadDigest, writeDigestSection } from "./digest-server"; +import type { DigestContext } from "./digestContext-server"; import { SORTFORMER_DIARIZATION_ENGINE, diarizationTarget, @@ -80,22 +91,29 @@ import { type DiarizationBackend, type DiarizationEngineId, type DiarizationFreshnessTarget, + type DiarizationRecord, } from "./diarization"; -import { loadDiarization } from "./diarization-server"; +import { loadDiarization, writeDiarization } from "./diarization-server"; import { ATTRIBUTION_FILENAME, isAttributionDowngrade, isAttributionFresh, type AttributionFreshnessTarget, type AttributionMethod, + type AttributionRecord, } from "./attribution"; -import { loadAttribution } from "./attribution-server"; +import { loadAttribution, writeAttribution } from "./attribution-server"; import { + CUES_JSON_FILENAME, + DIARIZATION_FILENAME, + META_FILENAME, + WHISPER_FILENAME, findSourceMedia, isVideoTranscribed, readVideoDurationSec, type VideoFiles, } from "./videoStatus"; +import { pickPreferredAudio } from "./mediaFiles"; import { SAVED_VIDEO_POINTER_FILENAME } from "./savedVideo"; import { resolveSavedVideo } from "./savedVideo-server"; import { diarizeOneVideo } from "../controller/diarizeOne"; @@ -230,6 +248,19 @@ export type BackfillRunOptions = { // which endpoint served a call can vary freely while everything that IS // identity (model, numCtx, timeouts) travels verbatim from the primary. appConfig?: DigestAppConfig; + // Per-kind settings injection for a UNIT EXECUTOR (see workerServer's + // startWorkerUnit). The primary resolves these and ships them in the unit + // envelope; each kind's run() threads its own field into the controller's + // existing `settings` override — so a bare executor's default settings + // (attribution disabled, empty app config = a DIFFERENT identity) can never + // leak into provenance. On the primary these stay unset and the controllers + // read live settings exactly as before. + attributionSettings?: AttributionSettings; + diarizationSettings?: DiarizationSettings; + // Channel context, injected rather than shipped as a file: the primary + // already computed the note + hash, and digest pins contextHash in its + // identity — injecting the object makes hash equality true by construction. + digestContext?: DigestContext; onLog?: (msg: string) => void; signal?: AbortSignal; }; @@ -378,14 +409,172 @@ export type BackfillKind = { resolveTarget(ctx: BackfillTargetContext): unknown | Promise<unknown>; state(probe: BackfillProbe): Promise<BackfillClassification>; run(opts: BackfillRunOptions): Promise<BackfillRunOutcome>; + + // --- The unit-executor contract: what one unit of this kind needs, what it + // produces, and how its output lands back on the primary. --- + + // Filenames (within the video dir) a unit executor must be shipped, derived + // from the listing the caller already has. ORDERED: the executor + // materializes them in this order, and every kind lists + // transcript.cues.json LAST — isCuesJsonFresh compares mtimes, and a cues + // file written before its metadata/raw transcript reads as stale, making + // the unit silently do nothing. + inputs(files: VideoFiles): string[]; + // Filenames the run writes — what the executor's result endpoint returns. + outputs: readonly string[]; + // Apply a unit's returned output files ON THE PRIMARY, through the guarded + // writers — never a raw file copy. Both sidecar writers are + // read-modify-write (writeDigestSection preserves the other section; + // attribution re-checks the downgrade rule against the primary's CURRENT + // disk, because the unit ran against a snapshot that is minutes old). + // "invalid" = the payload is not a usable record; "refused" = a guard said + // no (which is a success of the guard, not a failure of the unit). + applyResult( + videoDir: string, + payload: Record<string, string>, + ): Promise<ApplyResultOutcome>; }; +export type ApplyResultOutcome = "applied" | "refused" | "invalid"; + export type BackfillTargetContext = { settings: SiteSettings; paths: Paths; channelSlug: string; }; +// --------------------------------------------------------------------------- +// Unit-contract helpers, shared across the kinds. +// --------------------------------------------------------------------------- + +// The resolved primary raw transcript — needed by every transcript-derived +// kind because isCuesJsonFresh compares the cues sidecar's mtime against it. +function rawTranscriptOf(files: VideoFiles): string | null { + if (files.hasWhisper) return WHISPER_FILENAME; + return files.ytVttFile; +} + +// The common transcript-derived input set: metadata + raw transcript + the +// existing output sidecar (for the freshness/downgrade checks) + CUES LAST — +// see BackfillKind.inputs for why the order is load-bearing. +function transcriptUnitInputs( + files: VideoFiles, + extras: readonly string[], +): string[] { + const out: string[] = []; + if (files.entries.includes(META_FILENAME)) out.push(META_FILENAME); + const raw = rawTranscriptOf(files); + if (raw) out.push(raw); + for (const name of extras) { + if (files.entries.includes(name)) out.push(name); + } + if (files.entries.includes(CUES_JSON_FILENAME)) out.push(CUES_JSON_FILENAME); + return out; +} + +// Shared by both attribution kinds: parse, validate the same shape +// loadAttribution enforces, RE-CHECK the downgrade rule against the primary's +// current disk (the unit ran against a snapshot minutes old, and the diarized +// lane can have landed a better record meanwhile), then the guarded writer. +async function applyAttributionResult( + videoDir: string, + payload: Record<string, string>, +): Promise<ApplyResultOutcome> { + const raw = payload[ATTRIBUTION_FILENAME]; + if (!raw) return "invalid"; + let record: AttributionRecord; + try { + const parsed = JSON.parse(raw) as Partial<AttributionRecord>; + if ( + typeof parsed?.videoId !== "string" || + typeof parsed.generatedAt !== "string" || + !Array.isArray(parsed.speakers) || + !Array.isArray(parsed.segments) || + !parsed.provenance || + typeof parsed.provenance.method !== "string" + ) { + return "invalid"; + } + record = parsed as AttributionRecord; + } catch { + return "invalid"; + } + if ( + isAttributionDowngrade( + await loadAttribution(videoDir), + record.provenance.method as AttributionMethod, + ) + ) { + return "refused"; + } + await writeAttribution(videoDir, record); + return "applied"; +} + +// Digest results land SECTION-WISE through writeDigestSection, which preserves +// the other section, its warnings and the history from the record already on +// the primary's disk — a raw copy of the unit's file would clobber a section +// the primary wrote while the unit was in flight. +async function applyDigestResult( + videoDir: string, + payload: Record<string, string>, +): Promise<ApplyResultOutcome> { + const raw = payload[DIGEST_FILENAME]; + if (!raw) return "invalid"; + let parsed: Partial<DigestRecord>; + try { + parsed = JSON.parse(raw) as Partial<DigestRecord>; + } catch { + return "invalid"; + } + const sections = parsed?.sections; + if (!sections || typeof sections !== "object") return "invalid"; + let applied = false; + for (const [section, value] of Object.entries(sections)) { + if (!isDigestSectionKind(section)) continue; + const v = value as { + provenance?: DigestProvenance; + items?: DigestItem[]; + }; + if (!v?.provenance || !Array.isArray(v.items)) continue; + const warnings: DigestWarning[] = (parsed.warnings ?? []).filter( + (w) => w.section === section, + ); + await writeDigestSection(videoDir, { + section, + items: v.items, + provenance: v.provenance, + warnings, + }); + applied = true; + } + return applied ? "applied" : "invalid"; +} + +// Diarization is the ONE whole-file verbatim apply: the sidecar has a single +// writer and no sections to merge, so the shape check plus the atomic writer +// is the whole guard. +async function applyDiarizationResult( + videoDir: string, + payload: Record<string, string>, +): Promise<ApplyResultOutcome> { + const raw = payload[DIARIZATION_FILENAME]; + if (!raw) return "invalid"; + try { + const parsed = JSON.parse(raw) as Partial<DiarizationRecord>; + if ( + typeof parsed?.generatedAt !== "string" || + !Array.isArray(parsed.turns) + ) { + return "invalid"; + } + await writeDiarization(videoDir, parsed as DiarizationRecord); + return "applied"; + } catch { + return "invalid"; + } +} + // Diarization: the first entry, and the reason the table exists. const diarization: BackfillKind = { id: "diarization", @@ -457,6 +646,7 @@ const diarization: BackfillKind = { paths: opts.paths, videoDir: opts.videoDir, videoId: opts.videoId, + settings: opts.diarizationSettings, force: opts.force, onLog: opts.onLog, signal: opts.signal, @@ -468,6 +658,19 @@ const diarization: BackfillKind = { if (outcome === "disabled") return "disabled"; return "failed"; }, + // The one kind whose input is MEDIA: the preferred extracted audio, else the + // persisted source container. Same envelope as the transcript kinds, bigger + // files — and a reachable population of only ~836 videos, so modest use. + inputs(files) { + const out: string[] = []; + if (files.entries.includes(META_FILENAME)) out.push(META_FILENAME); + const audio = + pickPreferredAudio(files.audioFiles) ?? findSourceMedia(files.entries); + if (audio) out.push(audio); + return out; + }, + outputs: [DIARIZATION_FILENAME], + applyResult: applyDiarizationResult, }; // Is there anything on disk ffmpeg could read for this video? Mirrors @@ -606,12 +809,17 @@ const attributionText: BackfillKind = { channelSlug: opts.channelSlug, method: "text-only", force: opts.force, + settings: opts.attributionSettings, appConfig: opts.appConfig, + context: opts.digestContext, onLog: opts.onLog, signal: opts.signal, }), ); }, + inputs: (files) => transcriptUnitInputs(files, [ATTRIBUTION_FILENAME]), + outputs: [ATTRIBUTION_FILENAME], + applyResult: applyAttributionResult, }; const attributionDiarized: BackfillKind = { @@ -685,12 +893,18 @@ const attributionDiarized: BackfillKind = { channelSlug: opts.channelSlug, method: "diarized", force: opts.force, + settings: opts.attributionSettings, appConfig: opts.appConfig, + context: opts.digestContext, onLog: opts.onLog, signal: opts.signal, }), ); }, + inputs: (files) => + transcriptUnitInputs(files, [DIARIZATION_FILENAME, ATTRIBUTION_FILENAME]), + outputs: [ATTRIBUTION_FILENAME], + applyResult: applyAttributionResult, }; // One mapping, shared by both kinds, so the two lanes cannot report the same @@ -880,6 +1094,12 @@ const digest: BackfillKind = { channelSlug: opts.channelSlug, videoId: opts.videoId, sections, + // Unit-executor injection. On the primary both are unset and digestVideo + // resolves exactly as before; on an executor the injected config (baseUrl + // stripped — the executor localises its own endpoint) and the injected + // context are what keep the identity the primary's. + ...(opts.appConfig ? { config: opts.appConfig } : {}), + ...(opts.digestContext ? { context: opts.digestContext } : {}), force: opts.force, onLog: opts.onLog, signal: opts.signal, @@ -889,6 +1109,9 @@ const digest: BackfillKind = { if (outcome.status === "fresh") return "already-present"; return "skipped"; }, + inputs: (files) => transcriptUnitInputs(files, [DIGEST_FILENAME]), + outputs: [DIGEST_FILENAME], + applyResult: applyDigestResult, }; type DigestTarget = { diff --git a/common/lib/backfillUnit.test.ts b/common/lib/backfillUnit.test.ts @@ -0,0 +1,228 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { mkdtemp, readFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import path from "node:path"; +import { + BACKFILL_KINDS, + getBackfillKind, +} from "./backfillKinds"; +import { + ATTRIBUTION_FILENAME, + type AttributionRecord, +} from "./attribution"; +import { loadAttribution, writeAttribution } from "./attribution-server"; +import { DIGEST_FILENAME, type DigestProvenance } from "./digest"; +import { loadDigest, writeDigestSection } from "./digest-server"; +import { loadDiarization } from "./diarization-server"; +import { + CUES_JSON_FILENAME, + DIARIZATION_FILENAME, + META_FILENAME, + WHISPER_FILENAME, + type VideoFiles, +} from "./videoStatus"; + +// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test lib/backfillUnit.test.ts +// +// The unit-executor contract: every lane kind declares what a unit needs +// (inputs), what it produces (outputs), and how its output lands back on the +// primary (applyResult — through the GUARDED writers, never a raw copy). The +// two guards that must hold whatever a unit returns: the text-only lane can +// never overwrite a diarized attribution record, and applying one digest +// section never clobbers the other. + +function files(entries: string[]): VideoFiles { + return { + hasMeta: entries.includes(META_FILENAME), + hasYtVtt: entries.includes("transcript.en.vtt"), + ytVttFile: entries.includes("transcript.en.vtt") + ? "transcript.en.vtt" + : null, + hasNonCanonicalVtt: false, + hasWhisper: entries.includes(WHISPER_FILENAME), + hasCuesJson: entries.includes(CUES_JSON_FILENAME), + hasDiarization: entries.includes(DIARIZATION_FILENAME), + isUntranscribable: false, + audioFiles: entries.filter((e) => e.startsWith("audio.")), + partAudioFiles: [], + entries, + }; +} + +const FULL_DIR = files([ + META_FILENAME, + WHISPER_FILENAME, + CUES_JSON_FILENAME, + DIARIZATION_FILENAME, + ATTRIBUTION_FILENAME, + DIGEST_FILENAME, + "audio.mp3", +]); + +test("every lane kind declares the full unit contract", () => { + for (const kind of BACKFILL_KINDS) { + assert.equal(typeof kind.inputs, "function", `${kind.id} inputs`); + assert.ok(kind.outputs.length > 0, `${kind.id} outputs`); + assert.equal(typeof kind.applyResult, "function", `${kind.id} applyResult`); + // Outputs are the SIDECAR the kind writes — never the shipped transcript + // or metadata, which applying would silently rewrite on the primary. + for (const out of kind.outputs) { + assert.ok( + out !== CUES_JSON_FILENAME && + out !== META_FILENAME && + out !== WHISPER_FILENAME, + `${kind.id} output ${out} collides with a shipped input`, + ); + } + } +}); + +test("every transcript-derived kind ships transcript.cues.json LAST", () => { + // The executor materializes inputs in order and isCuesJsonFresh compares + // mtimes: a cues file written before its metadata/raw transcript reads as + // stale and the unit silently does nothing. + for (const id of ["attribution-text", "attribution-diarized", "digest"]) { + const inputs = getBackfillKind(id)!.inputs(FULL_DIR); + assert.equal( + inputs[inputs.length - 1], + CUES_JSON_FILENAME, + `${id} must list cues last, got: ${inputs.join(", ")}`, + ); + assert.ok(inputs.includes(META_FILENAME), `${id} ships metadata`); + assert.ok(inputs.includes(WHISPER_FILENAME), `${id} ships the raw transcript`); + } +}); + +test("the diarized attribution unit ships diarization.json; diarization ships audio", () => { + assert.ok( + getBackfillKind("attribution-diarized")! + .inputs(FULL_DIR) + .includes(DIARIZATION_FILENAME), + ); + assert.ok(getBackfillKind("diarization")!.inputs(FULL_DIR).includes("audio.mp3")); +}); + +function attributionPayload(method: "text-only" | "diarized"): string { + const record: AttributionRecord = { + videoId: "vid", + generatedAt: "2026-08-24T00:00:00.000Z", + speakers: [{ index: 0, label: "Host", seconds: 10 }], + segments: [{ start: 0, end: 10, speaker: 0 }], + provenance: { + method, + appId: "ollama-direct", + model: "m", + modelRequested: "m", + promptVersion: 2, + generatedAt: "2026-08-24T00:00:00.000Z", + durationMs: 1, + }, + }; + return JSON.stringify(record); +} + +test("applyResult refuses a text-only record over a diarized one (the downgrade rule)", async () => { + const dir = await mkdtemp(path.join(tmpdir(), "unit-attr-")); + await writeAttribution( + dir, + JSON.parse(attributionPayload("diarized")) as AttributionRecord, + ); + const outcome = await getBackfillKind("attribution-text")!.applyResult(dir, { + [ATTRIBUTION_FILENAME]: attributionPayload("text-only"), + }); + assert.equal(outcome, "refused"); + assert.equal( + (await loadAttribution(dir))?.provenance.method, + "diarized", + "the diarized record must survive", + ); +}); + +test("applyResult upgrades a text-only record to a diarized one", async () => { + const dir = await mkdtemp(path.join(tmpdir(), "unit-attr-")); + await writeAttribution( + dir, + JSON.parse(attributionPayload("text-only")) as AttributionRecord, + ); + const outcome = await getBackfillKind("attribution-diarized")!.applyResult( + dir, + { [ATTRIBUTION_FILENAME]: attributionPayload("diarized") }, + ); + assert.equal(outcome, "applied"); + assert.equal((await loadAttribution(dir))?.provenance.method, "diarized"); +}); + +test("applyResult rejects a malformed attribution payload", async () => { + const dir = await mkdtemp(path.join(tmpdir(), "unit-attr-")); + const kind = getBackfillKind("attribution-text")!; + assert.equal(await kind.applyResult(dir, {}), "invalid"); + assert.equal( + await kind.applyResult(dir, { [ATTRIBUTION_FILENAME]: "not json" }), + "invalid", + ); + assert.equal( + await kind.applyResult(dir, { [ATTRIBUTION_FILENAME]: "{}" }), + "invalid", + ); +}); + +function digestProvenance(model: string): DigestProvenance { + return { + appId: "ollama-direct", + model, + promptVersion: 2, + contextHash: "none", + generatedAt: "2026-08-24T00:00:00.000Z", + } as DigestProvenance; +} + +test("applying one digest section preserves the other (guarded read-modify-write)", async () => { + const dir = await mkdtemp(path.join(tmpdir(), "unit-digest-")); + // The primary already holds a tags section… + await writeDigestSection(dir, { + section: "tags", + items: [], + provenance: digestProvenance("local-model"), + warnings: [], + }); + // …and a unit returns a record carrying only chapters. + const payload = JSON.stringify({ + sections: { + chapters: { provenance: digestProvenance("m"), items: [] }, + }, + }); + const outcome = await getBackfillKind("digest")!.applyResult(dir, { + [DIGEST_FILENAME]: payload, + }); + assert.equal(outcome, "applied"); + const record = await loadDigest(dir); + assert.ok(record?.sections?.chapters, "the unit's section landed"); + assert.ok( + record?.sections?.tags, + "the section the primary wrote must survive the apply", + ); + assert.equal(record?.sections?.tags?.provenance.model, "local-model"); +}); + +test("diarization results apply verbatim (single-writer sidecar) after a shape check", async () => { + const dir = await mkdtemp(path.join(tmpdir(), "unit-diar-")); + const kind = getBackfillKind("diarization")!; + assert.equal( + await kind.applyResult(dir, { [DIARIZATION_FILENAME]: "{}" }), + "invalid", + ); + const outcome = await kind.applyResult(dir, { + [DIARIZATION_FILENAME]: JSON.stringify({ + videoId: "vid", + generatedAt: "2026-08-24T00:00:00.000Z", + speakers: 1, + turns: [{ start: 0, end: 5, speaker: 0 }], + }), + }); + assert.equal(outcome, "applied"); + const raw = await readFile(path.join(dir, DIARIZATION_FILENAME), "utf8"); + assert.equal((JSON.parse(raw) as { speakers: number }).speakers, 1); + // And the server-side loader accepts what applyResult wrote. + assert.notEqual(await loadDiarization(dir), null); +}); diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -1,6 +1,7 @@ # Changelog ## [Unreleased] +- **A machine with no corpus can now run whole backfill units for this one.** The remote-transcription protocol grew a general sibling: `/api/worker/unit` accepts one unit of any *backfill kind* — speaker attribution, diarization — as a small envelope of input files, runs it against a throwaway scratch corpus, and hands the produced sidecar back. Each kind now declares its own contract: what a unit needs (attribution ships the cue sidecar, metadata and raw transcript — the freshness gate compares their mtimes, so the executor writes the cues file *last* or the unit would silently do nothing), what it produces, and how the result lands back on the primary — always through the guarded writers, never a raw copy, so a unit's text-only record still cannot overwrite a diarized one and applying one digest section still preserves the other. Two refusals are load-bearing: the endpoint takes only backfill kinds (downloads and transcription are refused at the door, so download politeness stays one machine's promise), and a unit that reports "disabled" or "not configured" on the executor is an *error* — on a bare box that answer means the primary's injected model/prompt identity was dropped and the executor's default settings leaked in, which is exactly how a second machine writes permanently-stale records. Deploying an executor needs no new software: this app, `ARCHILYZER_IDLE_BOOT=1`, `WORKER_TOKEN`, and an empty transcripts dir — see RUNNING_IN_DOCKER.md. - **A second machine can now carry the AI backlog by running nothing but `ollama serve`.** The digest and speaker-attribution sweeps — about 194,000 and 78,000 model calls on this archive — bottom out in exactly one HTTP call per chunk; everything around that call (chunking, prompts, parsing, the guarded sidecar writes) is cheap and stays on this box. So the new **LLM endpoint** worker kind is just a URL: no repo, no editor, no token, no copy of the corpus on the other machine. The digest and attribution runners fan their calls across every free endpoint slot alongside the local one, and the lanes' limits rise to match — including while the digest lane is yielding the GPU to transcription, when the *local* term goes to zero and remote endpoints keep the lane moving (an operator pause and the metered spend cap still stop everything; intent and money are global). With no endpoints configured, nothing changes at all. - **An endpoint that can't serve the exact model is refused, loudly, instead of quietly poisoning the corpus.** What makes a record "current" pins the configured model string — and for digests the context size too — so a second machine running *almost* the right model would write records that this box marks stale and re-does forever, with no error anywhere. Only the endpoint URL may vary per call. Before an endpoint's first use the scheduler asks it (ollama's own `/api/tags`) whether the primary's exact tag is present, and one that lacks it — or is unreachable — is degraded with a log line naming the fix, at the cost of one probe rather than a failure per video. A separate tripwire logs when any engine reports having run a different model than was requested, because freshness deliberately compares the requested string and would never notice on its own. An LLM endpoint can never be handed a transcription, and a worker list of nothing but endpoints won't validate — they cannot transcribe, and letting them count as "enabled workers" would park every transcription forever. - **A remote worker is now worth what the remote can actually do, not one slot.** The worker model's own header has promised since it shipped that workers "carry a priority and a slot count" — no such field existed, so a remote box running four workers of its own took one transcription at a time from here unless you hand-copied its card four times. A remote card now has **Slots**: set a number, or leave it blank and the remote is asked directly — its health endpoint has always returned its per-worker summary, and every caller read one boolean off it and threw the rest away. The pool expands one remote into that many independently-schedulable slots (each can be enabled, drained or degraded on its own on the Workers page), re-checks a blank-slots remote about once a minute with no timer and no waiting — a stale answer costs at most one slot for one minute — and shrinking the count retires the surplus through the same drain-then-drop path a removed worker takes. diff --git a/editor/app/api/worker/unit/[id]/events/route.ts b/editor/app/api/worker/unit/[id]/events/route.ts @@ -0,0 +1,26 @@ +import { NextResponse } from "next/server"; +import { authorizeWorkerRequest } from "yt-dlp-transcript-common/lib/workerToken"; +import { getWorkerJobEvents } from "yt-dlp-transcript-common/controller/workerServer"; + +export const dynamic = "force-dynamic"; + +const ID_RE = /^[A-Za-z0-9_-]+$/; + +// Status + incremental log for a unit job. Polled ~1s by the requesting +// instance's remote-unit client. Unit jobs share the worker-job registry with +// transcriptions, so this is the same reader. +export async function GET( + request: Request, + { params }: { params: Promise<{ id: string }> }, +) { + const auth = authorizeWorkerRequest(request.headers.get("authorization")); + if (!auth.ok) { + return NextResponse.json({ error: auth.error }, { status: auth.status }); + } + const { id } = await params; + if (!ID_RE.test(id)) { + return NextResponse.json({ error: "invalid id" }, { status: 400 }); + } + const events = await getWorkerJobEvents(id); + return NextResponse.json(events); +} diff --git a/editor/app/api/worker/unit/[id]/result/route.ts b/editor/app/api/worker/unit/[id]/result/route.ts @@ -0,0 +1,31 @@ +import { NextResponse } from "next/server"; +import { authorizeWorkerRequest } from "yt-dlp-transcript-common/lib/workerToken"; +import { readWorkerUnitResult } from "yt-dlp-transcript-common/controller/workerServer"; + +export const dynamic = "force-dynamic"; + +const ID_RE = /^[A-Za-z0-9_-]+$/; + +// The unit's outcome plus its kind's declared output files (UTF-8 JSON text by +// name), once the run is done. 404 until then. The requester applies these on +// its side through the kind's guarded writers — never as a raw copy. +export async function GET( + request: Request, + { params }: { params: Promise<{ id: string }> }, +) { + const auth = authorizeWorkerRequest(request.headers.get("authorization")); + if (!auth.ok) { + return NextResponse.json({ error: auth.error }, { status: auth.status }); + } + const { id } = await params; + if (!ID_RE.test(id)) { + return NextResponse.json({ error: "invalid id" }, { status: 400 }); + } + const result = await readWorkerUnitResult(id); + if (!result) { + return NextResponse.json({ error: "no result yet" }, { status: 404 }); + } + return NextResponse.json(result, { + headers: { "Cache-Control": "no-store" }, + }); +} diff --git a/editor/app/api/worker/unit/[id]/route.ts b/editor/app/api/worker/unit/[id]/route.ts @@ -0,0 +1,30 @@ +import { NextResponse } from "next/server"; +import { authorizeWorkerRequest } from "yt-dlp-transcript-common/lib/workerToken"; +import { + cancelWorkerJob, + cleanupWorkerJob, +} from "yt-dlp-transcript-common/controller/workerServer"; + +export const dynamic = "force-dynamic"; + +const ID_RE = /^[A-Za-z0-9_-]+$/; + +// Cancel (if running) and clean up the unit's scratch corpus. Called when the +// requester has applied the result or given up. Idempotent — unit jobs share +// the worker-job registry with transcriptions. +export async function DELETE( + request: Request, + { params }: { params: Promise<{ id: string }> }, +) { + const auth = authorizeWorkerRequest(request.headers.get("authorization")); + if (!auth.ok) { + return NextResponse.json({ error: auth.error }, { status: auth.status }); + } + const { id } = await params; + if (!ID_RE.test(id)) { + return NextResponse.json({ error: "invalid id" }, { status: 400 }); + } + cancelWorkerJob(id); + await cleanupWorkerJob(id); + return NextResponse.json({ ok: true }); +} diff --git a/editor/app/api/worker/unit/route.ts b/editor/app/api/worker/unit/route.ts @@ -0,0 +1,51 @@ +import { NextResponse } from "next/server"; +import { authorizeWorkerRequest } from "yt-dlp-transcript-common/lib/workerToken"; +import { getBackfillKind } from "yt-dlp-transcript-common/lib/backfillKinds"; +import { + startWorkerUnit, + type StartWorkerUnitInput, +} from "yt-dlp-transcript-common/controller/workerServer"; + +export const dynamic = "force-dynamic"; + +// Accept ONE backfill unit (attribution, diarization, …) from another +// instance: a JSON envelope carrying the kind id, the input files (base64 by +// name), the primary's resolved freshness target and its identity config. +// Returns { remoteJobId } (202); the requester polls +// /api/worker/unit/<id>/events and pulls /result. Same token guard as the +// transcription protocol. +// +// Anything that is not a backfill KIND is refused at the door — download and +// transcription are ExternalOperations, and letting a second host run a +// download would double the per-platform politeness this repo is careful +// about. +export async function POST(request: Request) { + const auth = authorizeWorkerRequest(request.headers.get("authorization")); + if (!auth.ok) { + return NextResponse.json({ error: auth.error }, { status: auth.status }); + } + let body: StartWorkerUnitInput; + try { + body = (await request.json()) as StartWorkerUnitInput; + } catch { + return NextResponse.json({ error: "malformed JSON body" }, { status: 400 }); + } + if (!body.op || !body.channelSlug || !body.videoId) { + return NextResponse.json( + { error: "op, channelSlug and videoId are required" }, + { status: 400 }, + ); + } + if (!getBackfillKind(body.op)) { + return NextResponse.json( + { error: `"${body.op}" is not a backfill kind this executor can run` }, + { status: 400 }, + ); + } + try { + const { remoteJobId } = await startWorkerUnit(body); + return NextResponse.json({ remoteJobId }, { status: 202 }); + } catch (e) { + return NextResponse.json({ error: (e as Error).message }, { status: 500 }); + } +}