// THE BATCH'S FIFTH GUARD: a relocation that starts AFTER the run did. // // Run with: // node_modules/.bin/tsx --test common/controller/operationBatchRelocation.test.ts // // Its own file because it needs the SETTINGS SEAM — the diarization entry is // only in the backfill lane when settings say so, and getPaths() memoizes its // first answer at module scope, so TRANSCRIPTS_DIR and SETTINGS_FILE must be // set before anything can import the module under test. Same arrangement as // videoOperations.test.ts and laneForOperation.test.ts. // // WHAT IT PINS. The reachability guard at the top of runOperationBatch runs // ONCE, at start. The omnimirror incident (2026-09-13) is what that misses: the // relocate copy began while an operation batch was already running, and a unit // writing into `data/` mid-rsync left the copy's verify refusing. The marker is // now re-read on every candidate pull, and finding one ENDS the run. import { mkdtempSync, writeFileSync } from "node:fs"; import { mkdir, rm, 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(), "batch-relocation-")); process.env.TRANSCRIPTS_DIR = ROOT; const SETTINGS_FILE = path.join(ROOT, "settings.json"); process.env.SETTINGS_FILE = SETTINGS_FILE; // The backfill lane armed with ONE operation, and that operation's cheapest // classification: every seeded video has a transcript and no audio, so // diarization reads `missing-input` — counted, never dispatched, and (with // re-download off) needing no engine, no model and no network. The guard under // test lives BEFORE the classification, so what the classification answers only // has to be deterministic, not interesting. writeFileSync( SETTINGS_FILE, JSON.stringify({ // THE LANE'S GATE, SPELLED. It used to be spelled by `backfill.enabled: // true` — the inverted retired field, where `true` meant NOT held — // and S0-pause deleted it. The backfill lane's gate DEFAULTS shut // (`defaultHeldFor`, which is the reading that field always gave a file // naming no gate), so a fixture that needs the batch to dispatch says so // on the lane. Without it `runOperationBatch` holds and never resolves. autoQueue: { backfill: { enabled: true, held: false } }, backfill: { allowRedownload: false, concurrency: 1 }, diarization: { enabled: true, segModel: "/models/seg-1.onnx", embModel: "/models/emb-1.onnx", }, attribution: { enabled: false, diarizedEnabled: false, textOnlyEnabled: false }, }), "utf8", ); const { runOperationBatch } = await import("./operationBatch"); const { getPaths } = await import("../lib/paths"); const { relocationMarkerPath } = await import("../lib/channelMedia"); after(() => rm(ROOT, { recursive: true, force: true })); async function seed(slug: string, ids: string[]): Promise { 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 }); // Transcribed (so diarization is applicable) and no audio (so its input is // missing). A non-empty transcription array keeps it out of // `isUntranscribable`. await writeFile( path.join(videoDir, "transcript.json"), JSON.stringify({ transcription: [{ text: "hello", offsets: {} }] }), ); } } test("every candidate is classified when no relocation is in flight", async () => { await seed("quiet", ["v1", "v2"]); const lines: string[] = []; const result = await runOperationBatch({ lane: "backfill", channelSlug: "quiet", paths: getPaths(), operationIds: ["diarization"], onLog: (m) => lines.push(m), }); // The control for the case below: both videos are reached and counted. assert.equal(result.missingInput, 2); assert.equal(result.attempted, 0); assert.equal(result.stoppedForRelocation, false); assert.ok(!lines.some((l) => l.includes("Stopping: a relocation"))); }); test("a marker written mid-run ends the batch instead of racing the copy", async () => { await seed("omni", ["v1", "v2"]); const marker = relocationMarkerPath(getPaths(), "omni"); // setProgress fires once to seed the bar and then once per resolved // candidate, so the SECOND call is the instant after the first video was // classified and before the second is pulled — exactly the window a relocate // job starts in. Written synchronously, because that is how the disk would // have changed by the time the next pull reads it. let progressCalls = 0; const lines: string[] = []; const result = await runOperationBatch({ lane: "backfill", channelSlug: "omni", paths: getPaths(), operationIds: ["diarization"], onLog: (m) => lines.push(m), setProgress: () => { progressCalls++; if (progressCalls !== 2) return; writeFileSync( marker, JSON.stringify({ target: "/platter/omni/data", direction: "out", startedAt: new Date().toISOString(), phase: "copy", }), "utf8", ); }, }); // ONE video was classified, not two: the run stopped at the next pull rather // than working its way through the rest of the channel while rsync copied it. assert.equal(result.missingInput, 1); assert.equal(result.attempted, 0); // The flag is how the job's summary line says why it stopped — an early end // that the counters alone would render as a quiet, complete run. assert.equal(result.stoppedForRelocation, true); const stop = lines.find((l) => l.includes("Stopping: a relocation")); assert.ok(stop, `expected a stop line, got:\n${lines.join("\n")}`); assert.match(stop as string, /\(out\) of omni to \/platter\/omni\/data/); });