Archilyzer · Source

archilyzer

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

commit 3ffd166de250210552ea57c454b7064a65f90770
parent c4a0862c7f3138810f63bd05bfdf86d45fcc4a1d
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Thu,  1 Oct 2026 10:49:38 -0400

common: a media move refuses to start over a writer, holds the channel's writers while its marker stands, and mirrors its copy before verifying

- controller/channelWriters.ts (new): who is writing into a channel — running
  (and, for the editor's courtesy, queued) registry jobs on its slug plus the
  lanes' in-process units — named: "a transcription of v50t5yt is running
  (Transcribe all, job …)". The move's preview, its job's first step and a
  third ask right after the copy phase's marker is written refuse over one;
  a fresh move's marker is removed when the third ask refuses.
- relocateDir.ts: the mirror pass (`rsync -a --delete`) after the copy pass,
  only ever toward the copy under construction (assertMirrorDirection refuses
  any other direction before rsync is spawned); the verify's dry run carries
  --delete, one change between the mirror and the check gets one more pass,
  and a second refuses with the paths by kind (extra on the destination /
  missing on the destination / changed). copyMirrorVerify is the copy phase
  for both movers and both directions; a resume takes it, so stale files left
  on a destination by an interrupted move are removed. Reconcile skips the
  copy pass, lists the differences by kind, and mirrors.
- relocateChannelMedia: `reconcile` and `writers`; the swap's re-verify
  mirrors only while data/ is still the live directory, never from a parked
  copy. relocateSavedVideos: the same pipeline, and a `busy` first step.
- jobs/streamCommand: the needsMedia guard is asked again when the queue
  starts a job, and a move's marker refuses in the hold's words.
- jobs/jobKinds: needsMedia for the per-video writers that were not listed —
  whisper-video, transcribe-one, download-one-pipeline, transcode-audio and
  the three availability checks.
- lib/channelMediaHold: "its media is moving" leads the in-transition reason;
  mediaHoldText is the one wording. The lanes skip on isMediaHeld; the row
  view carries `mediaHold`.

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

Diffstat:
Mcommon/controller/autoRunner.test.ts | 76+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcommon/controller/autoRunner.ts | 23+++++++++++++++++------
Acommon/controller/channelWriters.test.ts | 172+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/controller/channelWriters.ts | 189+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/controller/relocateChannelMedia.test.ts | 349++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----
Mcommon/controller/relocateChannelMedia.ts | 216++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------------------
Mcommon/controller/relocateDir.test.ts | 315++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcommon/controller/relocateDir.ts | 375+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Mcommon/controller/relocateSavedVideos.test.ts | 52++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/controller/relocateSavedVideos.ts | 93++++++++++++++++++++++++++++++++++++++-----------------------------------------
Mcommon/jobs/jobKinds.ts | 57+++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/streamCommand.test.ts | 51++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcommon/jobs/streamCommand.ts | 38+++++++++++++++++++++++++++++++++++---
Mcommon/lib/channelMediaHold.ts | 17++++++++++++++++-
Mcommon/views/channelRow.ts | 8++++++++
15 files changed, 1889 insertions(+), 142 deletions(-)

diff --git a/common/controller/autoRunner.test.ts b/common/controller/autoRunner.test.ts @@ -1,6 +1,6 @@ import { test } from "node:test"; import assert from "node:assert/strict"; -import { mkdirSync, mkdtempSync, writeFileSync } from "node:fs"; +import { mkdirSync, mkdtempSync, rmSync, writeFileSync } from "node:fs"; import { tmpdir } from "node:os"; import path from "node:path"; import { @@ -33,6 +33,8 @@ import type { } from "../jobs/autoQueuePolicy"; import type { Paths } from "../lib/paths"; import type { SiteSettings } from "../lib/settings"; +import { forgetChannelMedia } from "../lib/channelMedia"; +import { bucketLaneOperationId } from "../lib/operations"; // THE RUNNER'S HALF OF S1 (plans/channel-priority.md), and it is here rather // than in `jobs/channelPriorityCompile.test.ts` for one reason: @@ -649,3 +651,75 @@ test("sharedAutoQueueState: a reset singleton reads afresh; a different state fi assert.equal(await sharedAutoQueueState(p2), y); freshSingletonForTest(); }); + +// --- The media hold: every lane skips a channel whose move marker stands ---- +// +// Release 16 slice RM. The lanes' one channel-level chokepoint is the +// projection computeLeafPending shares with the runner (buildChannelWork): a +// channel whose `.relocating.json` stands is held (lib/channelMediaHold.ts) and +// projects no work on any lane, and the hold lifts the moment the marker goes +// — the movers forget the inspect memo on every marker write and removal, as +// this test does by hand. + +async function laneTotal( + kind: (typeof LANES)[number], + paths: Paths, + configs: { slug: string; config: never }[], +): Promise<number> { + const r = await computeLeafPending(kind, paths, { + configs, + state: emptyAutoQueueState(), + }); + return Object.values(r.counts).reduce((a, b) => a + b, 0); +} + +test("a channel with a move marker projects no work on any lane, and the hold lifts with it", async () => { + const paths = pendingFixture(); + const slug = "moving"; + const channelDir = path.join(paths.channelsDir, slug); + mkdirSync(path.join(channelDir, "data", "v1"), { recursive: true }); + const config = { + url: "https://www.youtube.com/@moving", + handling: "transcribe", + }; + writeFileSync(path.join(channelDir, "config.json"), JSON.stringify(config)); + const opId = bucketLaneOperationId("transcription"); + writeFileSync( + path.join(channelDir, "snapshot.json"), + JSON.stringify({ + generatedAt: new Date().toISOString(), + totals: { videos: 1, downloaded: 1, transcribed: 0 }, + buckets: { downloadedNoTranscript: ["v1"], noTranscript: ["v1"] }, + ...(opId ? { backfill: { [opId]: { ids: ["v1"] } } } : {}), + }), + ); + const configs = [{ slug, config: config as never }]; + forgetChannelMedia(slug); + assert.ok( + (await laneTotal("transcription", paths, configs)) > 0, + "in place, the channel's video is the lane's work", + ); + + writeFileSync( + path.join(channelDir, ".relocating.json"), + JSON.stringify({ + target: "/mnt/platter/moving/data", + direction: "out", + startedAt: new Date().toISOString(), + phase: "copy", + }), + ); + forgetChannelMedia(slug); + for (const lane of LANES) { + assert.equal( + await laneTotal(lane, paths, configs), + 0, + `the ${lane} lane skips a channel whose media is moving`, + ); + } + + // The move completes (or is abandoned): the marker goes, the hold lifts. + rmSync(path.join(channelDir, ".relocating.json")); + forgetChannelMedia(slug); + assert.ok((await laneTotal("transcription", paths, configs)) > 0); +}); diff --git a/common/controller/autoRunner.ts b/common/controller/autoRunner.ts @@ -89,6 +89,7 @@ import { readRelocationMarker, type ChannelMediaStatus, } from "../lib/channelMedia"; +import { isMediaHeld, mediaHoldText } from "../lib/channelMediaHold"; import { type ChannelPriority, type FocusSummary, @@ -563,15 +564,19 @@ function noteSkippedForMedia( ): void { if (mediaSkipLogged.get(slug) === status) return; mediaSkipLogged.set(slug, status); + // The hold's own words first (lib/channelMediaHold.ts — "held: its media is + // moving …" for a channel with a relocation marker), then the status and the + // inspector's detail as they always were. console.log( - `[auto] skipping ${slug}: media ${status}${detail ? ` — ${detail}` : ""}`, + `[auto] skipping ${slug}, ${mediaHoldText(status) ?? "held"}: ` + + `media ${status}${detail ? ` — ${detail}` : ""}`, ); } function noteMediaReachable(slug: string): void { if (!mediaSkipLogged.has(slug)) return; mediaSkipLogged.delete(slug); - console.log(`[auto] ${slug}: media reachable again`); + console.log(`[auto] ${slug}: media reachable again — the hold is lifted`); } // For tests and for /api/test/invalidate-cache: the log-once memory is @@ -625,10 +630,15 @@ async function buildChannelWork( const snap = snaps[i]; if (!snap) continue; // A SKIP IS A SKIP. The lane keeps running every other channel: this is - // never a lane stop and it is not a hold — "a zero limit is a hold, never a - // stop" (pauseGates.ts) is a different mechanism and is untouched by it. + // never a lane stop, and it is not the LANE's hold — "a zero limit is a + // hold, never a stop" (pauseGates.ts) is a different mechanism and is + // untouched by it. It is the CHANNEL's media hold (lib/channelMediaHold.ts, + // the words both builds use): a relocation marker, an unmounted or stalled + // drive, a link and a config that disagree. It lifts by itself — the next + // tick after the marker is removed (a move completed or abandoned; the + // movers forget the memo) or the drive is back reads the channel again. const location = media[i]; - if (location && location.status !== "ok" && location.status !== "in-place") { + if (location && isMediaHeld(location.status)) { noteSkippedForMedia(slug, location.status, location.detail); continue; } @@ -1887,7 +1897,8 @@ async function runLoop( const marker = await readRelocationMarker(paths, channelSlug); if (marker) { onLog( - `Auto-${kind}: skipping ${channelSlug}/${pick.videoId} — a relocation ` + + `Auto-${kind}: skipping ${channelSlug}/${pick.videoId}, ` + + `${mediaHoldText("in-transition")} — a relocation ` + `(${marker.direction}) to ${marker.target} is in flight.`, ); // `return` still runs the `finally`, which is what releases the diff --git a/common/controller/channelWriters.test.ts b/common/controller/channelWriters.test.ts @@ -0,0 +1,172 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import type { JobRecord } from "../jobs/registry"; +import type { AutoQueueKind } from "../lib/autoQueueTypes"; +import type { AutoRunnerInFlight } from "./autoRunner"; +import { + channelWriters, + channelWritersRefusal, + describeChannelWriter, + type ChannelWritersSource, +} from "./channelWriters"; + +// Run with: +// pnpm --filter yt-dlp-transcript-common exec tsx --test controller/channelWriters.test.ts +// +// WHO IS WRITING INTO A CHANNEL, named (release 16 slice RM). The registry and +// the lanes are injected: what is pinned is the reading of them — which +// records count, in which order, and the sentence the move refuses with. + +function job(over: Partial<JobRecord>): JobRecord { + return { + id: "J1", + kind: "whisper-all", + queueKey: "transcription", + status: "running", + queuedAt: 1, + logPath: "/dev/null", + ...over, + }; +} + +function unit( + lane: AutoQueueKind, + over: Partial<AutoRunnerInFlight>, +): { lane: AutoQueueKind; unit: AutoRunnerInFlight } { + return { + lane, + unit: { + videoId: "v1", + leafId: "leaf", + channelSlug: "alpha", + startedAt: 1, + ...over, + }, + }; +} + +function source( + jobs: JobRecord[], + units: { lane: AutoQueueKind; unit: AutoRunnerInFlight }[] = [], +): ChannelWritersSource { + return { jobs: () => jobs, units: () => units }; +} + +test("the 2026-09-30 case: a running Transcribe all is named by the video it is on", () => { + const writers = channelWriters("realcandaceo", { + source: source([ + job({ + id: "01M3TGYDXA7WCBGN014A90J3JW", + channelSlug: "realcandaceo", + tasks: [ + { id: "v50t5yt", label: "v50t5yt", kind: "transcribe", startedAt: 1 }, + ], + }), + ]), + }); + assert.equal(writers.length, 1); + assert.equal( + describeChannelWriter(writers[0]), + "a transcription of v50t5yt is running (Transcribe all, job 01M3TGYDXA7WCBGN014A90J3JW)", + ); + assert.equal( + channelWritersRefusal("realcandaceo", writers), + 'Cannot move the media of "realcandaceo" now: a transcription of v50t5yt ' + + "is running (Transcribe all, job 01M3TGYDXA7WCBGN014A90J3JW) — wait for " + + "it or cancel it. Nothing has been touched.", + ); +}); + +test("only this channel's jobs, only running ones unless queued are asked for", () => { + const jobs = [ + job({ id: "A", channelSlug: "alpha", kind: "sync", status: "queued" }), + job({ id: "B", channelSlug: "beta", kind: "sync" }), + job({ id: "C", channelSlug: "alpha", kind: "sync", status: "done" }), + job({ id: "D", channelSlug: "alpha", kind: "download-missing" }), + ]; + assert.deepEqual( + channelWriters("alpha", { source: source(jobs) }).map((w) => + w.source === "job" ? w.jobId : "", + ), + ["D"], + ); + // Running before queued, whatever the registry order. + assert.deepEqual( + channelWriters("alpha", { source: source(jobs), includeQueued: true }).map( + (w) => (w.source === "job" ? w.jobId : ""), + ), + ["D", "A"], + ); + assert.equal( + describeChannelWriter( + channelWriters("alpha", { source: source(jobs), includeQueued: true })[1], + ), + "Sync is queued (job A)", + ); +}); + +test("the move's own kind is not a writer to refuse over", () => { + const jobs = [ + job({ id: "M", channelSlug: "alpha", kind: "relocate-channel-media" }), + ]; + assert.equal(channelWriters("alpha", { source: source(jobs) }).length, 1); + assert.deepEqual( + channelWriters("alpha", { + source: source(jobs), + ignoreKinds: ["relocate-channel-media"], + }), + [], + ); +}); + +test("a lane's in-process unit is a writer, and its way out is the lane", () => { + const writers = channelWriters("alpha", { + source: source([], [unit("transcription", { videoId: "v9" })]), + }); + assert.equal( + channelWritersRefusal("alpha", writers), + 'Cannot move the media of "alpha" now: a transcription of v9 is running ' + + "(the transcription lane) — wait for it, or hold the transcription lane. " + + "Nothing has been touched.", + ); +}); + +test("a download unit is named once, as the lane's, not again as its job", () => { + const writers = channelWriters("alpha", { + source: source( + [job({ id: "U1", kind: "auto-download-unit", channelSlug: "alpha" })], + [unit("download", { videoId: "v2", jobId: "U1" })], + ), + }); + assert.equal(writers.length, 1); + assert.equal( + describeChannelWriter(writers[0]), + "a download of v2 is running (the download lane)", + ); +}); + +test("a single-video job with no task is named by its record's video", () => { + const writers = channelWriters("alpha", { + source: source([ + job({ id: "W", kind: "whisper-video", channelSlug: "alpha", videoId: "v3" }), + ]), + }); + assert.equal( + describeChannelWriter(writers[0]), + "whisper-video of v3 is running (job W)", + ); +}); + +test("more than one writer: the first named, the rest counted; none is no refusal", () => { + const writers = channelWriters("alpha", { + source: source( + [job({ id: "D", kind: "download-missing", channelSlug: "alpha" })], + [unit("digest", { videoId: "v4" })], + ), + }); + assert.match( + channelWritersRefusal("alpha", writers) ?? "", + /Download missing is running \(job D\), and 1 more writer\(s\) — wait for it or cancel it\./, + ); + assert.equal(channelWritersRefusal("alpha", []), null); +}); diff --git a/common/controller/channelWriters.ts b/common/controller/channelWriters.ts @@ -0,0 +1,189 @@ +import { + getRegistry, + type JobRecord, + type JobStatus, + type JobTaskKind, +} from "../jobs/registry"; +import { jobKindLabel } from "../jobs/jobKinds"; +import { LANES, type AutoQueueKind } from "../lib/autoQueueTypes"; +import { getAutoRunnerStatus, type AutoRunnerInFlight } from "./autoRunner"; + +// WHO IS WRITING INTO A CHANNEL'S MEDIA RIGHT NOW — one answer, named. +// +// A media move must not start over a writer (release 16 slice RM). On +// 2026-09-30 a move of `realcandaceo` was queued at 21:16:42 and asked the +// question then, when nothing was running; a "Transcribe all" started at +// 21:23:30, and the move — which had waited twenty minutes behind other moves +// on the one relocation queue — started its copy at 21:36:28 without asking +// again. The transcription of `v50t5yt` finished at 21:38:07, after rsync had +// passed that directory, and the verify refused a 13.57 GB copy. So the move's +// PREVIEW and its JOB'S FIRST STEP ask here too, not only the action that +// enqueued it, and the refusal names the writer so the operator knows what to +// wait for or cancel. +// +// TWO HALVES, for the reason editor/app/channels/lib/mediaBusy.ts gives: the +// job registry is half the truth. The auto-queue lanes run their per-video +// units in-process and make no job record (the omnimirror incident, +// 2026-09-13), so `getAutoRunnerStatus(lane).inFlight` is the other half. Both +// are memory reads — no disk, no await. +// +// SERVER-SIDE ONLY: it reaches the runner, which imports execa transitively. + +export type ChannelWriter = + | { + source: "job"; + jobId: string; + kind: string; + // The kind's label (jobKindLabel), or the raw kind when it has none. + label: string; + status: Extract<JobStatus, "running" | "queued">; + // What the job is on now: its first in-flight task, else the video the + // record names (a single-video job). + videoId?: string; + taskKind?: JobTaskKind; + } + | { + source: "lane"; + lane: AutoQueueKind; + videoId: string; + // A unit that is not a video (the download lane's metadata scan) says + // what it is instead. + note?: string; + // The registry job running the unit, when it is one. + jobId?: string; + }; + +// Where the two halves are read from. Injected by the tests; production reads +// the live registry and the live runners. +export type ChannelWritersSource = { + jobs: () => ReadonlyArray<JobRecord>; + units: () => ReadonlyArray<{ lane: AutoQueueKind; unit: AutoRunnerInFlight }>; +}; + +export const liveChannelWritersSource: ChannelWritersSource = { + jobs: () => getRegistry().list(), + units: () => + LANES.flatMap((lane) => + getAutoRunnerStatus(lane).inFlight.map((unit) => ({ lane, unit })), + ), +}; + +export type ChannelWritersOptions = { + // Count queued jobs too. The editor's courtesy check does (a queued job is + // about to write); the move's own check does not — a queued media job that + // starts after the marker is written refuses itself at its start + // (streamCommand.ts), so refusing the move over it would trade the + // operator's twenty-minute wait in the relocation queue for nothing. + includeQueued?: boolean; + // Kinds that are not writers for this question. The move's own job passes + // `relocate-channel-media`: it is running on this channel's slug, it is the + // asker, and the relocation queue runs one move at a time. + ignoreKinds?: ReadonlyArray<string>; + source?: ChannelWritersSource; +}; + +// Every writer on `slug`, jobs first (running before queued), then lane units. +export function channelWriters( + slug: string, + opts: ChannelWritersOptions = {}, +): ChannelWriter[] { + const source = opts.source ?? liveChannelWritersSource; + const ignore = new Set(opts.ignoreKinds ?? []); + const running: ChannelWriter[] = []; + const queued: ChannelWriter[] = []; + for (const j of source.jobs()) { + if (j.channelSlug !== slug || ignore.has(j.kind)) continue; + if (j.status !== "running" && !(opts.includeQueued && j.status === "queued")) { + continue; + } + const task = (j.tasks ?? []).find((t) => t.kind !== "relocate"); + (j.status === "running" ? running : queued).push({ + source: "job", + jobId: j.id, + kind: j.kind, + label: jobKindLabel(j.kind), + status: j.status, + ...(task ? { videoId: task.id, taskKind: task.kind } : {}), + ...(!task && j.videoId ? { videoId: j.videoId } : {}), + }); + } + const jobs = [...running, ...queued]; + const units: ChannelWriter[] = source + .units() + .filter(({ unit }) => unit.channelSlug === slug) + .map(({ lane, unit }) => ({ + source: "lane" as const, + lane, + videoId: unit.videoId, + ...(unit.note ? { note: unit.note } : {}), + ...(unit.jobId ? { jobId: unit.jobId } : {}), + })); + // A lane's download unit is ALSO a registry job (auto-download-unit), so it + // would be named twice; the lane's wording is the one that says where it + // came from. + const unitJobs = new Set( + units.flatMap((u) => (u.source === "lane" && u.jobId ? [u.jobId] : [])), + ); + return [ + ...jobs.filter((w) => !(w.source === "job" && unitJobs.has(w.jobId))), + ...units, + ]; +} + +const TASK_NOUN: Record<JobTaskKind, string> = { + download: "a download", + transcribe: "a transcription", + digest: "a digest", + backfill: "a backfill", + relocate: "a media move", +}; + +const LANE_NOUN: Record<AutoQueueKind, string> = { + transcription: "a transcription", + download: "a download", + digest: "a digest", + backfill: "a backfill", +}; + +// One writer, as a clause: "a transcription of v50t5yt is running (Transcribe +// all, job 01M3…)", "Sync is queued (job 01M3…)", "a download of abc is +// running (the download lane)". +export function describeChannelWriter(w: ChannelWriter): string { + if (w.source === "lane") { + const what = w.note ?? `${LANE_NOUN[w.lane]} of ${w.videoId}`; + return `${what} is running (the ${w.lane} lane)`; + } + if (w.taskKind && w.videoId) { + return ( + `${TASK_NOUN[w.taskKind]} of ${w.videoId} is ${w.status} ` + + `(${w.label}, job ${w.jobId})` + ); + } + const what = w.videoId ? `${w.label} of ${w.videoId}` : w.label; + return `${what} is ${w.status} (job ${w.jobId})`; +} + +// What the operator does about the first writer: a job is waited for or +// cancelled on /jobs; a lane's unit is waited for, or its lane held. +function waysOut(w: ChannelWriter): string { + return w.source === "lane" + ? `wait for it, or hold the ${w.lane} lane` + : "wait for it or cancel it"; +} + +// The first writer named, the rest counted, and what to do — or null when +// nothing is writing. `doing` is the clause the refusal opens with ("move the +// media of"), so the same sentence reads right for each mover. +export function channelWritersRefusal( + slug: string, + writers: ReadonlyArray<ChannelWriter>, + doing = "move the media of", +): string | null { + if (writers.length === 0) return null; + const [first, ...rest] = writers; + const more = rest.length > 0 ? `, and ${rest.length} more writer(s)` : ""; + return ( + `Cannot ${doing} "${slug}" now: ${describeChannelWriter(first)}${more} — ` + + `${waysOut(first)}. Nothing has been touched.` + ); +} diff --git a/common/controller/relocateChannelMedia.test.ts b/common/controller/relocateChannelMedia.test.ts @@ -15,6 +15,7 @@ import { utimes, writeFile, } from "node:fs/promises"; +import { writeFileSync } from "node:fs"; import { tmpdir } from "node:os"; import path from "node:path"; import type { Paths } from "../lib/paths"; @@ -24,6 +25,7 @@ import { readRelocationMarker, relocatedDataDir, } from "../lib/channelMedia"; +import type { ChannelWriter } from "./channelWriters"; import { assertRelocationRootPresent, previewRelocation, @@ -252,19 +254,50 @@ test("abort from an onLog hook leaves the source intact, and the rerun completes }); }); -test("a verify failure keeps the source and does not write the config", async () => { +// A STRAY FILE ON THE DESTINATION IS MIRRORED AWAY, NOT REFUSED (release 16 +// slice RM). Until this slice the copy never deleted, so a stray at the target +// passed the itemized dry run and failed the counts — forever: no rerun could +// settle it. The mirror pass removes it from the COPY; the source is never the +// target of --delete. +test("a stray file on the destination is removed by the mirror pass", async () => { await withTmp(async (paths, root) => { const channelDir = await seed(paths, "alpha", { v1: { "audio.m4a": "one" }, }); const target = relocatedDataDir(root, "alpha"); - // A stray file at the target from some earlier, unrelated write. rsync - // without --delete leaves it, so the itemize dry-run is clean and only the - // two-sided measurement catches the drift — which is exactly why the verify - // is both checks and not just the dry run. await mkdir(path.join(target, "v9"), { recursive: true }); await writeFile(path.join(target, "v9", "stray.m4a"), "not ours"); + const res = await relocateChannelMedia({ + io: TEST_IO, + paths, + slug: "alpha", + direction: "out", + root, + onLog: () => {}, + }); + assert.equal(res.files, 1); + await assert.rejects(() => stat(path.join(target, "v9"))); + assert.equal( + await readFile(path.join(target, "v1", "audio.m4a"), "utf8"), + "one", + ); + assert.equal((await readChannelConfig(paths, "alpha"))?.dataDir, target); + assert.deepEqual(await leftoverCopies(channelDir), []); + }); +}); + +// WHAT STILL REFUSES: a source that keeps changing. A writer that lands a file +// before every check is seen by the first verify (one more mirror pass) and +// again by the second, and the move refuses naming the paths by kind — with the +// source untouched, the config unwritten, and the marker left for a Reconcile +// and resume once the writer has stopped. +test("a verify failure keeps the source and does not write the config", async () => { + await withTmp(async (paths, root) => { + const channelDir = await seed(paths, "alpha", { + v1: { "audio.m4a": "one" }, + }); + let n = 0; await assert.rejects( () => relocateChannelMedia({ @@ -273,9 +306,17 @@ test("a verify failure keeps the source and does not write the config", async () slug: "alpha", direction: "out", root, - onLog: () => {}, + onLog: (m) => { + if (m.startsWith("$ ") && m.includes("--dry-run")) { + n++; + writeFileSync( + path.join(channelDir, "data", "v1", `chunk-${n}.json`), + "{}", + ); + } + }, }), - /Verification failed/, + /Verification failed: after a second mirror pass .*1 missing on the destination \(v1\/chunk-2\.json\)/, ); assert.ok((await lstat(path.join(channelDir, "data"))).isDirectory()); @@ -284,6 +325,7 @@ test("a verify failure keeps the source and does not write the config", async () "one", ); assert.equal((await readChannelConfig(paths, "alpha"))?.dataDir, undefined); + assert.equal((await readRelocationMarker(paths, "alpha"))?.phase, "copy"); }); }); @@ -1019,7 +1061,13 @@ test("out @ swap: a directory timestamp is settled by one more pass, not refused }); }); -test("out @ swap: a file the target is missing is still a refusal", async () => { +// CONTENT drift at the swap's re-verify, arriving the way the timestamp did — +// a sidecar written into the source after the copy. Until release 16 slice RM +// any file line refused; now the re-verify runs in mirror mode while `data/` is +// still the live directory, and one change gets one more mirror pass, exactly +// as in the copy phase. A second change is still a refusal ("a verify failure +// keeps the source …" above). +test("out @ swap: a file the target is missing is mirrored by one more pass", async () => { await withTmp(async (paths, root) => { const channelDir = await seed(paths, "alpha", { v1: { "audio.m4a": "one".repeat(500) }, @@ -1027,13 +1075,141 @@ test("out @ swap: a file the target is missing is still a refusal", async () => const dataDir = path.join(channelDir, "data"); const target = relocatedDataDir(root, "alpha"); await copyTree(dataDir, target); - // CONTENT drift, arriving the same way the timestamp did — a sidecar - // written into the source after the copy. It bumps v1/'s mtime too, so the - // drift list is a `.d..t` line AND a file line: one non-directory line is - // all it takes, and the retry must not swallow it. await writeFile(path.join(dataDir, "v1", "diarization.json"), "{}"); await seedMarker(paths, "alpha", { target, direction: "out", phase: "swap" }); + const res = await relocateChannelMedia({ + io: TEST_IO, + paths, + slug: "alpha", + direction: "out", + root, + onLog: () => {}, + }); + assert.equal(res.retried, true); + assert.equal( + await readFile(path.join(target, "v1", "diarization.json"), "utf8"), + "{}", + ); + assert.equal((await inspectChannelMedia(paths, "alpha")).status, "ok"); + assert.equal((await readChannelConfig(paths, "alpha"))?.dataDir, target); + }); +}); + +// A PARKED COPY IS NEVER MIRRORED FROM. Once the rename has committed, the +// link is what readers follow and `data.relocated-*` is not live media any +// more: the re-verify against it is the strict one, as it always was — a +// difference refuses, and nothing on the target is deleted to match a stale +// copy. +test("out @ swap: after the rename, the re-verify against the parked copy is strict", async () => { + await withTmp(async (paths, root) => { + const channelDir = await seed(paths, "alpha", { + v1: { "audio.m4a": "one".repeat(500) }, + }); + const dataDir = path.join(channelDir, "data"); + const target = relocatedDataDir(root, "alpha"); + await copyTree(dataDir, target); + // A file on the target the parked copy lacks. + await writeFile(path.join(target, "v1", "newer.json"), "{}"); + await rename(dataDir, path.join(channelDir, "data.relocated-1")); + await seedMarker(paths, "alpha", { target, direction: "out", phase: "swap" }); + + await assert.rejects( + () => + relocateChannelMedia({ + io: TEST_IO, + paths, + slug: "alpha", + direction: "out", + root, + onLog: () => {}, + }), + /Verification failed/, + ); + assert.equal( + await readFile(path.join(target, "v1", "newer.json"), "utf8"), + "{}", + "nothing on the target was deleted to match the parked copy", + ); + }); +}); + +// --------------------------------------------------------------------------- +// Release 16 slice RM — a move holds the writers and mirrors its copy +// --------------------------------------------------------------------------- + +const TRANSCRIBING: ChannelWriter = { + source: "job", + jobId: "J1", + kind: "whisper-all", + label: "Transcribe all", + status: "running", + videoId: "v50t5yt", + taskKind: "transcribe", +}; + +// THE REFUSAL OVER A RUNNING JOB, the job's first step. The registry is +// injected: what is pinned is that the move asks, refuses naming the writer, +// and has touched nothing — no marker, no target directory, the source as it +// was. +test("a move refuses to start over a running job, naming it, and touches nothing", async () => { + await withTmp(async (paths, root) => { + const channelDir = await seed(paths, "alpha", { + v1: { "audio.m4a": "one" }, + }); + await assert.rejects( + () => + relocateChannelMedia({ + io: TEST_IO, + paths, + slug: "alpha", + direction: "out", + root, + onLog: () => {}, + writers: () => [TRANSCRIBING], + }), + /Cannot move the media of "alpha" now: a transcription of v50t5yt is running \(Transcribe all, job J1\) — wait for it or cancel it\. Nothing has been touched\./, + ); + assert.equal(await readRelocationMarker(paths, "alpha"), null); + await assert.rejects(() => stat(path.join(root, "alpha"))); + assert.ok((await lstat(path.join(channelDir, "data"))).isDirectory()); + }); +}); + +test("the preview gives the same refusal before the operator commits", async () => { + await withTmp(async (paths, root) => { + await seed(paths, "alpha", { v1: { "audio.m4a": "one" } }); + await assert.rejects( + () => + previewRelocation({ + paths, + slug: "alpha", + root, + writers: () => [TRANSCRIBING], + }), + /a transcription of v50t5yt is running/, + ); + // And with nothing writing, the preview answers. + const ok = await previewRelocation({ + paths, + slug: "alpha", + root, + writers: () => [], + }); + assert.equal(ok.files, 1); + }); +}); + +// A WRITER THAT STARTED AFTER THE FIRST STEP, before the marker landed. The +// third ask, right after the marker is written, sees it; the refusal removes the +// marker this run created, so the channel is not left held by a move that never +// copied a byte. +test("a writer seen once the marker is written refuses, and a fresh move's marker is removed", async () => { + await withTmp(async (paths, root) => { + const channelDir = await seed(paths, "alpha", { + v1: { "audio.m4a": "one" }, + }); + let asked = 0; await assert.rejects( () => relocateChannelMedia({ @@ -1043,18 +1219,155 @@ test("out @ swap: a file the target is missing is still a refusal", async () => direction: "out", root, onLog: () => {}, + writers: () => (++asked === 1 ? [] : [TRANSCRIBING]), }), - /The source has NOT been touched/, + /a transcription of v50t5yt is running/, + ); + assert.equal(asked, 2); + assert.equal(await readRelocationMarker(paths, "alpha"), null); + assert.equal( + (await inspectChannelMedia(paths, "alpha", undefined, { fresh: true })).status, + "in-place", + ); + assert.deepEqual( + await readdir(path.join(relocatedDataDir(root, "alpha"))), + [], + "nothing was copied", ); + assert.ok((await lstat(path.join(channelDir, "data"))).isDirectory()); + }); +}); + +// THE 2026-10-01 CASE. A move killed mid-copy left a transcriber's scratch dir +// on the destination that has since gone from the source. The resume used to +// copy everything else and refuse on the counts (1755 against 1750), and no +// rerun could settle it; the mirror pass now does. +test("a resume with a stale extra dir on the destination completes", async () => { + await withTmp(async (paths, root) => { + const channelDir = await seed(paths, "alpha", { + v50t5yt: { "audio.mp3": "a".repeat(64), "transcript.json": "{}" }, + v51fpcd: { "audio.mp3": "b".repeat(64) }, + }); + const dataDir = path.join(channelDir, "data"); + const target = relocatedDataDir(root, "alpha"); + await copyTree(dataDir, target); + const scratch = path.join(target, "v50t5yt", ".audio.mp3.parakeet"); + await mkdir(scratch, { recursive: true }); + for (const f of ["meta", "win-0000", "win-0001", "win-0002", "win-0003"]) { + await writeFile(path.join(scratch, `${f}.json`), "{}"); + } + await seedMarker(paths, "alpha", { target, direction: "out", phase: "copy" }); + + const res = await relocateChannelMedia({ + io: TEST_IO, + paths, + slug: "alpha", + direction: "out", + root, + onLog: () => {}, + }); + assert.equal(res.resumed, true); + assert.equal(res.files, 3); + await assert.rejects(() => stat(scratch)); + assert.equal((await inspectChannelMedia(paths, "alpha")).status, "ok"); + assert.equal(await readRelocationMarker(paths, "alpha"), null); + }); +}); + +test("back: a stale extra on the copy coming home is removed on resume", async () => { + await withTmp(async (paths, root) => { + const channelDir = await seed(paths, "alpha", { + v1: { "audio.m4a": "one".repeat(100) }, + }); + await relocateChannelMedia({ + io: TEST_IO, + paths, + slug: "alpha", + direction: "out", + root, + onLog: () => {}, + }); + const target = relocatedDataDir(root, "alpha"); + // An interrupted move back: `data.incoming` holds the copy, plus a file + // the target (the source of this direction) no longer has. + const incoming = path.join(channelDir, "data.incoming"); + await copyTree(target, incoming); + await writeFile(path.join(incoming, "v1", "gone.json"), "{}"); + await seedMarker(paths, "alpha", { target, direction: "back", phase: "copy" }); - // Nothing committed: `data/` is still the real directory, holding both files. + const res = await relocateChannelMedia({ + io: TEST_IO, + paths, + slug: "alpha", + direction: "back", + onLog: () => {}, + }); + assert.equal(res.resumed, true); + const dataDir = path.join(channelDir, "data"); assert.ok((await lstat(dataDir)).isDirectory()); - assert.ok(!(await lstat(dataDir)).isSymbolicLink()); + assert.deepEqual(await readdir(path.join(dataDir, "v1")), ["audio.m4a"]); + assert.equal(await readRelocationMarker(paths, "alpha"), null); + }); +}); + +// RECONCILE AND RESUME — the remediation (the ruling's last bullet). An extra +// file and a changed one on the destination: the job says what it found, by +// kind, makes the copy match the source and finishes the move. +test("reconcile: an extra and a changed file on the destination are settled, and the move completes", async () => { + await withTmp(async (paths, root) => { + const channelDir = await seed(paths, "alpha", { + v1: { "audio.m4a": "one".repeat(100), "transcript.json": '{"v":2}' }, + }); + const dataDir = path.join(channelDir, "data"); + const target = relocatedDataDir(root, "alpha"); + await copyTree(dataDir, target); + await writeFile(path.join(target, "v1", "stale.json"), "{}"); + await writeFile(path.join(target, "v1", "transcript.json"), '{"v":1}'); + await seedMarker(paths, "alpha", { target, direction: "out", phase: "copy" }); + + const lines: string[] = []; + const res = await relocateChannelMedia({ + io: TEST_IO, + paths, + slug: "alpha", + direction: "out", + root, + reconcile: true, + onLog: (m) => lines.push(m), + }); + assert.equal(res.resumed, true); + const said = lines.find((l) => l.startsWith("Reconciling:")) ?? ""; + assert.match(said, /1 extra on the destination \(v1\/stale\.json\)/); + assert.match(said, /1 changed \(v1\/transcript\.json\)/); + assert.deepEqual( + (await readdir(path.join(target, "v1"))).sort(), + ["audio.m4a", "transcript.json"], + ); assert.equal( - await readFile(path.join(dataDir, "v1", "audio.m4a"), "utf8"), - "one".repeat(500), + await readFile(path.join(target, "v1", "transcript.json"), "utf8"), + '{"v":2}', + ); + assert.equal((await inspectChannelMedia(paths, "alpha")).status, "ok"); + assert.equal(await readRelocationMarker(paths, "alpha"), null); + }); +}); + +test("reconcile: with no marker there is nothing to reconcile", async () => { + await withTmp(async (paths, root) => { + await seed(paths, "alpha", { v1: { "audio.m4a": "one" } }); + await assert.rejects( + () => + relocateChannelMedia({ + io: TEST_IO, + paths, + slug: "alpha", + direction: "out", + root, + reconcile: true, + onLog: () => {}, + }), + /no relocation marker — there is no interrupted move to reconcile/, ); - assert.equal((await readChannelConfig(paths, "alpha"))?.dataDir, undefined); }); }); diff --git a/common/controller/relocateChannelMedia.ts b/common/controller/relocateChannelMedia.ts @@ -14,18 +14,22 @@ import { import type { Paths } from "../lib/paths"; import { clearDirMarker, - COPY_ARGS, + copyMirrorVerify, isDirectory, linkOrDirState, makeProgressSink, measureTree, pathExists, readDirMarker, - rsyncTree, verifyCopy, writeDirMarker, type RelocationProgress, } from "./relocateDir"; +import { + channelWriters, + channelWritersRefusal, + type ChannelWriter, +} from "./channelWriters"; import { isSocialChannel } from "../lib/channelConfig"; import { locationOfDataDir, @@ -91,10 +95,11 @@ export type RelocateChannelMediaResult = { // True when the run picked up an interrupted one from its marker rather than // starting from scratch. resumed: boolean; - // True when the verify found ONLY directory-mtime drift and one extra - // `rsync -a` pass settled it — worth saying on the job's done line, because - // an operator who has seen this refuse a copy should be told it did not this - // time. Content drift is still a refusal, never this. + // True when the verify found the source had changed after the mirror pass + // (a directory timestamp, or a file written or removed) and one more mirror + // pass settled it — worth saying on the job's done line, because an operator + // who has seen this refuse a copy should be told it did not this time. A + // second change is a refusal, never this. retried: boolean; }; @@ -129,8 +134,39 @@ type RelocateOpts = { // // Production passes nothing and gets `getSettings()`, unchanged. io?: { read: () => SiteSettings }; + // RECONCILE AND RESUME: finish an interrupted or refused move by making the + // destination copy match the source — the mirror pass, the verify, then the + // swap and the reclaim — with no copy pass first. Refused without a marker + // to finish. The panel's remediation for a move whose verify failed; the + // operator never deletes a file on the destination by hand. + reconcile?: boolean; + // WHO IS WRITING INTO THE CHANNEL, injectable for the tests. Production asks + // the live job registry and the live lane runners (channelWriters.ts), + // leaving out this job itself. + writers?: (slug: string) => ChannelWriter[]; }; +// The move's own kind is not a writer to refuse over: it is the asker, running +// on this channel's slug, and the relocation queue runs one move at a time. +const RELOCATE_KIND = "relocate-channel-media"; + +function liveWriters(slug: string): ChannelWriter[] { + return channelWriters(slug, { ignoreKinds: [RELOCATE_KIND] }); +} + +// A MOVE NEVER STARTS OVER A WRITER, and never waits silently for one either: it +// refuses, naming the job (release 16 slice RM). Asked by the preview, by the +// job's first step, and again the moment the marker is written — after which +// the lanes skip the channel and a media job refuses to start on it, so the +// third answer is the one that cannot go stale. +function assertNoWriters( + slug: string, + writers: (slug: string) => ChannelWriter[], +): void { + const refusal = channelWritersRefusal(slug, writers(slug)); + if (refusal) throw new Error(refusal); +} + export type RelocationPreview = { slug: string; target: string; @@ -433,10 +469,14 @@ export async function previewRelocation({ paths, slug, root, + writers = (s) => channelWriters(s), }: { paths: Paths; slug: string; root: string; + // Test seam; the live registry and lanes otherwise. A running move of this + // channel counts here — the preview is not that move. + writers?: (slug: string) => ChannelWriter[]; }): Promise<RelocationPreview> { // The preview REFUSES rather than reporting numbers for a root the job will // reject: its whole job is to answer "is this root usable" before the @@ -444,6 +484,8 @@ export async function previewRelocation({ // the worst possible answer. const problem = await relocationRootProblem({ paths, slug, root }); if (problem) throw new Error(problem); + // The same refusal the job's first step gives, before the operator commits. + assertNoWriters(slug, writers); // EXISTENCE, HERE AS WELL AS IN THE JOB. getFreeBytes walks up to the nearest // existing ancestor, so a typo'd root statfs's its parent and previews with // perfectly plausible numbers — and the move is then refused by the job, after @@ -496,9 +538,23 @@ export async function relocateChannelMedia( ); } + // THE JOB'S FIRST STEP: nothing may be writing into the channel. The action + // that enqueued this asked too, but a move can wait a long time in the + // relocation queue, and on 2026-09-30 a transcription started in exactly + // that wait. + const writers = opts.writers ?? liveWriters; + assertNoWriters(slug, writers); + const channelDir = path.join(paths.channelsDir, slug); const dataDir = path.join(channelDir, "data"); const existingMarker = await readMarkerRaw(paths, slug); + if (opts.reconcile && !existingMarker) { + throw new Error( + `Channel "${slug}" has no relocation marker — there is no interrupted ` + + `move to reconcile.`, + ); + } + const reconcile = opts.reconcile === true; if (direction === "back") { // THE MARKER IS THE SECOND SOURCE OF TRUTH FOR THE TARGET, and without it @@ -542,6 +598,8 @@ export async function relocateChannelMedia( signal, resumed: Boolean(resume), phase: resume?.phase ?? "copy", + writers, + reconcile, }); } @@ -595,9 +653,29 @@ export async function relocateChannelMedia( // "this run continues an interrupted move OUT" and nothing looser. resumed: Boolean(resume), phase: resume?.phase ?? "copy", + writers, + reconcile, }); } +// THE THIRD ASK, right after the copy phase's marker is written. From that +// write on, the lanes skip the channel (their next inspect reads the marker) +// and a media job refuses to start on it, so a writer seen now was already +// running when the marker landed — one that started between the first step and +// here. The refusal puts the channel back as it was: a marker this run created +// is removed (nothing has been copied under it), one it resumed is left. +async function assertNoWritersUnderMarker(args: { + paths: Paths; + slug: string; + resumed: boolean; + writers: (slug: string) => ChannelWriter[]; +}): Promise<void> { + const refusal = channelWritersRefusal(args.slug, args.writers(args.slug)); + if (!refusal) return; + if (!args.resumed) await clearMarker(args.paths, args.slug); + throw new Error(refusal); +} + async function moveOut(args: { paths: Paths; slug: string; @@ -612,6 +690,8 @@ async function moveOut(args: { signal?: AbortSignal; resumed: boolean; phase: RelocationPhase; + writers: (slug: string) => ChannelWriter[]; + reconcile: boolean; }): Promise<RelocateChannelMediaResult> { const { paths, slug, channelDir, dataDir, root, target, log, signal } = args; @@ -708,33 +788,36 @@ async function moveOut(args: { startedAt: new Date().toISOString(), phase: "copy", }); - // --partial keeps an aborted transfer resumable; -a preserves mtimes, which - // is what makes the LMDB index a no-op afterwards. - const { exitCode } = await rsyncTree({ - rsyncBin: paths.rsyncBin, - src: dataDir, - dest: target, - args: COPY_ARGS, - log, - progress: makeProgressSink({ - totalBytes: measured.bytes, - alreadyBytes: already, - log, - onProgress: args.onProgress, - }), - signal, + await assertNoWritersUnderMarker({ + paths, + slug, + resumed: args.resumed, + writers: args.writers, }); - if (signal?.aborted) { - throw new Error( - `Cancelled. ${dataDir} is untouched and the partial copy at ${target} ` + - `is resumable — rerun to continue.`, - ); - } - if (exitCode !== 0) throw new Error(`rsync failed (exit ${exitCode})`); - - log("Verifying the copy…"); + // Copy (--partial keeps an aborted transfer resumable; -a preserves mtimes, + // which is what makes the LMDB index a no-op afterwards), mirror toward + // the target, verify — relocateDir.ts's copyMirrorVerify, the same for a + // fresh move, a resume and a reconcile. verifyRetried ||= ( - await verifyCopy({ rsyncBin: paths.rsyncBin, src: dataDir, dest: target, log, signal }) + await copyMirrorVerify({ + rsyncBin: paths.rsyncBin, + src: dataDir, + dest: target, + log, + progress: makeProgressSink({ + totalBytes: measured.bytes, + alreadyBytes: already, + log, + onProgress: args.onProgress, + }), + signal, + reconcile: args.reconcile, + cancelled: () => + new Error( + `Cancelled. ${dataDir} is untouched and the partial copy at ${target} ` + + `is resumable — rerun to continue.`, + ), + }) ).retried; phase = "swap"; } @@ -759,12 +842,27 @@ async function moveOut(args: { // must not take the interrupted one's word for it. The source to verify // against is whichever copy of it still exists; once the swap has committed // there is none, and there is nothing left to check. + // + // MIRRORED ONLY FROM THE LIVE MEDIA. While `data/` is still the real + // directory it is the source, and a difference gets the mirror pass a copy + // phase would give it. A parked `data.relocated-*` is not live any more — + // the link already points at the target — so the target is never mirrored + // FROM it: that verify is the strict one, as it always was. const verifySrc = state.kind === "real-dir" ? dataDir : (parked[0] ?? null); if (verifySrc) { log("Verifying the copy…"); verifyRetried ||= ( - await verifyCopy({ rsyncBin: paths.rsyncBin, src: verifySrc, dest: target, log, signal }) + await verifyCopy({ + rsyncBin: paths.rsyncBin, + src: verifySrc, + dest: target, + log, + signal, + ...(verifySrc === dataDir + ? { mirror: { underConstruction: target } } + : {}), + }) ).retried; } @@ -856,6 +954,8 @@ async function moveBack(args: { signal?: AbortSignal; resumed: boolean; phase: RelocationPhase; + writers: (slug: string) => ChannelWriter[]; + reconcile: boolean; }): Promise<RelocateChannelMediaResult> { const { paths, slug, channelDir, dataDir, target, log, signal } = args; const incoming = path.join(channelDir, "data.incoming"); @@ -954,32 +1054,36 @@ async function moveBack(args: { startedAt: new Date().toISOString(), phase: "copy", }); - const { exitCode } = await rsyncTree({ - rsyncBin: paths.rsyncBin, - src: target, - dest: incoming, - args: COPY_ARGS, - log, - progress: makeProgressSink({ - totalBytes: measured.bytes, - // The same figure the space check above is priced in — what is already - // in `data.incoming` from an interrupted run. - alreadyBytes: already, - log, - onProgress: args.onProgress, - }), - signal, + await assertNoWritersUnderMarker({ + paths, + slug, + resumed: args.resumed, + writers: args.writers, }); - if (signal?.aborted) { - throw new Error( - `Cancelled. ${target} is untouched and ${incoming} is resumable — ` + - `rerun to continue.`, - ); - } - if (exitCode !== 0) throw new Error(`rsync failed (exit ${exitCode})`); - log("Verifying the copy…"); + // The copy under construction is `data.incoming`; the target on the other + // drive is the source, and is never the target of the mirror's --delete. verifyRetried ||= ( - await verifyCopy({ rsyncBin: paths.rsyncBin, src: target, dest: incoming, log, signal }) + await copyMirrorVerify({ + rsyncBin: paths.rsyncBin, + src: target, + dest: incoming, + log, + progress: makeProgressSink({ + totalBytes: measured.bytes, + // The same figure the space check above is priced in — what is + // already in `data.incoming` from an interrupted run. + alreadyBytes: already, + log, + onProgress: args.onProgress, + }), + signal, + reconcile: args.reconcile, + cancelled: () => + new Error( + `Cancelled. ${target} is untouched and ${incoming} is resumable — ` + + `rerun to continue.`, + ), + }) ).retried; phase = "swap"; } diff --git a/common/controller/relocateDir.test.ts b/common/controller/relocateDir.test.ts @@ -1,6 +1,30 @@ import { test } from "node:test"; import assert from "node:assert/strict"; -import { makeProgressSink, type RelocationProgress } from "./relocateDir"; +import { + chmod, + mkdir, + mkdtemp, + readFile, + readdir, + rm, + stat, + utimes, + writeFile, +} from "node:fs/promises"; +import { rmSync, writeFileSync } from "node:fs"; +import { tmpdir } from "node:os"; +import path from "node:path"; +import { + assertMirrorDirection, + classifyDrift, + copyMirrorVerify, + CopyVerificationError, + describeDifferences, + mirrorTree, + MirrorDirectionError, + makeProgressSink, + type RelocationProgress, +} from "./relocateDir"; // Run with: // pnpm --filter yt-dlp-transcript-common exec tsx --test controller/relocateDir.test.ts @@ -89,3 +113,292 @@ test("a line that is not a frame is ignored entirely", () => { assert.deepEqual(seen, []); assert.deepEqual(lines, []); }); + +// --------------------------------------------------------------------------- +// The mirror pass (release 16 slice RM) +// --------------------------------------------------------------------------- +// +// The REAL rsync, over temp dirs, as relocateChannelMedia.test.ts does. The +// copy-phase tests reach "after the first pass" through the log: rsyncTree +// logs each command line (`$ rsync …`) synchronously before it spawns the +// child, so a hook that changes the source on the mirror pass's line changes +// it after the copy pass has finished and before the mirror pass starts — +// deterministically, with no fake rsync. + +const MTIME = new Date("2021-03-04T05:06:07.000Z"); + +async function withTrees( + fn: (t: { src: string; dest: string; dir: string }) => Promise<void>, +): Promise<void> { + const dir = await mkdtemp(path.join(tmpdir(), "ttb-mirror-")); + const src = path.join(dir, "src"); + const dest = path.join(dir, "dest"); + await mkdir(src, { recursive: true }); + await mkdir(dest, { recursive: true }); + try { + await fn({ src, dest, dir }); + } finally { + await rm(dir, { recursive: true, force: true }); + } +} + +async function put(root: string, rel: string, content: string): Promise<void> { + const file = path.join(root, rel); + await mkdir(path.dirname(file), { recursive: true }); + await writeFile(file, content); + await utimes(file, MTIME, MTIME); +} + +async function tree(root: string): Promise<string[]> { + const out: string[] = []; + const walk = async (rel: string) => { + for (const e of await readdir(path.join(root, rel), { withFileTypes: true })) { + const r = rel ? `${rel}/${e.name}` : e.name; + if (e.isDirectory()) await walk(r); + else out.push(r); + } + }; + await walk(""); + return out.sort(); +} + +function cancelled(): Error { + return new Error("cancelled"); +} + +test("a file added to the source after the copy pass is carried by the mirror pass", async () => { + await withTrees(async ({ src, dest }) => { + await put(src, "v1/audio.mp3", "one".repeat(100)); + let mirrored = false; + const res = await copyMirrorVerify({ + rsyncBin: "rsync", + src, + dest, + cancelled, + log: (m) => { + // The transcription that finished after rsync had passed v1/. + if (m.startsWith("$ ") && m.includes("--delete") && !m.includes("--dry-run") && !mirrored) { + mirrored = true; + writeFileSync(path.join(src, "v1", "transcript.json"), "{}"); + } + }, + }); + assert.equal(mirrored, true, "the hook ran before the mirror pass"); + assert.equal(res.retried, false, "the mirror pass, not a retry, carried it"); + assert.equal(res.files, 2); + assert.equal(await readFile(path.join(dest, "v1", "transcript.json"), "utf8"), "{}"); + }); +}); + +test("a file removed from the source after the copy pass is removed from the destination", async () => { + await withTrees(async ({ src, dest }) => { + await put(src, "v1/audio.mp3", "one".repeat(100)); + await put(src, "v1/.audio.mp3.parakeet/meta.json", "{}"); + await put(src, "v1/.audio.mp3.parakeet/win-0000.json", "[]"); + let removed = false; + await copyMirrorVerify({ + rsyncBin: "rsync", + src, + dest, + cancelled, + log: (m) => { + // The transcriber cleaning up its scratch dir once it finished. + if (m.startsWith("$ ") && m.includes("--delete") && !m.includes("--dry-run") && !removed) { + removed = true; + rmSync(path.join(src, "v1", ".audio.mp3.parakeet"), { recursive: true }); + } + }, + }); + assert.equal(removed, true); + assert.deepEqual(await tree(dest), ["v1/audio.mp3"]); + // The source is what it was left as — never the target of --delete. + assert.deepEqual(await tree(src), ["v1/audio.mp3"]); + }); +}); + +test("a stale extra dir on the destination (the realcandaceo case) is removed, the source untouched", async () => { + await withTrees(async ({ src, dest }) => { + await put(src, "v50t5yt/audio.mp3", "a".repeat(64)); + await put(src, "v50t5yt/transcript.json", "{}"); + // What the interrupted move left: the full copy, plus a transcriber's + // scratch dir copied mid-transcription and since deleted from the source. + await put(dest, "v50t5yt/audio.mp3", "a".repeat(64)); + await put(dest, "v50t5yt/transcript.json", "{}"); + for (const f of ["meta", "win-0000", "win-0001", "win-0002", "win-0003"]) { + await put(dest, `v50t5yt/.audio.mp3.parakeet/${f}.json`, "{}"); + } + const res = await copyMirrorVerify({ + rsyncBin: "rsync", + src, + dest, + cancelled, + log: () => {}, + }); + assert.equal(res.files, 2); + assert.deepEqual(await tree(dest), ["v50t5yt/audio.mp3", "v50t5yt/transcript.json"]); + assert.deepEqual(await tree(src), ["v50t5yt/audio.mp3", "v50t5yt/transcript.json"]); + }); +}); + +test("one change between the mirror and the check gets one more pass", async () => { + await withTrees(async ({ src, dest }) => { + await put(src, "v1/audio.mp3", "one"); + let wrote = false; + const lines: string[] = []; + const res = await copyMirrorVerify({ + rsyncBin: "rsync", + src, + dest, + cancelled, + log: (m) => { + lines.push(m); + if (m.startsWith("$ ") && m.includes("--dry-run") && !wrote) { + wrote = true; + writeFileSync(path.join(src, "v1", "diarization.json"), "{}"); + } + }, + }); + assert.equal(res.retried, true); + assert.ok(lines.some((l) => l.includes("running one more mirror pass"))); + assert.ok(lines.some((l) => l.includes("1 missing on the destination (v1/diarization.json)"))); + assert.deepEqual(await tree(dest), ["v1/audio.mp3", "v1/diarization.json"]); + }); +}); + +test("a source that changes again after the second pass refuses, naming the paths by kind", async () => { + await withTrees(async ({ src, dest }) => { + await put(src, "v1/audio.mp3", "one"); + let n = 0; + await assert.rejects( + () => + copyMirrorVerify({ + rsyncBin: "rsync", + src, + dest, + cancelled, + log: (m) => { + // Something that keeps writing: a new file before every check. + if (m.startsWith("$ ") && m.includes("--dry-run")) { + n++; + writeFileSync(path.join(src, "v1", `chunk-${n}.json`), "{}"); + } + }, + }), + (err: unknown) => { + assert.ok(err instanceof CopyVerificationError); + assert.deepEqual(err.differences.missing, ["v1/chunk-2.json"]); + assert.match(err.message, /1 missing on the destination \(v1\/chunk-2\.json\)/); + assert.match(err.message, /0 extra on the destination/); + assert.match(err.message, /Reconcile and resume/); + assert.match(err.message, /The source has NOT been touched/); + return true; + }, + ); + // Every source file is still there. + assert.deepEqual(await tree(src), ["v1/audio.mp3", "v1/chunk-1.json", "v1/chunk-2.json"]); + }); +}); + +test("a reconcile lists the differences by kind first, then mirrors without a copy pass", async () => { + await withTrees(async ({ src, dest }) => { + await put(src, "v1/audio.mp3", "one".repeat(10)); + await put(src, "v1/transcript.json", '{"new":true}'); + await put(dest, "v1/audio.mp3", "one".repeat(10)); + await put(dest, "v1/transcript.json", "{}"); + await put(dest, "v1/stale.json", "{}"); + const lines: string[] = []; + await copyMirrorVerify({ + rsyncBin: "rsync", + src, + dest, + cancelled, + reconcile: true, + log: (m) => lines.push(m), + }); + const said = lines.find((l) => l.startsWith("Reconciling:")) ?? ""; + assert.match(said, /1 extra on the destination \(v1\/stale\.json\)/); + assert.match(said, /0 missing on the destination/); + assert.match(said, /1 changed \(v1\/transcript\.json\)/); + // No copy pass: the only non-dry-run rsync is the mirror. + assert.equal( + lines.filter((l) => l.startsWith("$ ") && l.includes("--partial")).length, + 0, + ); + assert.deepEqual(await tree(dest), ["v1/audio.mp3", "v1/transcript.json"]); + assert.equal(await readFile(path.join(dest, "v1", "transcript.json"), "utf8"), '{"new":true}'); + }); +}); + +// THE ONE DIRECTION --delete MAY RUN IN, asserted before rsync is spawned. The +// fake rsync here records that it ran; it must never be reached. +test("--delete is refused toward anything but the copy under construction", async () => { + await withTrees(async ({ src, dest, dir }) => { + await put(src, "v1/audio.mp3", "one"); + await put(dest, "v1/other.mp3", "two"); + const ran = path.join(dir, "rsync-ran"); + const fake = path.join(dir, "fake-rsync.sh"); + await writeFile(fake, `#!/bin/sh\ntouch "${ran}"\nexit 0\n`); + await chmod(fake, 0o755); + + // SWAPPED: the source named as the destination — the live media as the + // target of --delete. + await assert.rejects( + () => + mirrorTree({ + rsyncBin: fake, + src: dest, + dest: src, + underConstruction: dest, + log: () => {}, + }), + MirrorDirectionError, + ); + // One inside the other, either way round. + await assert.rejects( + () => + assertMirrorDirection({ + src, + dest: path.join(src, "v1"), + underConstruction: path.join(src, "v1"), + }), + /one is inside the other/, + ); + await assert.rejects( + () => assertMirrorDirection({ src: path.join(dest, "x"), dest, underConstruction: dest }), + /one is inside the other/, + ); + await assert.rejects(() => stat(ran), "the fake rsync was never spawned"); + // Both trees as they were. + assert.deepEqual(await tree(src), ["v1/audio.mp3"]); + assert.deepEqual(await tree(dest), ["v1/other.mp3"]); + + // The right way round is allowed, and does what it says. + await assertMirrorDirection({ src, dest, underConstruction: dest }); + }); +}); + +test("rsync's itemized lines sort into extra, missing and changed", () => { + const d = classifyDrift([ + "*deleting v50t5yt/.audio.mp3.parakeet/meta.json", + "*deleting v50t5yt/.audio.mp3.parakeet/", + ">f+++++++++ v51/transcript.json", + "cd+++++++++ v52/", + ">f.st...... v1/transcript.json", + ".d..t...... v1/", + "rsync: something it wanted to say", + ]); + assert.deepEqual(d.extra, [ + "v50t5yt/.audio.mp3.parakeet/meta.json", + "v50t5yt/.audio.mp3.parakeet/", + ]); + assert.deepEqual(d.missing, ["v51/transcript.json", "v52/"]); + assert.deepEqual(d.changed, [ + "v1/transcript.json", + "v1/", + "rsync: something it wanted to say", + ]); + assert.equal( + describeDifferences({ extra: [], missing: ["a", "b", "c", "d", "e", "f", "g"], changed: [] }), + "0 extra on the destination, 7 missing on the destination (a, b, c, d, e, … 2 more), 0 changed", + ); +}); diff --git a/common/controller/relocateDir.ts b/common/controller/relocateDir.ts @@ -4,6 +4,7 @@ import { readdir, readFile, readlink, + realpath, rm, stat, } from "node:fs/promises"; @@ -33,9 +34,11 @@ import type { // // What IS shared is the part the omnimirror incident and three rounds of resume // bugs were paid for: measure the tree, copy with `--partial` so an abort is -// resumable, verify with a dry run AND a re-measure, tolerate directory-mtime -// drift exactly once, and dispatch every step past the copy on what is ON DISK -// rather than on what the previous step was supposed to have done. +// resumable, MIRROR the copy (`--delete` toward the destination only — the +// realcandaceo incident, 2026-09-30), verify with a dry run AND a re-measure, +// allow the source one change between the mirror and the check, and dispatch +// every step past the copy on what is ON DISK rather than on what the previous +// step was supposed to have done. // // THE SOURCE IS NEVER TOUCHED UNTIL THE COPY VERIFIES. That is the invariant // every caller inherits and none may weaken. @@ -321,6 +324,29 @@ export async function rsyncTree(opts: { // the sink above reads. export const COPY_ARGS = ["-a", "--partial", "--info=progress2"]; +// THE MIRROR PASS: the destination copy is made to match the source exactly — +// what is missing is sent, what changed is re-sent, and what the source no +// longer has is DELETED FROM THE DESTINATION (release 16 slice RM). +// +// Why a copy needs it. The copy pass (`COPY_ARGS`) never deletes, so a file +// that existed on the source while rsync passed it and was removed afterwards +// stays on the destination forever. On 2026-09-30 that was a transcriber's +// scratch directory (`v50t5yt/.audio.mp3.parakeet/`, five files) copied mid- +// transcription and deleted from the source when the transcription finished: +// the resume the next day copied everything else and refused on the counts, +// 1755 files against 1750, with the source untouched — and nothing in the move +// could ever settle it. `--delete` is what can. +// +// `--delete` IS ONLY EVER POINTED AT THE COPY UNDER CONSTRUCTION, never at the +// source: run the other way round it would delete from the channel's live +// media every file the partial copy has not reached yet. `mirrorTree` refuses +// any other direction before rsync is spawned (assertMirrorDirection). +export const MIRROR_ARGS = ["-a", "--delete"]; + +// The verify's dry run: what the mirror pass WOULD change, itemized, deletions +// included. +const VERIFY_ARGS = ["-a", "--dry-run", "--itemize-changes", "--delete"]; + // A DIRECTORY MTIME IS NOT CONTENT. `.d..t` is rsync's itemization for "this is // a directory and only its modification time differs" — nothing to send, and no // byte of the copy is in question. It is what the omnimirror move hit @@ -338,17 +364,276 @@ function driftLines(output: string): string[] { .filter((l) => !/^(sent|total size|$)/.test(l)); } +// --------------------------------------------------------------------------- +// What differs, by kind +// --------------------------------------------------------------------------- + +// The differences a verify's dry run reports, sorted into what the operator can +// act on: a path the destination has and the source does not ("extra on the +// destination" — the realcandaceo scratch dir), one the destination lacks +// ("missing on the destination"), and one both have that differs ("changed", +// a directory's timestamp included). Paths are relative to the tree. +export type CopyDifferences = { + extra: string[]; + missing: string[]; + changed: string[]; +}; + +// rsync's itemize line is `YXcstpoguax path`: Y the update type, X the file +// type, then one column per attribute, `+` in every column for an item that +// does not exist on the receiving side yet. A deletion is `*deleting path`. +// Anything else (a warning rsync printed) is reported whole, as changed. +export function classifyDrift(lines: ReadonlyArray<string>): CopyDifferences { + const out: CopyDifferences = { extra: [], missing: [], changed: [] }; + for (const line of lines) { + if (line.startsWith("*deleting")) { + out.extra.push(line.replace(/^\S+\s+/, "")); + } else if (/^[<>ch.][fdLDS]\+{5,}\s/.test(line)) { + out.missing.push(line.replace(/^\S+\s+/, "")); + } else if (/^[<>ch.][fdLDS]\S*\s/.test(line)) { + out.changed.push(line.replace(/^\S+\s+/, "")); + } else { + out.changed.push(line); + } + } + return out; +} + +function countDifferences(d: CopyDifferences): number { + return d.extra.length + d.missing.length + d.changed.length; +} + +// "2 extra on the destination (v1/.x/meta.json, v1/.x/win-0000.json), 0 +// missing on the destination, 1 changed (v2/transcript.json)" — at most five +// paths of each kind, the rest counted. +export function describeDifferences(d: CopyDifferences): string { + const kind = (n: string[], words: string) => { + if (n.length === 0) return `0 ${words}`; + const shown = n.slice(0, 5).join(", "); + const more = n.length > 5 ? `, … ${n.length - 5} more` : ""; + return `${n.length} ${words} (${shown}${more})`; + }; + return [ + kind(d.extra, "extra on the destination"), + kind(d.missing, "missing on the destination"), + kind(d.changed, "changed"), + ].join(", "); +} + +// A verify that the second mirror pass did not settle. Carries the differences +// by kind, so a caller (and a test) can read them without parsing prose. +export class CopyVerificationError extends Error { + readonly differences: CopyDifferences; + constructor(message: string, differences: CopyDifferences) { + super(message); + this.name = "CopyVerificationError"; + this.differences = differences; + } +} + +// --------------------------------------------------------------------------- +// The mirror pass, and the one direction it may run in +// --------------------------------------------------------------------------- + +async function realOrResolved(p: string): Promise<string> { + return await realpath(p).catch(() => path.resolve(p)); +} + +function isWithin(parent: string, child: string): boolean { + const rel = path.relative(parent, child); + return rel === "" || (!rel.startsWith("..") && !path.isAbsolute(rel)); +} + +export class MirrorDirectionError extends Error { + constructor(message: string) { + super(message); + this.name = "MirrorDirectionError"; + } +} + +// THE ONE RULE `--delete` LIVES UNDER, asserted before rsync is spawned: the +// destination is the copy under construction — the directory the move itself +// created to copy into (`<root>/<slug>/data` out, `data.incoming` back, +// `<root>/saved-videos` for the store) — and it is neither the source nor +// inside it, nor the other way round. A caller that swaps `src` and `dest` +// names the live media as the destination, which is not the copy it is +// building, and is refused here with nothing deleted. +export async function assertMirrorDirection(opts: { + src: string; + dest: string; + underConstruction: string; +}): Promise<void> { + if (path.resolve(opts.dest) !== path.resolve(opts.underConstruction)) { + throw new MirrorDirectionError( + `Refusing to mirror into ${opts.dest}: --delete may only target the copy ` + + `under construction (${opts.underConstruction}), never the source.`, + ); + } + const [src, dest] = await Promise.all([ + realOrResolved(opts.src), + realOrResolved(opts.dest), + ]); + if (isWithin(src, dest) || isWithin(dest, src)) { + throw new MirrorDirectionError( + `Refusing to mirror ${opts.src} into ${opts.dest}: one is inside the ` + + `other, so --delete would reach the source.`, + ); + } +} + +// One `rsync -a --delete` pass, source → the copy under construction. With a +// progress sink it reports like the copy pass (Reconcile and resume's pass is +// the copy). +export async function mirrorTree(opts: { + rsyncBin: string; + src: string; + dest: string; + underConstruction: string; + log: (m: string) => void; + progress?: (line: string) => void; + signal?: AbortSignal; +}): Promise<{ exitCode: number; output: string }> { + await assertMirrorDirection(opts); + return rsyncTree({ + rsyncBin: opts.rsyncBin, + src: opts.src, + dest: opts.dest, + args: opts.progress ? [...MIRROR_ARGS, "--info=progress2"] : MIRROR_ARGS, + log: opts.log, + progress: opts.progress, + signal: opts.signal, + }); +} + +// What the destination copy differs from the source by, now, without changing +// anything: Reconcile and resume says this before it mirrors. +export async function diffTrees(opts: { + rsyncBin: string; + src: string; + dest: string; + log: (m: string) => void; + signal?: AbortSignal; +}): Promise<CopyDifferences> { + const { exitCode, output } = await rsyncTree({ + ...opts, + args: VERIFY_ARGS, + // The listing is what this returns; the log gets the command line only. + log: (m) => { + if (m.startsWith("$ ")) opts.log(m); + }, + }); + if (exitCode !== 0) { + throw new Error(`The comparison rsync failed (exit ${exitCode})`); + } + return classifyDrift(driftLines(output)); +} + +// --------------------------------------------------------------------------- +// The copy phase, for both movers and both directions +// --------------------------------------------------------------------------- + +// COPY, MIRROR, VERIFY — the whole copy phase, in one place so the channel +// mover's two directions and the store mover's two cannot drift apart. +// +// 1. The copy pass (`COPY_ARGS`, progress-reporting, `--partial`): the bulk of +// the bytes, resumable. Skipped by a reconcile, whose mirror pass reports +// the progress instead. +// 2. The mirror pass (`MIRROR_ARGS`) toward the copy under construction: what +// changed on the source while pass 1 ran is sent, and what the source no +// longer has is removed from the copy. +// 3. The verify, in mirror mode: an empty itemized dry run AND equal counts; +// one more mirror pass if the source moved under it, and a refusal naming +// the differences by kind if it moved again. +// +// A resumed move takes the same path, so a destination left with stale files +// by an interrupted move is mirrored clean rather than refused forever. +export async function copyMirrorVerify(opts: { + rsyncBin: string; + src: string; + // The copy under construction, which the caller has already created. + dest: string; + log: (m: string) => void; + progress?: (line: string) => void; + signal?: AbortSignal; + // Reconcile and resume: no copy pass; the mirror pass carries the progress, + // and the differences it is about to settle are logged by kind first. + reconcile?: boolean; + // The caller's sentence for a cancel, which names its own two directories. + cancelled: () => Error; +}): Promise<{ bytes: number; files: number; retried: boolean }> { + const { rsyncBin, src, dest, log, signal } = opts; + if (opts.reconcile) { + const before = await diffTrees({ rsyncBin, src, dest, log, signal }); + log( + countDifferences(before) === 0 + ? "Reconciling: the destination copy already matches the source." + : `Reconciling: the destination copy differs from the source — ` + + `${describeDifferences(before)}. The mirror pass makes it match ` + + `the source; the source is not changed.`, + ); + } else { + const { exitCode } = await rsyncTree({ + rsyncBin, + src, + dest, + args: COPY_ARGS, + log, + progress: opts.progress, + signal, + }); + if (signal?.aborted) throw opts.cancelled(); + if (exitCode !== 0) throw new Error(`rsync failed (exit ${exitCode})`); + } + log("Mirroring the copy (rsync --delete toward the destination)…"); + const mirrored = await mirrorTree({ + rsyncBin, + src, + dest, + underConstruction: dest, + log, + progress: opts.reconcile ? opts.progress : undefined, + signal, + }); + if (signal?.aborted) throw opts.cancelled(); + if (mirrored.exitCode !== 0) { + throw new Error( + `The mirror pass failed (exit ${mirrored.exitCode}). The source has NOT ` + + `been touched.`, + ); + } + log("Verifying the copy…"); + return verifyCopy({ + rsyncBin, + src, + dest, + log, + signal, + mirror: { underConstruction: dest }, + }); +} + // What a copy has to clear before the swap: rsync itself agrees there is // nothing left to send, AND the two trees measure the same. The dry run alone // would accept a target that is byte-identical for the wrong reason; the counts // alone would accept two trees of equal size with different contents. +// +// IN MIRROR MODE (`mirror` given; every copy phase, and a swap-phase re-verify +// whose source is still the live media) the dry run carries `--delete`, so an +// extra file on the destination is a difference too, and ANY difference gets +// one more mirror pass — a source change seen between the pass and the check — +// before a second difference refuses with the paths by kind. Without `mirror` +// (a re-verify against a parked copy, which is no longer the live media and +// must never be mirrored from) it is the verify as it always was: only +// directory timestamps get a second pass, plain `rsync -a`. export async function verifyCopy(opts: { rsyncBin: string; src: string; dest: string; log: (m: string) => void; signal?: AbortSignal; + mirror?: { underConstruction: string }; }): Promise<{ bytes: number; files: number; retried: boolean }> { + if (opts.mirror) return verifyMirrored({ ...opts, mirror: opts.mirror }); let retried = false; // ONE retry, never a loop: if a second pass does not settle it, something is // still writing into the tree and the answer is to refuse, not to chase it. @@ -382,10 +667,17 @@ export async function verifyCopy(opts: { `(first: ${drift[0]}). The source has NOT been touched.`, ); } - const [a, b] = await Promise.all([ - measureTree(opts.src), - measureTree(opts.dest), - ]); + return measuredEqual(opts.src, opts.dest, retried); +} + +// The counts half of every verify: the two trees hold the same number of files +// and the same bytes. +async function measuredEqual( + src: string, + dest: string, + retried: boolean, +): Promise<{ bytes: number; files: number; retried: boolean }> { + const [a, b] = await Promise.all([measureTree(src), measureTree(dest)]); if (a.files !== b.files || a.bytes !== b.bytes) { throw new Error( `Verification failed: source has ${a.files} file(s)/${formatBytes(a.bytes)}, ` + @@ -394,3 +686,72 @@ export async function verifyCopy(opts: { } return { ...a, retried }; } + +async function verifyMirrored(opts: { + rsyncBin: string; + src: string; + dest: string; + log: (m: string) => void; + signal?: AbortSignal; + mirror: { underConstruction: string }; +}): Promise<{ bytes: number; files: number; retried: boolean }> { + const { rsyncBin, src, dest, log, signal } = opts; + let retried = false; + // ONE more pass, never a loop: a source that changes again between that + // pass and the next check is still being written to, and the answer is to + // refuse and name what moved, not to chase it. + for (;;) { + const { exitCode, output } = await rsyncTree({ + rsyncBin, + src, + dest, + args: VERIFY_ARGS, + log, + signal, + }); + if (exitCode !== 0) { + throw new Error(`Verification rsync failed (exit ${exitCode})`); + } + const drift = driftLines(output); + if (drift.length === 0) break; + const differences = classifyDrift(drift); + if (!retried) { + retried = true; + log( + drift.every((l) => DIR_MTIME_ONLY.test(l)) + ? `Verification found ${drift.length} directory timestamp(s) differing and ` + + `no content drift — running one more mirror pass to settle them.` + : `Verification found the source changed during the copy — ` + + `${describeDifferences(differences)} — running one more mirror pass.`, + ); + const again = await mirrorTree({ + rsyncBin, + src, + dest, + underConstruction: opts.mirror.underConstruction, + log, + signal, + }); + if (again.exitCode !== 0) { + throw new Error( + `Verification rsync failed (exit ${again.exitCode}). ` + + `The source has NOT been touched.`, + ); + } + continue; + } + const first = [ + ...differences.extra, + ...differences.missing, + ...differences.changed, + ][0]; + throw new CopyVerificationError( + `Verification failed: after a second mirror pass the copy still differs ` + + `from the source — ${describeDifferences(differences)} (first: ${first}). ` + + `Something is still writing into ${src}: stop it, then Reconcile and ` + + `resume. The source has NOT been touched.`, + differences, + ); + } + return measuredEqual(src, dest, retried); +} diff --git a/common/controller/relocateSavedVideos.test.ts b/common/controller/relocateSavedVideos.test.ts @@ -212,6 +212,58 @@ test("an interrupted copy resumes from its marker", async () => { }); }); +// THE CHANNEL MOVER'S PIPELINE, SHARED (release 16 slice RM): a container on +// the partial copy that the store no longer has — unpersisted while the move +// was down — is mirrored away on resume instead of failing the counts forever. +test("a resume removes from the copy what the store no longer has", async () => { + await withTmp(async (h) => { + await seedStore(h.paths); + const target = relocatedSavedVideosDir(h.root); + await mkdir(path.join(target, "chan", "gone"), { recursive: true }); + await writeFile(path.join(target, "chan", "gone", "source-media.mp4"), "z"); + await writeDirMarker(savedVideosMarkerPath(h.paths), { + target, + direction: "out", + startedAt: new Date().toISOString(), + phase: "copy", + }); + const result = await relocateSavedVideos({ + paths: h.paths, + locationId: "cold", + io: h.io, + onLog: () => {}, + }); + assert.equal(result.files, 2); + assert.deepEqual((await readdir(path.join(target, "chan"))).sort(), [ + "vid1", + "vid2", + ]); + }); +}); + +// The job's first step: a store move never starts over a writer. +test("a store move refuses to start while something is writing into the store", async () => { + await withTmp(async (h) => { + await seedStore(h.paths); + await assert.rejects( + () => + relocateSavedVideos({ + paths: h.paths, + locationId: "cold", + io: h.io, + onLog: () => {}, + busy: () => "1 running download job(s) (sync).", + }), + /1 running download job\(s\) \(sync\)\. Nothing has been touched\./, + ); + assert.equal( + await readDirMarkerExists(savedVideosMarkerPath(h.paths)), + false, + ); + assert.equal((await lstat(h.paths.savedVideosDir)).isDirectory(), true); + }); +}); + // SAME TARGET, OPPOSITE INTENT. Finishing a move-out as a move-back would swap // the wrong way round, so it is refused by name rather than resumed. test("a marker for the other direction is refused, not resumed", async () => { diff --git a/common/controller/relocateSavedVideos.ts b/common/controller/relocateSavedVideos.ts @@ -32,13 +32,12 @@ import { } from "../lib/savedVideoStore"; import { clearDirMarker, - COPY_ARGS, + copyMirrorVerify, isDirectory, linkOrDirState, makeProgressSink, measureTree, readDirMarker, - rsyncTree, verifyCopy, writeDirMarker, type RelocationProgress, @@ -238,6 +237,13 @@ type Opts = { signal?: AbortSignal; // Injectable for unit tests, exactly as the storage controller does it. io?: { read: () => SiteSettings; write: (next: SiteSettings) => Promise<void> }; + // WHAT IS WRITING INTO THE STORE, asked as the job's first step (release 16 + // slice RM, the channel mover's rule): a sentence naming it, or null. The + // editor passes its store check (storage/lib/storeBusy.ts) — a store has no + // slug, so the question is the whole machine's and lives with the editor + // that knows its kinds. A bin script passes nothing; the marker remains the + // guard either way (`assertSavedVideosStoreWritable`). + busy?: () => string | null; }; const BYTES_PER_GB = 1024 ** 3; @@ -248,6 +254,10 @@ export async function relocateSavedVideos( const io = opts.io ?? { read: getSettings, write: writeSettings }; const log = opts.onLog ?? ((m: string) => console.log(m)); const paths = opts.paths; + // The job's first step: a move never starts over a writer, and never waits + // silently for one either. + const busy = opts.busy?.() ?? null; + if (busy) throw new Error(`${busy} Nothing has been touched.`); const markerFile = savedVideosMarkerPath(paths); const existing = await readDirMarker(markerFile); const wantBack = opts.locationId.trim() === ""; @@ -401,35 +411,28 @@ async function moveStoreOut(a: Inner): Promise<SavedVideosRelocateResult> { // What a previous attempt already landed — see the channel mover. const already = (await measureTree(target)).bytes; await stampMarker(markerFile, target, "out", "copy"); - const { exitCode } = await rsyncTree({ - rsyncBin: paths.rsyncBin, - src: store, - dest: target, - args: COPY_ARGS, - log, - progress: makeProgressSink({ - totalBytes: measured.bytes, - alreadyBytes: already, - log, - onProgress: a.onProgress, - }), - signal: a.signal, - }); - if (a.signal?.aborted) { - throw new Error( - `Cancelled. ${store} is untouched and the partial copy at ${target} is ` + - `resumable — rerun to continue.`, - ); - } - if (exitCode !== 0) throw new Error(`rsync failed (exit ${exitCode})`); - log("Verifying the copy…"); + // Copy, mirror toward the target, verify — the channel mover's pipeline + // (relocateDir.ts's copyMirrorVerify), so a persist that landed in the + // store during the copy is carried, and a container unpersisted during it + // is removed from the copy rather than failing the counts. retried ||= ( - await verifyCopy({ + await copyMirrorVerify({ rsyncBin: paths.rsyncBin, src: store, dest: target, log, + progress: makeProgressSink({ + totalBytes: measured.bytes, + alreadyBytes: already, + log, + onProgress: a.onProgress, + }), signal: a.signal, + cancelled: () => + new Error( + `Cancelled. ${store} is untouched and the partial copy at ${target} is ` + + `resumable — rerun to continue.`, + ), }) ).retried; phase = "swap"; @@ -455,6 +458,9 @@ async function moveStoreOut(a: Inner): Promise<SavedVideosRelocateResult> { dest: target, log, signal: a.signal, + // The store is still the live directory here, so a difference gets + // the mirror pass a copy phase would give it. + mirror: { underConstruction: target }, }) ).retried; await rename(store, parked); @@ -617,35 +623,26 @@ async function moveStoreBack(a: Inner): Promise<SavedVideosRelocateResult> { ); await mkdir(incoming, { recursive: true }); await stampMarker(markerFile, target, "back", "copy"); - const { exitCode } = await rsyncTree({ - rsyncBin: paths.rsyncBin, - src: target, - dest: incoming, - args: COPY_ARGS, - log, - progress: makeProgressSink({ - totalBytes: measured.bytes, - // The figure the space check above is priced in. - alreadyBytes: already, - log, - onProgress: a.onProgress, - }), - signal: a.signal, - }); - if (a.signal?.aborted) { - throw new Error( - `Cancelled. ${target} is untouched and ${incoming} is resumable — rerun to continue.`, - ); - } - if (exitCode !== 0) throw new Error(`rsync failed (exit ${exitCode})`); - log("Verifying the copy…"); + // Copy, mirror toward `incoming` (the copy under construction — the target + // on the drive is the source here, never the target of --delete), verify. retried ||= ( - await verifyCopy({ + await copyMirrorVerify({ rsyncBin: paths.rsyncBin, src: target, dest: incoming, log, + progress: makeProgressSink({ + totalBytes: measured.bytes, + // The figure the space check above is priced in. + alreadyBytes: already, + log, + onProgress: a.onProgress, + }), signal: a.signal, + cancelled: () => + new Error( + `Cancelled. ${target} is untouched and ${incoming} is resumable — rerun to continue.`, + ), }) ).retried; phase = "swap"; diff --git a/common/jobs/jobKinds.ts b/common/jobs/jobKinds.ts @@ -566,6 +566,63 @@ const JOB_KINDS: Record<string, JobKindMeta> = { replayable: false, queueKeyStrategy: "custom", }, + // THE PER-VIDEO WRITERS THAT WERE NOT IN THIS TABLE (release 16 slice RM). + // Each runs with a channelSlug and writes under `data/<id>/` — a single + // video's transcription (the video page's two Transcribe buttons, and its + // worker path), its download, its audio transcode, and the availability + // checks (`availability.json` per video; the quick check diffs the playlist + // against what is on disk) — and each, being absent, was never refused for + // a channel whose media is moving or unreachable. Labels stay absent so /jobs + // shows them exactly as before. + "whisper-video": { + kind: "whisper-video", + drainable: false, + replayable: false, + queueKeyStrategy: "custom", + needsMedia: true, + }, + "transcribe-one": { + kind: "transcribe-one", + drainable: false, + replayable: false, + queueKeyStrategy: "custom", + needsMedia: true, + }, + "download-one-pipeline": { + kind: "download-one-pipeline", + drainable: false, + replayable: false, + queueKeyStrategy: "custom", + needsMedia: true, + }, + "transcode-audio": { + kind: "transcode-audio", + drainable: false, + replayable: false, + queueKeyStrategy: "custom", + needsMedia: true, + }, + "check-availability": { + kind: "check-availability", + drainable: false, + replayable: false, + queueKeyStrategy: "platform", + needsMedia: true, + }, + "quick-availability-check": { + kind: "quick-availability-check", + drainable: false, + replayable: false, + queueKeyStrategy: "platform", + needsMedia: true, + }, + "check-maybe-missing": { + kind: "check-maybe-missing", + drainable: false, + replayable: false, + queueKeyStrategy: "platform", + needsMedia: true, + }, // Replayable kinds that never had a JOB_KIND_LABELS entry: label omitted so // jobKindLabel() keeps falling back to the raw kind (unchanged behavior). "store-playlist": { diff --git a/common/jobs/streamCommand.test.ts b/common/jobs/streamCommand.test.ts @@ -1,6 +1,6 @@ import { test } from "node:test"; import assert from "node:assert/strict"; -import { mkdtemp, readFile, rm } from "node:fs/promises"; +import { mkdir, mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; import { tmpdir } from "node:os"; import path from "node:path"; import type { Paths } from "../lib/paths"; @@ -168,6 +168,55 @@ test("serialWriter runs writes one at a time, in order, past a rejection", async ]); }); +// THE MEDIA GUARD, ASKED AGAIN AT THE START (release 16 slice RM). A media job +// queued while its channel was in place, behind other work on its queue, and +// started after a move wrote the channel's marker, must not write into the +// tree being copied: the guard that passed at enqueue is asked again when the +// queue starts the job, and the refusal — the hold's words — fails it before +// `fn` runs. +test("a media job queued before a move's marker and started after it refuses at its start", async () => { + const { paths: base, root } = await jobsDir(); + const paths = { ...base, channelsDir: path.join(root, "channels") } as Paths; + await mkdir(path.join(paths.channelsDir, "alpha", "data"), { recursive: true }); + const queueKey = `test:start-guard:${newJobId()}`; + const holder = await hold(paths, queueKey); + let ran = false; + try { + const job = ok( + await runManagedFunction({ + kind: "whisper-all", + queueKey, + paths, + channelSlug: "alpha", + fn: async () => { + ran = true; + }, + }), + ); + assert.equal(getRegistry().get(job.jobId)?.status, "queued"); + // The move starts while the job waits. + await writeFile( + path.join(paths.channelsDir, "alpha", ".relocating.json"), + JSON.stringify({ + target: "/mnt/platter/alpha/data", + direction: "out", + startedAt: new Date().toISOString(), + phase: "copy", + }), + ); + await holder.release(); + assert.equal((await job.done).status, "failed"); + assert.equal(ran, false, "fn never ran"); + const log = await readFile(path.join(paths.jobsDir, `${job.jobId}.log`), "utf8"); + assert.match( + log, + /Channel "alpha" is held: its media is moving \(a move is in progress or was interrupted\)/, + ); + } finally { + await rm(root, { recursive: true, force: true }); + } +}); + // THE ONE CANCEL THAT MUST NOT: the graceful-shutdown reaper cancels every // queued job only so the exit cannot promote one into a child. Nobody cancelled // it, and its `queued` sidecar is what the boot pass (bootQueuedJobs.ts) diff --git a/common/jobs/streamCommand.ts b/common/jobs/streamCommand.ts @@ -21,7 +21,11 @@ import { writeJobMeta } from "./jobMeta"; import { maybePruneJobLogs } from "./listJobs"; import type { JobSpec } from "./jobSpec"; import { kindNeedsMedia } from "./jobKinds"; -import { assertChannelMediaReachable } from "../lib/channelMedia"; +import { + assertChannelMediaReachable, + ChannelMediaUnreachableError, +} from "../lib/channelMedia"; +import { mediaHoldText } from "../lib/channelMediaHold"; // Mark a job's channel report dirty so the debounced scheduler regenerates the // snapshot — called both on each completed sub-operation and on the job's @@ -310,6 +314,16 @@ export async function runManagedCommand( // false, so a bookkeeping kind is never refused for a drive it does not read, // and the relocate job itself — the thing that FIXES an unreachable channel — // must never declare it. +// +// ASKED TWICE: here, before the record exists, and again when the queue STARTS +// the job (`start` below). A job can wait hours in its queue — a "Transcribe +// all" behind another channel's on the one transcription queue — and a move +// can begin in that wait; the answer it got at enqueue is then stale, and the +// job would write into a tree being copied. At the start the refusal fails the +// job with the sentence in its log. +// +// A MOVE IS A HOLD, and the refusal says so in the hold's words ("held: its +// media is moving …") rather than calling the media unreachable. async function refuseForUnreachableMedia( opts: CommonOpts, ): Promise<string | null> { @@ -318,6 +332,17 @@ async function refuseForUnreachableMedia( await assertChannelMediaReachable(opts.paths, opts.channelSlug); return null; } catch (err) { + if ( + err instanceof ChannelMediaUnreachableError && + err.status === "in-transition" + ) { + return ( + `Channel "${opts.channelSlug}" is ${mediaHoldText(err.status)} — ` + + `${err.location.detail ?? "a relocation marker is present"}. Its media ` + + `jobs start again when the move completes, or when its marker is ` + + `cleared on the channel's Storage panel.` + ); + } return (err as Error).message; } } @@ -398,8 +423,15 @@ export async function runManagedFunction( }, }; - opts - .fn(onLog, abort.signal, setProgress, ctx) + // THE GUARD AGAIN, NOW THAT THE QUEUE HAS STARTED THE JOB (see + // refuseForUnreachableMedia): a refusal here fails the job before `fn` + // touches anything, with the sentence as its log's `[error]` line. + const run = async () => { + const refusal = await refuseForUnreachableMedia(opts); + if (refusal) throw new Error(refusal); + await opts.fn(onLog, abort.signal, setProgress, ctx); + }; + run() .then(() => { if (record.status === "cancelled") { registry.finalize(id, "cancelled"); diff --git a/common/lib/channelMediaHold.ts b/common/lib/channelMediaHold.ts @@ -29,15 +29,30 @@ export function isMediaHeld(status: ChannelMediaStatus): boolean { // Why a channel is held, without the paths inspectChannelMedia's `detail` // carries. +// +// "Its media is moving" leads the in-transition reason because that is what the +// operator is looking at: a `.relocating.json` marker is a move, running or +// interrupted, and while it stands every writer of the channel is held — the +// four lanes skip it, a media job refuses to start on it, and both builds keep +// what they last read of it (release 16 slice RM). export const HELD_REASON: Record<ChannelMediaStatus, string> = { unreachable: "its media is not reachable (drive not mounted?)", - "in-transition": "a move of its media is in progress or was interrupted", + "in-transition": "its media is moving (a move is in progress or was interrupted)", inconsistent: "its data link and its config disagree", stalled: "its drive is not answering (a stalled disk)", ok: "reachable", "in-place": "reachable", }; +// THE HOLD AS A SURFACE SAYS IT — "held: its media is moving (…)" — or null for +// a channel whose media is read. One wording for the rack's chip, the channel +// page's Storage panel and a refused job, so the three cannot drift. It lifts +// when the status does: a move that completes or is abandoned removes the +// marker, and the next inspect reads the channel again. +export function mediaHoldText(status: ChannelMediaStatus): string | null { + return isMediaHeld(status) ? `held: ${HELD_REASON[status]}` : null; +} + // The reason, and the storage location's label when the channel's media is on // one. `dataDir` is the channel config's; the inspector's target stands in when // the config names none (a move in flight). diff --git a/common/views/channelRow.ts b/common/views/channelRow.ts @@ -1,6 +1,7 @@ import type { ChannelConfig } from "../lib/channelConfig"; import type { ChannelSnapshot } from "../controller/channelSnapshot"; import type { ChannelMediaLocation } from "../lib/channelMedia"; +import { mediaHoldText } from "../lib/channelMediaHold"; import type { PriorityOperation, StoredChannelTier, @@ -105,6 +106,12 @@ export type ChannelRowView = { // Null for an in-place channel — the overwhelming majority — so the badge // marks only the rows whose other numbers may not be trustworthy. media: ChannelRowMedia | null; + // THE CHANNEL'S MEDIA HOLD, in its one wording (lib/channelMediaHold.ts): + // "held: its media is moving (…)" while a relocation marker stands, and the + // drive's words while it is unreachable, stalled or inconsistent. Null when + // the lanes read the channel. The rack draws it beside the tier, because it + // is the reason the lanes are skipping the row (release 16 slice RM). + mediaHold: string | null; // WHICH VOLUME, as an id a filter can name: a location id, "internal" for the // corpus volume, or "" for a dataDir under a root nobody named. volumeId: string; @@ -211,6 +218,7 @@ export function buildChannelRowView(i: ChannelRowInput): ChannelRowView { : i.volume.label, } : null, + mediaHold: i.media ? mediaHoldText(i.media.status) : null, volumeId: i.volume.id, volumeLabel: i.volume.label, mediaBytes: snap?.totalMediaBytes ?? null,