commit c9776bb5d55a59c8d8bedb5966aef61ea17007ea
parent 0512803e29a0d64c2b0c75f8b5aba6bb1d37a33f
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Thu, 1 Oct 2026 12:46:29 -0400
Merge r16/move-holds-writers (release 16 slice RM) — a media move refuses to start over a running job (naming it), holds the channel's writers while its marker stands, mirrors its copy before verifying (the --delete direction guarded against the live media, every removal logged), and the Storage panel offers Reconcile and resume; the writer that got past the guard on 2026-09-30 was a Transcribe all batch whose per-video guard ran only at enqueue — a queued media job now re-checks when it starts; a cancelled job counts as a writer until it has stopped; reviewed SHIP
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Diffstat:
30 files changed, 2923 insertions(+), 173 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,216 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import { getRegistry, newJobId, 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);
+});
+
+// A CANCEL IS A REQUEST, NOT AN EXIT (the review's L2). A running job's cancel
+// marks it `cancelled` at once while its function winds down; it is a writer,
+// "stopping", until the registry stamps `endedAt` when the function returns.
+test("a job cancelled but still winding down is a writer until it has stopped", () => {
+ const winding = job({ id: "S", channelSlug: "alpha", status: "cancelled" });
+ const writers = channelWriters("alpha", { source: source([winding]) });
+ assert.equal(writers.length, 1);
+ assert.equal(
+ channelWritersRefusal("alpha", writers),
+ 'Cannot move the media of "alpha" now: Transcribe all is stopping (job S) ' +
+ "— wait for it to stop. Nothing has been touched.",
+ );
+ const stopped = { ...winding, endedAt: 2 };
+ assert.deepEqual(channelWriters("alpha", { source: source([stopped]) }), []);
+});
+
+test("through the live registry: cancel leaves it stopping, the job's end stamps endedAt", () => {
+ const registry = getRegistry();
+ const id = newJobId();
+ const record: JobRecord = {
+ id,
+ kind: "whisper-all",
+ queueKey: `test:writers:${id}`,
+ channelSlug: `writers-${id}`,
+ status: "queued",
+ queuedAt: Date.now(),
+ logPath: "/dev/null",
+ };
+ registry.register(record);
+ registry.enqueue(record, { start: () => {}, onCancel: () => {} });
+ const slug = record.channelSlug as string;
+ assert.equal(channelWriters(slug).length, 1, "running");
+ assert.equal(registry.cancel(id), true);
+ assert.equal(record.status, "cancelled");
+ assert.equal(record.endedAt, undefined);
+ const stopping = channelWriters(slug);
+ assert.equal(stopping.length, 1);
+ assert.equal(stopping[0].source === "job" && stopping[0].status, "stopping");
+ // The job's function returns: streamCommand finalizes it as cancelled.
+ registry.finalize(id, "cancelled");
+ assert.equal(typeof record.endedAt, "number");
+ assert.deepEqual(channelWriters(slug), []);
+});
diff --git a/common/controller/channelWriters.ts b/common/controller/channelWriters.ts
@@ -0,0 +1,202 @@
+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;
+ // "stopping": cancelled, but its function has not returned yet — a
+ // transcriber finishing its current window still writes. The registry
+ // stamps `endedAt` only when the job has actually stopped.
+ status: Extract<JobStatus, "running" | "queued"> | "stopping";
+ // 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;
+ // A CANCEL IS A REQUEST, NOT AN EXIT. The registry marks a running job
+ // `cancelled` the moment it is asked to stop, and stamps `endedAt` when
+ // its function has returned (registry.ts `finalize`). Until then it is
+ // still a writer, "stopping" — an operator who cancels and moves at once
+ // would otherwise start a copy under it.
+ const stopping = j.status === "cancelled" && j.endedAt === undefined;
+ if (
+ j.status !== "running" &&
+ !stopping &&
+ !(opts.includeQueued && j.status === "queued")
+ ) {
+ continue;
+ }
+ const task = (j.tasks ?? []).find((t) => t.kind !== "relocate");
+ (j.status === "queued" ? queued : running).push({
+ source: "job",
+ jobId: j.id,
+ kind: j.kind,
+ label: jobKindLabel(j.kind),
+ status: stopping ? "stopping" : (j.status as "running" | "queued"),
+ ...(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 (one already stopping, only waited for); a lane's unit is
+// waited for, or its lane held.
+function waysOut(w: ChannelWriter): string {
+ if (w.source === "lane") return `wait for it, or hold the ${w.lane} lane`;
+ return w.status === "stopping" ? "wait for it to stop" : "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,45 @@ 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({
@@ -1044,17 +1124,278 @@ test("out @ swap: a file the target is missing is still a refusal", async () =>
root,
onLog: () => {},
}),
- /The source has NOT been touched/,
+ /Verification failed/,
);
+ assert.equal(
+ await readFile(path.join(target, "v1", "newer.json"), "utf8"),
+ "{}",
+ "nothing on the target was deleted to match the parked copy",
+ );
+ });
+});
- // Nothing committed: `data/` is still the real directory, holding both files.
+// ---------------------------------------------------------------------------
+// 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({
+ io: TEST_IO,
+ paths,
+ slug: "alpha",
+ direction: "out",
+ root,
+ onLog: () => {},
+ writers: () => (++asked === 1 ? [] : [TRANSCRIBING]),
+ }),
+ /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" });
+
+ 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);
+ assert.equal(res.reconciled, 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);
+ });
+});
+
+// A marker past the copy phase has nothing to reconcile: the run is a plain
+// resume, and says so (the review's L3).
+test("reconcile: a marker past the copy phase resumes, and does not claim a reconcile", async () => {
+ await withTmp(async (paths, root) => {
+ const channelDir = await seed(paths, "alpha", {
+ v1: { "audio.m4a": "one".repeat(100) },
+ });
+ const target = relocatedDataDir(root, "alpha");
+ await copyTree(path.join(channelDir, "data"), target);
+ await seedMarker(paths, "alpha", { target, direction: "out", phase: "swap" });
+ 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);
+ assert.equal(res.reconciled, false);
+ assert.equal(lines.some((l) => l.startsWith("Reconciling:")), false);
+ assert.equal((await inspectChannelMedia(paths, "alpha")).status, "ok");
+ });
+});
+
+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,11 +95,17 @@ 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;
+ // True only when a Reconcile and resume actually reconciled: it found the
+ // move in its copy phase and made the destination copy match the source. A
+ // reconcile of a marker already past the copy phase has nothing to
+ // reconcile and is a plain resume — the done line says "resumed".
+ reconciled: boolean;
};
// The progress shape is the shared mover's — re-exported because callers
@@ -129,8 +139,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 +474,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 +489,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 +543,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 +603,8 @@ export async function relocateChannelMedia(
signal,
resumed: Boolean(resume),
phase: resume?.phase ?? "copy",
+ writers,
+ reconcile,
});
}
@@ -595,9 +658,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 +695,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;
@@ -661,6 +746,7 @@ async function moveOut(args: {
// Sticky across both verify points below: a retry at either one is the fact
// the caller wants reported, and neither overwrites the other's answer.
let verifyRetried = false;
+ let reconciled = false;
log(
`Relocating ${slug}: ${measured.files} file(s), ${formatBytes(measured.bytes)} ` +
`-> ${target}`,
@@ -708,34 +794,40 @@ 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,
+ // The channel's media, which --delete may never reach.
+ live: dataDir,
+ 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;
+ reconciled = args.reconcile;
phase = "swap";
}
@@ -759,12 +851,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: { live: dataDir } }
+ : {}),
+ })
).retried;
}
@@ -840,6 +947,7 @@ async function moveOut(args: {
files: measured.files,
resumed: args.resumed,
retried: verifyRetried,
+ reconciled,
};
}
@@ -856,6 +964,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");
@@ -900,6 +1010,7 @@ async function moveBack(args: {
}
// Measured off whichever copy still exists, in the order they stop existing.
let verifyRetried = false;
+ let reconciled = false;
const measured = (await isDirectory(target))
? await measureTree(target)
: (await isDirectory(incoming))
@@ -954,33 +1065,41 @@ 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.
+ // `live` is `data/` — the link, which resolves to that target — so a call
+ // with source and destination swapped is refused before rsync runs.
verifyRetried ||= (
- await verifyCopy({ rsyncBin: paths.rsyncBin, src: target, dest: incoming, log, signal })
+ await copyMirrorVerify({
+ rsyncBin: paths.rsyncBin,
+ src: target,
+ dest: incoming,
+ live: dataDir,
+ 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;
+ reconciled = args.reconcile;
phase = "swap";
}
@@ -1055,5 +1174,6 @@ async function moveBack(args: {
files: measured.files,
resumed: args.resumed,
retried: verifyRetried,
+ reconciled,
};
}
diff --git a/common/controller/relocateDir.test.ts b/common/controller/relocateDir.test.ts
@@ -1,6 +1,31 @@
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,
+ symlink,
+ 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 +114,401 @@ 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,
+ live: src,
+ 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,
+ live: src,
+ 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,
+ live: src,
+ 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,
+ live: src,
+ 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,
+ live: src,
+ 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,
+ live: src,
+ 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 any rsync is spawned.
+// `live` comes from the caller (the channel's `data/`, the store), never from
+// `src` or `dest`, so the guard is not satisfied by construction. The fake
+// rsync records that it ran; it must never be reached.
+async function fakeRsync(dir: string): Promise<{ bin: string; ran: string }> {
+ const ran = path.join(dir, "rsync-ran");
+ const bin = path.join(dir, "fake-rsync.sh");
+ await writeFile(bin, `#!/bin/sh\ntouch "${ran}"\nexit 0\n`);
+ await chmod(bin, 0o755);
+ return { bin, ran };
+}
+
+test("source and destination swapped across two roots: refused, nothing spawned", async () => {
+ // Two separate roots, as a real move has: the corpus and the platter.
+ const corpus = await mkdtemp(path.join(tmpdir(), "ttb-live-"));
+ const platter = await mkdtemp(path.join(tmpdir(), "ttb-platter-"));
+ try {
+ const live = path.join(corpus, "channels", "alpha", "data");
+ const copy = path.join(platter, "alpha", "data");
+ await put(live, "v1/audio.mp3", "one");
+ await put(copy, "v1/other.mp3", "two");
+ const { bin, ran } = await fakeRsync(platter);
+
+ // OUT, swapped: the platter's copy named as the source, `data/` as the
+ // destination.
+ await assert.rejects(
+ () =>
+ copyMirrorVerify({
+ rsyncBin: bin,
+ src: copy,
+ dest: live,
+ live,
+ cancelled,
+ log: () => {},
+ }),
+ (err: unknown) =>
+ err instanceof MirrorDirectionError &&
+ /is the live media/.test(err.message),
+ );
+ // BACK, swapped: `data/` is a link to the platter copy by then, and the
+ // destination named is the copy it resolves to.
+ const linked = path.join(corpus, "channels", "beta", "data");
+ await mkdir(path.dirname(linked), { recursive: true });
+ await symlink(copy, linked);
+ const incoming = path.join(corpus, "channels", "beta", "data.incoming");
+ await mkdir(incoming, { recursive: true });
+ await assert.rejects(
+ () =>
+ copyMirrorVerify({
+ rsyncBin: bin,
+ src: incoming,
+ dest: copy,
+ live: linked,
+ cancelled,
+ log: () => {},
+ }),
+ MirrorDirectionError,
+ );
+ await assert.rejects(() => stat(ran), "the fake rsync was never spawned");
+ assert.deepEqual(await tree(live), ["v1/audio.mp3"]);
+ assert.deepEqual(await tree(copy), ["v1/other.mp3"]);
+
+ // The right way round, both directions, is allowed.
+ await assertMirrorDirection({ src: live, dest: copy, live });
+ await assertMirrorDirection({ src: copy, dest: incoming, live: linked });
+ } finally {
+ await rm(corpus, { recursive: true, force: true });
+ await rm(platter, { recursive: true, force: true });
+ }
+});
+
+test("a destination inside the live media, or containing it, is refused", async () => {
+ await withTrees(async ({ src, dest, dir }) => {
+ await put(src, "v1/audio.mp3", "one");
+ const { bin, ran } = await fakeRsync(dir);
+ // Inside: a copy built under the channel's own data/.
+ await assert.rejects(
+ () =>
+ mirrorTree({
+ rsyncBin: bin,
+ src: dest,
+ dest: path.join(src, "v1"),
+ live: src,
+ log: () => {},
+ }),
+ /is inside the live media/,
+ );
+ // Containing: a destination one level above the media.
+ await assert.rejects(
+ () => assertMirrorDirection({ src: dest, dest: dir, live: src }),
+ /contains the live media/,
+ );
+ await assert.rejects(() => stat(ran), "the fake rsync was never spawned");
+ assert.deepEqual(await tree(src), ["v1/audio.mp3"]);
+ });
+});
+
+test("--delete is refused when source and destination contain one another", async () => {
+ await withTrees(async ({ src, dest, dir }) => {
+ await put(src, "v1/audio.mp3", "one");
+ await put(dest, "v1/other.mp3", "two");
+ const { bin, ran } = await fakeRsync(dir);
+ const elsewhere = path.join(dir, "elsewhere");
+ await assert.rejects(
+ () =>
+ mirrorTree({
+ rsyncBin: bin,
+ src,
+ dest: path.join(src, "v1"),
+ live: elsewhere,
+ log: () => {},
+ }),
+ /one is inside the other/,
+ );
+ await assert.rejects(
+ () =>
+ assertMirrorDirection({
+ src: path.join(dest, "x"),
+ dest,
+ live: elsewhere,
+ }),
+ /one is inside the other/,
+ );
+ await assert.rejects(() => stat(ran), "the fake rsync was never spawned");
+ assert.deepEqual(await tree(src), ["v1/audio.mp3"]);
+ assert.deepEqual(await tree(dest), ["v1/other.mp3"]);
+ await assertMirrorDirection({ src, dest, live: src });
+ });
+});
+
+test("every removal by the mirror pass is in the log", async () => {
+ await withTrees(async ({ src, dest }) => {
+ await put(src, "v1/audio.mp3", "one");
+ await put(dest, "v1/audio.mp3", "one");
+ await put(dest, "v1/stale.json", "{}");
+ const lines: string[] = [];
+ await copyMirrorVerify({
+ rsyncBin: "rsync",
+ src,
+ dest,
+ live: src,
+ cancelled,
+ log: (m) => lines.push(m),
+ });
+ assert.ok(
+ lines.some((l) => /deleting v1\/stale\.json/.test(l)),
+ `the mirror's deletion is logged: ${JSON.stringify(lines)}`,
+ );
+ });
+});
+
+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,34 @@ 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 NEVER POINTED AT THE LIVE MEDIA: run the other way round it
+// would delete from the channel's media every file the partial copy has not
+// reached yet. `mirrorTree` and `copyMirrorVerify` refuse, before rsync is
+// spawned, a destination that is, contains or sits inside the live media the
+// CALLER names (assertMirrorDirection).
+//
+// `--info=del` puts every removal in the job log (`deleting <path>`): a file
+// taken off a copy is a fact the operator can read afterwards, not only a
+// count.
+export const MIRROR_ARGS = ["-a", "--delete", "--info=del"];
+
+// 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 +369,290 @@ 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 not the live media. `live` is supplied by the caller from
+// what it KNOWS is live — a channel's `channels/<slug>/data` (resolved through
+// its link, so on the way back it is the relocated target) or the saved-video
+// store — and never derived from `src` or `dest`, so a call with the two
+// swapped is caught: out, the destination would be `data/` itself; back, it
+// would be the target `data/` points at. Refused when the destination's real
+// path is the live media's, contains it, or sits inside it; and, belt and
+// braces, when the source and the destination contain one another.
+export async function assertMirrorDirection(opts: {
+ src: string;
+ dest: string;
+ live: string;
+}): Promise<void> {
+ const [src, dest, live] = await Promise.all([
+ realOrResolved(opts.src),
+ realOrResolved(opts.dest),
+ realOrResolved(opts.live),
+ ]);
+ if (isWithin(live, dest) || isWithin(dest, live)) {
+ throw new MirrorDirectionError(
+ `Refusing to mirror into ${opts.dest}: it ${
+ dest === live
+ ? "is"
+ : isWithin(live, dest)
+ ? "is inside"
+ : "contains"
+ } the live media at ${opts.live}${
+ live !== path.resolve(opts.live) ? ` (${live})` : ""
+ }, and --delete would remove from it every file the copy lacks. ` +
+ `Nothing has been touched.`,
+ );
+ }
+ 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. Nothing has been touched.`,
+ );
+ }
+}
+
+// One `rsync -a --delete` pass, source → the copy. 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;
+ live: 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;
+ // The live media — see assertMirrorDirection. Asked before the first rsync,
+ // so a swapped call copies nothing either.
+ live: 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, live, log, signal } = opts;
+ await assertMirrorDirection({ src, dest, live });
+ 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,
+ live,
+ 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: { live },
+ });
+}
+
// 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?: { live: 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 +686,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 +705,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: { live: 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,
+ live: opts.mirror.live,
+ 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,29 @@ 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,
+ live: store,
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 +459,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: { live: store },
})
).retried;
await rename(store, parked);
@@ -617,35 +624,29 @@ 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,
+ // The store's path, a link to the target by now: --delete may never
+ // reach what it resolves to.
+ live: store,
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/registry.ts b/common/jobs/registry.ts
@@ -268,6 +268,15 @@ class JobRegistry {
const job = this.jobs.get(id);
if (!job) return;
markTerminal(job, status, exitCode);
+ // A RUNNING job's cancel sets `cancelled` at once and leaves the record
+ // without `endedAt` while its function winds down (a transcriber finishing
+ // its window); markTerminal then skips it, because the status is already
+ // terminal. This is where the job has actually stopped, so it is stamped
+ // here — and "cancelled with no endedAt" means "still stopping" to anyone
+ // asking whether the job is still writing (controller/channelWriters.ts).
+ if (job.endedAt === undefined && job.status === "cancelled") {
+ job.endedAt = Date.now();
+ }
// Release per-task and drain references on terminal jobs.
job.tasks = [];
job.drainController = undefined;
@@ -337,6 +346,10 @@ class JobRegistry {
if (job.status === "queued" || job.status === "running") {
job.status = "cancelled";
job.endedAt = Date.now();
+ } else if (job.status === "cancelled" && job.endedAt === undefined) {
+ // Cancelled and never finalized (the wedged case this exists for): the
+ // operator has declared it stopped, so it stops counting as a writer.
+ job.endedAt = Date.now();
}
job.tasks = [];
job.drainController = undefined;
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,
diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md
@@ -14,6 +14,7 @@
- **A job cancelled before it started now stays cancelled.** Its record on disk kept saying "queued", so a restart could put a job you had just cancelled back in its queue, and a clip fetch cancelled while waiting could be reported as still queued. Jobs still waiting when the editor shuts down are handled as before: the next start settles or re-queues them.
- **A site's Charts tab gives a sixth series its own colour.** The dashboard's charts coloured their series from five colours and started again at the sixth, so a chart broken down by six or more channels drew the sixth in the first one's colour. The sixth now takes the palette's sixth colour, and each from the seventh on a hue of its own; the first five are unchanged. The published sites' charts get the same change with their next build.
- **A form whose save is refused keeps what you typed.** Every editor form put its plain fields back to the stored values when its save was refused — a site's ID rejected, a page size out of range, a slug already taken — so everything typed had to be typed again. A refused save now leaves every field as you left it, beside the reason: **Settings**; a site's form (new and existing); the hub's config on `/sites`; **Cut release**; a channel's form (new and **Configure**), **Rename** and **Delete**; a video's **Delete directory**; **Drive health timing** on `/storage`; the backup config on `/saved-videos`; the sync operation's controls; the **Digest**, **Diarization**, **Speaker attribution** and **Speaker work lane** settings; and the worker list on `/workers`. A save that succeeds behaves as before, with one difference you may notice: a drop-down, and a checkbox or choice that the page tracks as you change it (a cadence, a worker's **Enabled**, a social link's **Keep in header**, a site membership, a site's accent), now shows what was saved. A form's own drop-downs used to go back to what the page had loaded with until a reload, and a second save from the same page sent that old choice again; the others went back until the page next refreshed itself (every 5 seconds by default).
+- **A media move no longer starts over a job that is writing into the channel, holds the channel's writers while it runs, and makes its copy match the source before it verifies — so a transcription or a download during a move cannot fail it.** A move that has waited its turn behind other moves now checks again when it starts: if a job is running on the channel, or an auto-queue lane is working on one of its videos, it stops at once and says which ("a transcription of abc123 is running (Transcribe all, job …) — wait for it or cancel it"), with nothing copied — and a job you have just cancelled counts until it has actually stopped ("is stopping … — wait for it to stop"); **Preview** says the same, and the Storage panel's blocked message now names the job too. While a move's marker stands, the channel is held: every lane skips it, and every job that reads or writes its media (single-video transcriptions, downloads and transcodes and the availability checks now included) refuses to start, including one that was already queued when the move began. The rack shows a **media held** chip in the channel's Tier cell and the Storage panel says "Held: its media is moving"; both go when the move finishes or its marker is cleared. The copy is now followed by a pass that makes the destination copy match the source — files the source no longer has are removed from the copy, never from the source — so a file written or deleted during the copy (a transcriber's scratch folder, say) no longer fails the check, and **Resume move** finishes a move whose copy holds such leftovers. Every file removed from a copy is listed in the move's log, and **Preview** says so when a copy from an earlier attempt is already there. If the source keeps changing, the move stops and lists what differs: extra on the destination, missing there, or changed. A new **Reconcile and resume** button beside **Resume move** lists those differences, makes the copy match and finishes the move, so no file has to be deleted by hand. The saved-video store's move does the same matching and the same check before it starts. Needs a rebuild and restart of the editor.
## [0.11.0] - 2026-09-30
- **Transcripts that arrived after a video was first seen are counted.** The stats behind the homepage, the hub and every site's charts were cached per video and refreshed only when the video's metadata changed, so a transcript that came later — a Whisper run days after the download, or a video downloaded after the last index build — never reached them, and a video with YouTube captions alone had no transcription date. Counts and charts were low; the homepage could show a site with 0 transcripts, 0 channels and 0 hours while it served its videos. A stat is now also redone whenever the index re-reads the video, every transcript has a date, and a captioned video is dated by when its captions arrived rather than by a later Normalize run, so its place on "Transcribed over time" can move. **After updating, rebuild and restart the editor before anything else:** until then, **Build stats dataset** runs the old code and would undo the new stats, while a site, hub or homepage build already runs the new code — and the first stats build of any kind re-reads every video once (about 10–30 minutes on a large archive; it can be stopped and picks up where it stopped). Then build the index, the stats, the homepage, the hub, and the sites.
diff --git a/editor/app/api/test/stuck-job/route.ts b/editor/app/api/test/stuck-job/route.ts
@@ -27,6 +27,12 @@ export async function GET(request: Request) {
if (denied) return denied;
const url = new URL(request.url);
const queueKey = url.searchParams.get("queue") || "stuck-queue";
+ // `slug` (+ `task`): the fake job runs ON A CHANNEL, transcribing one video —
+ // a writer the media move must refuse over and name (release 16 slice RM's
+ // relocate e2e). Without them the record is the slug-less holder it always
+ // was.
+ const slug = url.searchParams.get("slug") || undefined;
+ const task = url.searchParams.get("task") || undefined;
const registry = getRegistry();
const paths = getPaths();
@@ -36,6 +42,7 @@ export async function GET(request: Request) {
id,
kind: "whisper-all",
queueKey,
+ ...(slug ? { channelSlug: slug } : {}),
status: "queued",
queuedAt: Date.now(),
logPath,
@@ -47,7 +54,9 @@ export async function GET(request: Request) {
// and clear tasks so it reads as an idle, long-held slot → "possibly-stalled".
registry.enqueue(record, { start: () => {}, onCancel: () => {} });
record.startedAt = Date.now() - 20 * 60 * 1000;
- record.tasks = [];
+ record.tasks = task
+ ? [{ id: task, label: task, kind: "transcribe", startedAt: Date.now() }]
+ : [];
// A log line so the page's lastLogLine tail has something to show.
try {
diff --git a/editor/app/channels/[slug]/components/stages/StorageStage.tsx b/editor/app/channels/[slug]/components/stages/StorageStage.tsx
@@ -13,6 +13,7 @@ import {
clearRelocationMarkerAction,
moveChannelMediaBackAction,
previewRelocationAction,
+ reconcileRelocationAction,
relocateChannelMediaAction,
resumeRelocationAction,
} from "../../storageActions";
@@ -112,6 +113,10 @@ type Props = {
destinations: StorageDestination[];
// `settings.storage.defaultLocationId` — which one the select opens on.
defaultLocationId: string;
+ // The channel's media hold in its one wording ("held: its media is moving
+ // (…)", lib/channelMediaHold.ts), or null. Computed on the server: the hold
+ // module is not imported into a client file.
+ mediaHold: string | null;
};
export function StorageStage({
@@ -126,6 +131,7 @@ export function StorageStage({
canClearMarker,
destinations,
defaultLocationId,
+ mediaHold,
}: Props) {
return (
<div className="flex flex-col gap-6">
@@ -161,6 +167,24 @@ export function StorageStage({
{location.detail && (
<p className="text-sm text-muted-foreground">{location.detail}</p>
)}
+ {/* THE HOLD, said where the operator looks for it (release 16 slice
+ RM): the rack carries the same sentence on its "media held" chip.
+ It lifts with the status — the move completing, its marker cleared,
+ the drive back. */}
+ {mediaHold && (
+ <p
+ role="status"
+ aria-label="media hold"
+ className="text-sm rounded border border-warning/40 bg-warning-soft px-3 py-2"
+ >
+ {mediaHold.charAt(0).toUpperCase() + mediaHold.slice(1)}. The lanes
+ skip this channel and its media jobs refuse to start until the hold
+ lifts
+ {location.status === "in-transition"
+ ? " — when the move completes, or its marker is cleared below."
+ : "."}
+ </p>
+ )}
{mediaBytes === null && (
<p className="text-xs text-muted-foreground">
No audio total in this channel’s report yet — refresh the
@@ -514,7 +538,9 @@ function MoveMedia({
{preview?.existingPartial && (
<p className="text-sm text-muted-foreground">
A partial copy from an earlier attempt is already at the target; this
- run resumes it rather than starting over.
+ run resumes it rather than starting over, and files there that the
+ channel no longer has are removed from that copy (never from the
+ channel).
</p>
)}
{!confirmed && named && !movingBack && (
@@ -605,8 +631,17 @@ function StaleMarker({
<p className="text-sm text-muted-foreground">
<strong>Resume move</strong> runs the same move again from where
it stopped — whatever already copied correctly is not copied
- twice, and the source is not touched until the copy verifies. That
- is the usual answer.
+ twice, the copy is mirrored to the source, and the source is not
+ touched until the copy verifies. That is the usual answer.
+ </p>
+ <p className="text-sm text-muted-foreground">
+ <strong>Reconcile and resume</strong> is the answer to a move
+ whose verification failed: it first lists how the copy on the
+ destination differs from the source — extra there, missing there,
+ changed — then makes the copy match the source (removing from the
+ copy what the source no longer has; the source itself is never
+ changed), verifies, and finishes the move. No file needs deleting
+ by hand.
</p>
<p className="text-sm text-muted-foreground">
<strong>Clear marker</strong> removes the marker file and nothing
@@ -657,6 +692,21 @@ function StaleMarker({
</button>
}
/>
+ {/* THE REMEDIATION, with its own log so each run reads as itself. Same
+ gate as Resume: a marker stands and nothing is running on the
+ channel. */}
+ <StreamActionLog
+ key="reconcile-move-log"
+ trigger={() => {
+ setRanHere(true);
+ return reconcileRelocationAction(slug);
+ }}
+ cancelAction={cancelJobAction}
+ buttonLabel="Reconcile and resume"
+ runningLabel="Reconciling…"
+ label="Reconcile and resume"
+ disabled={!canAct || busy}
+ />
</section>
);
}
diff --git a/editor/app/channels/[slug]/page.tsx b/editor/app/channels/[slug]/page.tsx
@@ -44,6 +44,7 @@ import {
channelMediaStall,
inspectChannelMedia,
} from "yt-dlp-transcript-common/lib/channelMedia";
+import { mediaHoldText } from "yt-dlp-transcript-common/lib/channelMediaHold";
import { MediaNotAnswering } from "./components/MediaNotAnswering";
import {
isDriveNotAnswering,
@@ -638,7 +639,7 @@ export default async function ChannelDetailPage({
const blockedReason =
busy ??
(marker
- ? `A relocation (${marker.direction}) to ${marker.target} is in flight, or was interrupted at phase "${marker.phase}". A channel in transition is not moved again from here — the running job finishes it, and an interrupted one is finished by "Resume move" below.`
+ ? `A relocation (${marker.direction}) to ${marker.target} is in flight, or was interrupted at phase "${marker.phase}". A channel in transition is not moved again from here — the running job finishes it, and an interrupted one is finished by "Resume move" or "Reconcile and resume" below.`
: null);
// THE DESTINATIONS, EACH WITH ITS DRIVE'S CURRENT STATE. One probe per
// configured location, memoised for 10 s inside the controller — so a
@@ -682,6 +683,9 @@ export default async function ChannelDetailPage({
};
})}
defaultLocationId={settings.storage.defaultLocationId}
+ // THE HOLD, in the rack's words: "held: its media is moving (…)"
+ // while the marker stands (release 16 slice RM).
+ mediaHold={mediaHoldText(media.status)}
/>
);
}
diff --git a/editor/app/channels/[slug]/storageActions.ts b/editor/app/channels/[slug]/storageActions.ts
@@ -122,6 +122,28 @@ export async function relocateChannelMediaAction(
export async function resumeRelocationAction(
slug: string,
): Promise<StreamActionResult> {
+ return resumeRelocation(slug, false);
+}
+
+// RECONCILE AND RESUME — the remediation for a move whose verification failed
+// (release 16 slice RM). The same job as Resume move, asked to make the
+// destination copy MATCH the source before it verifies: one mirror pass
+// (`rsync -a --delete` toward the destination copy, never the source), the
+// verify, the swap, the reclaim. What it found is logged by kind first — extra
+// on the destination, missing there, changed — so the operator reads what was
+// settled; a difference that will not settle (something still writing) is
+// refused again and named. The operator never deletes a file by hand.
+export async function reconcileRelocationAction(
+ slug: string,
+): Promise<StreamActionResult> {
+ return resumeRelocation(slug, true);
+}
+
+// Not exported: in a "use server" file every export is an endpoint.
+async function resumeRelocation(
+ slug: string,
+ reconcile: boolean,
+): Promise<StreamActionResult> {
const paths = getPaths();
const marker = await readRelocationMarker(paths, slug);
if (!marker) {
@@ -129,16 +151,19 @@ export async function resumeRelocationAction(
ok: false,
error:
`Channel "${slug}" has no relocation marker — there is no ` +
- `interrupted move to resume.`,
+ `interrupted move to ${reconcile ? "reconcile" : "resume"}.`,
};
}
// The same guard the two moves use, and for the same reason: a marker with a
// LIVE run behind it is not interrupted, it is in progress, and a second job
// would copy into the directory the first one is writing.
- const refusal = channelMediaBusyReason(slug, "resuming its move");
+ const refusal = channelMediaBusyReason(
+ slug,
+ reconcile ? "reconciling its move" : "resuming its move",
+ );
if (refusal) return { ok: false, error: refusal };
if (marker.direction === "back") {
- return enqueueRelocation({ slug, direction: "back" });
+ return enqueueRelocation({ slug, direction: "back", reconcile });
}
const root = path.dirname(path.dirname(marker.target));
if (marker.target !== relocatedDataDir(root, slug)) {
@@ -150,7 +175,7 @@ export async function resumeRelocationAction(
`it. Clear the marker and start the move again.`,
};
}
- return enqueueRelocation({ slug, direction: "out", root });
+ return enqueueRelocation({ slug, direction: "out", root, reconcile });
}
export async function moveChannelMediaBackAction(
diff --git a/editor/app/channels/components/ChannelTierSelect.tsx b/editor/app/channels/components/ChannelTierSelect.tsx
@@ -57,6 +57,12 @@ export type ChannelTierSelectProps = {
// says the tier ITSELF was set by something other than the operator, and
// will be set back.
autoPausedReason?: string | null;
+ // THE CHANNEL'S MEDIA HOLD — "held: its media is moving (…)" while a
+ // relocation marker stands — or null (common/views/channelRow.ts'
+ // `mediaHold`). Distinct from both chips above: the lanes skip this channel
+ // whatever its tier and whatever the focus, until the move completes or its
+ // marker is cleared (release 16 slice RM).
+ mediaHold?: string | null;
disabled?: boolean;
};
@@ -86,6 +92,7 @@ export default function ChannelTierSelect({
focused = false,
heldReason = null,
autoPausedReason = null,
+ mediaHold = null,
disabled = false,
}: ChannelTierSelectProps): React.ReactNode {
const [pending, startTransition] = useTransition();
@@ -187,6 +194,20 @@ export default function ChannelTierSelect({
<span className="sr-only"> — {autoPausedReason}</span>
</span>
)}
+ {/* THE MEDIA HOLD. The chip says the fact in the rack's width; the
+ sentence ("held: its media is moving …") is on `title` and in the
+ screen-reader text, as for the two chips above. */}
+ {mediaHold && (
+ <span
+ role="note"
+ aria-label={`media hold for ${slug}`}
+ title={mediaHold}
+ className="rounded-full border border-warning/30 bg-warning-soft px-1.5 text-[10px] uppercase tracking-wide text-warning"
+ >
+ media held
+ <span className="sr-only"> — {mediaHold}</span>
+ </span>
+ )}
{error && (
<span
role="alert"
diff --git a/editor/app/channels/components/channelColumns.tsx b/editor/app/channels/components/channelColumns.tsx
@@ -178,6 +178,7 @@ export const CHANNEL_COLUMNS: Record<
focused={c.priority.focused}
heldReason={c.priority.heldReason}
autoPausedReason={c.priority.autoPausedReason}
+ mediaHold={c.mediaHold}
/>
),
}),
diff --git a/editor/app/channels/lib/mediaBusy.ts b/editor/app/channels/lib/mediaBusy.ts
@@ -23,31 +23,37 @@
// transitively), so nothing reachable from a `"use client"` file may import it;
// `next build` is what proves that.
-import { getRegistry } from "yt-dlp-transcript-common/jobs/registry";
-import { getAutoRunnerStatus } from "yt-dlp-transcript-common/controller/autoRunner";
-import { LANES } from "yt-dlp-transcript-common/lib/autoQueueTypes";
+import {
+ channelWriters,
+ describeChannelWriter,
+} from "yt-dlp-transcript-common/controller/channelWriters";
// A sentence naming what is holding the channel, or null when nothing is. The
// caller supplies the verb (`"moving its media"`) so the same reason reads as
// an instruction in every panel it appears in.
+//
+// THE COUNTS, THEN THE FIRST WRITER BY NAME (release 16 slice RM): "Finish or
+// cancel 1 running/queued job(s) for this channel before moving its media.
+// Now: a transcription of v50t5yt is running (Transcribe all, job 01M…)." The
+// list is common/controller/channelWriters.ts's — the one the move's preview
+// and its job's first step ask — with queued jobs counted here, because this is
+// the courtesy before anything is enqueued. A lane's download unit is also a
+// registry job; it is counted once, as the unit.
export function channelMediaBusyReason(
slug: string,
what?: string,
): string | null {
- const jobs = getRegistry()
- .list()
- .filter(
- (j) =>
- j.channelSlug === slug &&
- (j.status === "running" || j.status === "queued"),
- ).length;
- const units = LANES.flatMap((kind) => getAutoRunnerStatus(kind).inFlight)
- .filter((u) => u.channelSlug === slug).length;
- if (jobs === 0 && units === 0) return null;
+ const writers = channelWriters(slug, { includeQueued: true });
+ if (writers.length === 0) return null;
+ const jobs = writers.filter((w) => w.source === "job").length;
+ const units = writers.length - jobs;
const parts: string[] = [];
if (jobs > 0) parts.push(`${jobs} running/queued job(s)`);
if (units > 0) parts.push(`${units} auto-queue unit(s) in flight`);
const subject = `${parts.join(" and ")} for this channel`;
- return what ? `Finish or cancel ${subject} before ${what}.` : subject;
+ const now = `Now: ${describeChannelWriter(writers[0])}.`;
+ return what
+ ? `Finish or cancel ${subject} before ${what}. ${now}`
+ : `${subject}. ${now}`;
}
diff --git a/editor/app/channels/lib/relocationJob.ts b/editor/app/channels/lib/relocationJob.ts
@@ -36,8 +36,11 @@ export async function enqueueRelocation(opts: {
slug: string;
direction: RelocationDirection;
root?: string;
+ // Reconcile and resume (release 16 slice RM): finish an interrupted move by
+ // mirroring the destination copy to the source, verifying, and completing.
+ reconcile?: boolean;
}): Promise<StreamActionResult> {
- const { slug, direction, root } = opts;
+ const { slug, direction, root, reconcile } = opts;
const paths = getPaths();
return runManagedFunction({
kind: "relocate-channel-media",
@@ -64,6 +67,7 @@ export async function enqueueRelocation(opts: {
slug,
direction,
root,
+ reconcile,
onLog,
onProgress: (p) => {
if (!started) {
@@ -93,12 +97,18 @@ export async function enqueueRelocation(opts: {
onLog(
`${direction === "out" ? "Moved" : "Moved back"} ${result.files} file(s) / ` +
`${formatBytes(result.bytes)} — ${result.target}` +
- (result.resumed ? " (resumed an interrupted move)" : "") +
+ // "Reconciled" only when the copy phase ran as a reconcile: a
+ // marker already past the copy phase had nothing to reconcile.
+ (result.resumed
+ ? result.reconciled
+ ? " (reconciled and resumed an interrupted move)"
+ : " (resumed an interrupted move)"
+ : "") +
// Said out loud because the operator has watched this refuse a move:
// a directory timestamp left by a sidecar written mid-copy used to
// fail the verify with 131 GB correctly on the far side.
(result.retried
- ? " (one directory timestamp settled by a second pass)"
+ ? " (a change made during the copy was settled by a second mirror pass)"
: ""),
);
// No snapshot regen — deliberately, and `relocate-channel-media` is in
diff --git a/editor/app/storage/lib/savedVideosJob.ts b/editor/app/storage/lib/savedVideosJob.ts
@@ -7,6 +7,7 @@ import {
} from "yt-dlp-transcript-common/jobs/streamCommand";
import { formatBytes } from "yt-dlp-transcript-common/lib/format";
import { relocateSavedVideos } from "yt-dlp-transcript-common/controller/relocateSavedVideos";
+import { savedVideosStoreBusyReason } from "./storeBusy";
// ONE ENQUEUE OF THE SAVED-VIDEO STORE MOVE, for the two callers that have one:
// the Move control on /storage and its Resume.
@@ -39,6 +40,12 @@ export async function enqueueSavedVideosRelocation(opts: {
const result = await relocateSavedVideos({
paths,
locationId: opts.locationId,
+ // Asked again as the job's first step: the action asked at enqueue,
+ // and the move may have waited behind others on the relocation queue.
+ busy: () =>
+ savedVideosStoreBusyReason("moving the saved-video store", {
+ runningOnly: true,
+ }),
onLog,
// Added on the FIRST frame, not at start: the preflight and a resumed
// run's verify pass transfer nothing, and a bar sitting at 0 % through
@@ -66,7 +73,7 @@ export async function enqueueSavedVideosRelocation(opts: {
`${formatBytes(result.bytes)} — ${result.target}` +
(result.resumed ? " (resumed an interrupted move)" : "") +
(result.retried
- ? " (one directory timestamp settled by a second pass)"
+ ? " (a change made during the copy was settled by a second mirror pass)"
: ""),
);
// No snapshot regen — the kind is in NO_REGEN_KINDS. The move changes
diff --git a/editor/app/storage/lib/storeBusy.ts b/editor/app/storage/lib/storeBusy.ts
@@ -54,13 +54,23 @@ const STORE_TOUCHING_KINDS = new Set([
// A sentence naming what is holding the store, or null. The caller supplies the
// verb, so the same reason reads as an instruction wherever it appears.
-export function savedVideosStoreBusyReason(what?: string): string | null {
+//
+// `runningOnly` is the store move's own first step (relocateSavedVideos'
+// `busy`): by then the move has waited its turn on the relocation queue, and a
+// download merely QUEUED behind other work is not writing — when it persists
+// later, the marker refuses the persist and leaves the container in its data
+// dir. The enqueue-time courtesy counts queued jobs too.
+export function savedVideosStoreBusyReason(
+ what?: string,
+ opts: { runningOnly?: boolean } = {},
+): string | null {
const jobs = getRegistry()
.list()
.filter(
(j) =>
STORE_TOUCHING_KINDS.has(j.kind) &&
- (j.status === "running" || j.status === "queued"),
+ (j.status === "running" ||
+ (!opts.runningOnly && j.status === "queued")),
);
// THE DOWNLOAD LANE'S UNITS MAKE NO JOB RECORD — the omnimirror lesson, and
// the reason `channelMediaBusyReason` exists in the shape it does. A unit
@@ -72,7 +82,10 @@ export function savedVideosStoreBusyReason(what?: string): string | null {
const parts: string[] = [];
if (jobs.length > 0) {
const kinds = [...new Set(jobs.map((j) => j.kind))].slice(0, 3).join(", ");
- parts.push(`${jobs.length} running/queued download job(s) (${kinds})`);
+ parts.push(
+ `${jobs.length} ${opts.runningOnly ? "running" : "running/queued"} ` +
+ `download job(s) (${kinds})`,
+ );
}
if (units > 0) parts.push(`${units} auto-download unit(s) in flight`);
const subject = `${parts.join(" and ")} — any of them can persist a source video into the store`;
diff --git a/editor/e2e/channel-storage.spec.ts b/editor/e2e/channel-storage.spec.ts
@@ -4,6 +4,7 @@ import {
mkdir,
readdir,
rm,
+ stat,
symlink,
utimes,
writeFile,
@@ -574,6 +575,15 @@ test("Resume move finishes an interrupted move and clears its marker", async ({
),
join(target, VIDEO, "transcript.en.vtt"),
);
+ // AND A STALE SCRATCH DIR the source no longer has — the 2026-10-01 case: a
+ // transcriber's `.audio.mp3.parakeet/` copied mid-transcription, then deleted
+ // from the source when the transcription finished. The resume used to copy
+ // everything else and refuse on the counts, every time; its mirror pass
+ // (release 16 slice RM) removes it from the copy.
+ const scratch = join(target, VIDEO, ".audio.mp3.parakeet");
+ await mkdir(scratch, { recursive: true });
+ await writeFile(join(scratch, "meta.json"), "{}");
+ await writeFile(join(scratch, "win-0000.json"), "[]");
await writeFile(
resolvePath(`test-transcripts/channels/${SLUG}/.relocating.json`),
JSON.stringify(
@@ -627,6 +637,293 @@ test("Resume move finishes an interrupted move and clears its marker", async ({
`${baseUrl}/api/channels/${SLUG}/videos/${VIDEO}/files/transcript.en.vtt`,
);
expect(file.status()).toBe(200);
+ // The stale scratch dir is gone from the copy, and the copy is the source's
+ // two files and nothing else.
+ expect(await existsAbs(scratch)).toBe(false);
+ expect((await readdir(join(target, VIDEO))).sort()).toEqual([
+ "metadata.info.json",
+ "transcript.en.vtt",
+ ]);
+});
+
+// ---------------------------------------------------------------------------
+// Release 16 slice RM — a move holds the channel's writers and mirrors its copy
+// ---------------------------------------------------------------------------
+
+const OPS_AUTH = { authorization: "Bearer test-worker-token" };
+
+async function existsAbs(p: string): Promise<boolean> {
+ return stat(p).then(
+ () => true,
+ () => false,
+ );
+}
+
+// The inspect memo (5 s) and the registry, dropped the way resetData drops
+// them: a marker planted by hand after a page has rendered is otherwise read
+// as "in place" for up to five seconds.
+async function forgetCaches(): Promise<void> {
+ await fetch(`${baseUrl}/api/test/invalidate-cache`).catch(() => {});
+}
+
+async function jobLog(
+ page: Page,
+ id: string,
+): Promise<{ status: string; content: string }> {
+ const res = await page.request.get(`${baseUrl}/api/jobs/${id}/log`);
+ return (await res.json()) as { status: string; content: string };
+}
+
+// THE 2026-09-30 SHAPE. A move asked for while the channel is quiet waits in
+// the relocation queue behind another; while it waits, a transcription starts
+// on the channel. The move used to start its copy over it. Now its first step
+// asks again and refuses, naming the job — and touches nothing.
+//
+// The holder is /api/test/stuck-job on the relocation queue, released after
+// 4 s; the writer is a fake running "Transcribe all" on this channel, on a
+// video (`?slug=&task=`). The move goes in through the ops API so the two
+// enqueues are milliseconds apart, well inside the holder's 4 s.
+test("a move that starts while a job writes into the channel refuses, naming the job", async ({
+ page,
+}, testInfo) => {
+ test.setTimeout(90_000);
+ await resetData("one-youtube-channel-with-data");
+ await generateReport(page, SLUG);
+ await quiet(page);
+ const root = testInfo.outputPath("busy-root");
+ await mkdir(root, { recursive: true });
+
+ const holder = await page.request.get(
+ `${baseUrl}/api/test/stuck-job?queue=relocate&releaseAfterMs=4000`,
+ );
+ expect(holder.ok()).toBe(true);
+ const queued = await page.request.post(`${baseUrl}/api/ops/relocate`, {
+ headers: OPS_AUTH,
+ data: { slugs: [SLUG], root },
+ });
+ expect(queued.ok()).toBe(true);
+ const { jobIds } = (await queued.json()) as { jobIds: string[] };
+ expect(jobIds).toHaveLength(1);
+ const moveId = jobIds[0];
+
+ const writer = await page.request.get(
+ `${baseUrl}/api/test/stuck-job?queue=${encodeURIComponent("rm-writer")}&slug=${SLUG}&task=${VIDEO}`,
+ );
+ expect(writer.ok()).toBe(true);
+ const writerId = ((await writer.json()) as { id: string }).id;
+
+ await expect
+ .poll(async () => (await jobLog(page, moveId)).status, { timeout: 30_000 })
+ .toBe("failed");
+ expect((await jobLog(page, moveId)).content).toContain(
+ `Cannot move the media of "${SLUG}" now: a transcription of ${VIDEO} is ` +
+ `running (Transcribe all, job ${writerId}) — wait for it or cancel it. ` +
+ `Nothing has been touched.`,
+ );
+
+ // NOTHING WAS TOUCHED: no marker, no copy, the media a real directory.
+ expect(
+ await pathExists(`test-transcripts/channels/${SLUG}/.relocating.json`),
+ ).toBe(false);
+ expect(await existsAbs(join(root, SLUG))).toBe(false);
+ expect((await lstat(dataDir())).isDirectory()).toBe(true);
+
+ // And the panel says it before anyone clicks: the move is blocked, and the
+ // sentence names the writer. `.first()`: the clip-window card on the same
+ // panel repeats the panel's blocked reason.
+ await page.goto(channelStage(SLUG, "storage"));
+ await expect(
+ page.getByText(
+ new RegExp(
+ `before moving its media\\. Now: a transcription of ${VIDEO} is running`,
+ ),
+ ).first(),
+ ).toBeVisible();
+});
+
+// A CANCEL IS A REQUEST, NOT AN EXIT (the review's L2). A job asked to stop
+// keeps writing until its function returns — a transcriber finishes its
+// window — and the registry stamps its end only then. Until that moment it is
+// a writer "stopping", and the move waits for it; after it, the move is free.
+//
+// The writer is the stuck-job route's fake Transcribe all on this channel's
+// video, cancelled from /jobs; its function "returns" when the route releases
+// it, 25 s after it was made.
+test("a job cancelled but still stopping holds the move until it has stopped", async ({
+ page,
+}, testInfo) => {
+ test.setTimeout(120_000);
+ await resetData("one-youtube-channel-with-data");
+ await generateReport(page, SLUG);
+ await quiet(page);
+ const root = testInfo.outputPath("stopping-root");
+ await mkdir(root, { recursive: true });
+
+ const writer = await page.request.get(
+ `${baseUrl}/api/test/stuck-job?queue=${encodeURIComponent("rm-stopping")}&slug=${SLUG}&task=${VIDEO}&releaseAfterMs=25000`,
+ );
+ expect(writer.ok()).toBe(true);
+ const writerId = ((await writer.json()) as { id: string }).id;
+
+ await page.goto("/jobs");
+ const row = page.locator(`tr[data-job-id="${writerId}"]`).first();
+ const cancel = row.getByRole("button", { name: /^Cancel$/ });
+ await expect(async () => {
+ if (await cancel.isVisible()) await cancel.click({ timeout: 2_000 });
+ await expect(row).toContainText("cancelled", { timeout: 2_000 });
+ }).toPass({ timeout: 20_000 });
+
+ // Cancelled, and still stopping: the panel names it, and Preview refuses.
+ await page.goto(channelStage(SLUG, "storage"));
+ await expect(
+ page
+ .getByText(
+ `Now: a transcription of ${VIDEO} is stopping (Transcribe all, job ${writerId}).`,
+ )
+ .first(),
+ ).toBeVisible();
+ await page.getByLabel("destination root").fill(root);
+ await page.getByRole("button", { name: "Preview", exact: true }).click();
+ await expect(
+ page.getByRole("alert").filter({ hasText: /is stopping/ }),
+ ).toContainText(
+ `a transcription of ${VIDEO} is stopping (Transcribe all, job ${writerId}) — wait for it to stop`,
+ { timeout: 15_000 },
+ );
+
+ // Once it has stopped, nothing holds the move.
+ await expect(async () => {
+ await page.goto(channelStage(SLUG, "storage"));
+ await expect(page.getByText(/is stopping/)).toHaveCount(0);
+ }).toPass({ timeout: 45_000 });
+ await page.getByLabel("destination root").fill(root);
+ await page.getByRole("button", { name: "Preview", exact: true }).click();
+ await expect(page.getByLabel("relocation preview")).toBeVisible({
+ timeout: 15_000,
+ });
+});
+
+// THE HOLD, WHERE THE OPERATOR LOOKS. A marker stands — a move running, or one
+// interrupted — so every lane skips the channel and its media jobs refuse; the
+// rack's tier cell and the channel's Storage panel say so in one sentence, and
+// both stop saying it when the marker goes (here: Clear marker, the abandon).
+test("the rack and the Storage panel show the hold while a marker stands, and it lifts with it", async ({
+ page,
+}, testInfo) => {
+ test.setTimeout(90_000);
+ await resetData("one-youtube-channel-with-data");
+ await generateReport(page, SLUG);
+ await quiet(page);
+ const target = join(testInfo.outputPath("hold-root"), SLUG, "data");
+ await writeFile(
+ resolvePath(`test-transcripts/channels/${SLUG}/.relocating.json`),
+ JSON.stringify({
+ target,
+ direction: "out",
+ startedAt: new Date().toISOString(),
+ phase: "copy",
+ }),
+ );
+ await forgetCaches();
+
+ await page.goto("/channels");
+ const chip = page.getByLabel(`media hold for ${SLUG}`);
+ await expect(chip).toContainText("media held");
+ await expect(chip).toContainText(
+ "held: its media is moving (a move is in progress or was interrupted)",
+ );
+
+ await page.goto(channelStage(SLUG, "storage"));
+ const hold = page.getByLabel("media hold", { exact: true });
+ await expect(hold).toContainText(
+ "Held: its media is moving (a move is in progress or was interrupted)",
+ );
+ await expect(hold).toContainText(
+ "when the move completes, or its marker is cleared below",
+ );
+
+ // Abandoned: the marker goes, and the hold with it. The click is retried
+ // until the marker is gone: one that lands before hydration does nothing.
+ const clear = page.getByRole("button", { name: "Clear marker" });
+ await expect(async () => {
+ if (await clear.isVisible()) await clear.click({ timeout: 2_000 });
+ expect(
+ await pathExists(`test-transcripts/channels/${SLUG}/.relocating.json`),
+ ).toBe(false);
+ }).toPass({ timeout: 20_000 });
+ await forgetCaches();
+ await page.goto(channelStage(SLUG, "storage"));
+ await expect(page.getByLabel("media hold", { exact: true })).toHaveCount(0);
+ await page.goto("/channels");
+ await expect(page.getByLabel(`media hold for ${SLUG}`)).toHaveCount(0);
+});
+
+// RECONCILE AND RESUME — the remediation for a move whose verification failed
+// (the ruling's last bullet). The destination copy holds an extra the source
+// does not (a transcriber's scratch file) and a file that differs from the
+// source's: the job lists them by kind, makes the copy match the source, and
+// finishes the move. The operator deletes nothing by hand.
+test("Reconcile and resume settles an extra and a changed file on the destination", async ({
+ page,
+}, testInfo) => {
+ test.setTimeout(90_000);
+ await resetData("one-youtube-channel-with-data");
+ const root = testInfo.outputPath("reconcile-root");
+ const target = join(root, SLUG, "data");
+ await mkdir(join(target, VIDEO), { recursive: true });
+ // The full copy, timestamps and all, as a copy pass leaves it…
+ for (const name of ["metadata.info.json", "transcript.en.vtt"]) {
+ const src = join(dataDir(), VIDEO, name);
+ await copyFile(src, join(target, VIDEO, name));
+ const { atime, mtime } = await stat(src);
+ await utimes(join(target, VIDEO, name), atime, mtime);
+ }
+ // …then an extra, and a changed file.
+ const scratch = join(target, VIDEO, ".audio.mp3.parakeet");
+ await mkdir(scratch, { recursive: true });
+ await writeFile(join(scratch, "meta.json"), "{}");
+ await writeFile(join(target, VIDEO, "transcript.en.vtt"), "WEBVTT\n\nstale\n");
+ await writeFile(
+ resolvePath(`test-transcripts/channels/${SLUG}/.relocating.json`),
+ JSON.stringify({
+ target,
+ direction: "out",
+ startedAt: new Date().toISOString(),
+ phase: "copy",
+ }),
+ );
+
+ await quiet(page);
+ await page.goto(channelStage(SLUG, "storage"));
+ const reconcile = page.getByRole("button", { name: "Reconcile and resume" });
+ await expect(reconcile).toBeEnabled();
+ await reconcile.click();
+
+ const out = page.getByLabel("Reconcile and resume output");
+ await expect(out).toContainText("reconciled and resumed an interrupted move", {
+ timeout: 60_000,
+ });
+ // WHAT IT FOUND, BY KIND, before it changed anything.
+ await expect(out).toContainText(
+ /Reconciling: the destination copy differs from the source — \d+ extra on the destination \([^)]*\.audio\.mp3\.parakeet/,
+ );
+ await expect(out).toContainText(/changed \([^)]*transcript\.en\.vtt/);
+
+ // Finished: link, config, no marker; the copy is the source's two files.
+ expect(
+ await pathExists(`test-transcripts/channels/${SLUG}/.relocating.json`),
+ ).toBe(false);
+ expect((await lstat(dataDir())).isSymbolicLink()).toBe(true);
+ expect(await existsAbs(scratch)).toBe(false);
+ expect((await readdir(join(target, VIDEO))).sort()).toEqual([
+ "metadata.info.json",
+ "transcript.en.vtt",
+ ]);
+ const file = await page.request.get(
+ `${baseUrl}/api/channels/${SLUG}/videos/${VIDEO}/files/transcript.en.vtt`,
+ );
+ expect(file.status()).toBe(200);
+ expect(await file.text()).not.toContain("stale");
});
// A MOVE MUST NOT MATERIALISE THE MOUNTPOINT.
diff --git a/plans/FACTS.md b/plans/FACTS.md
@@ -8090,3 +8090,106 @@ source mirror (homepage)". Anchors are at the branch.
(a read never throws). `forms-keep-input.spec.ts` `withSettingsUnwritable` restores the file in a
`finally`. A background writer in that window fails the same way, so the spec keeps it to one
submit.
+
+## A move holds the channel's writers and mirrors its copy (verified 2026-10-01, branch `r16/move-holds-writers`)
+
+The record is [`release-16.md`](release-16.md), "Slice RM, as shipped". Anchors are at the branch tip.
+This section amends "A channel's media may be on another drive" (above): **a move is now serialized
+against the channel's own writers by refusal at the job's start as well as at enqueue**, and the copy
+phase deletes from the destination.
+
+- **Which writer got past the guard on 2026-09-30 (`realcandaceo`, from the job metas, read only).**
+ Not a lane: the lanes skip a channel with a marker twice over (guard 2 in `buildChannelWork`, and the
+ pick→run backstop that reads the marker before a unit launches), and both held. It was a manual
+ **Transcribe all** (`whisper-all`, job `01M3TGYDXA7WCBGN014A90J3JW`, queue `transcription`, spec
+ `params: {}` — the scan over the channel's 434 items). The move (`01M3TGHZSMQHHYXH44ATXXFKSQ`) was
+ queued 21:16:42 and passed the action's busy check (nothing running); the Transcribe all was queued
+ and started at 21:23:30 — no marker existed, so its `needsMedia` guard passed, correctly; the move
+ started at 21:36:28 after waiting behind other moves on the one relocation queue, and asked nothing
+ about the channel's jobs. `v50t5yt`'s transcription (749.78 s) began ≈21:25:37, its parakeet scratch
+ dir (`.audio.mp3.parakeet/`, meta + four windows) was copied, and its outputs landed at
+ 21:38:06–07, after rsync had passed the directory; the verify at 22:02:37 refused (4 differ). The
+ batch then **kept writing while the marker stood** — `v511otv` at 22:40:29, `v519k6c` at 23:13:25,
+ `v51fpcd`'s scratch from 23:21 — until it was cancelled at 23:29:44: a job's guard ran only at
+ enqueue, and its per-video loop never re-asks. The resume (`01M3VX740RAK1H24RZ2CRKYZ79`, 10:17) copied
+ the new files and refused on the counts, 1755 against 1750 — the five scratch files on the
+ destination, which a copy that never deletes could not settle.
+- **The three closes.** (1) The move refuses over a writer at its preview, its job's FIRST STEP and a
+ third time right after the copy phase's marker is written (`relocateChannelMedia.ts:167`
+ `assertNoWriters`, `:493`, `:551`, `:672` `assertNoWritersUnderMarker` — the third removes a marker
+ this run created). (2) A `needsMedia` job's guard is asked again when the queue STARTS it
+ (`jobs/streamCommand.ts:429`, `refuseForUnreachableMedia` `:327`), so a job queued before a marker
+ and started after it fails before `fn` runs. (3) The per-video writers that were not in
+ `JOB_KINDS` at all — `whisper-video`, `transcribe-one`, `download-one-pipeline`, `transcode-audio`,
+ `check-availability`, `quick-availability-check`, `check-maybe-missing` (`jobs/jobKinds.ts:577`
+ on) — declare `needsMedia`; absent meant false, so none was ever refused. Labels stay absent
+ (`/jobs` unchanged). Not in the table and still not refused: `worker-transcribe`/`worker-unit` (this
+ box as a worker: a scratch dir or another machine's mount), `refresh-report` (asserts reachability
+ itself), and the corpus-wide `normalize-live-chat`/`archive-*` (no slug).
+- **Who is a writer: `controller/channelWriters.ts`.** `channelWriters(slug, opts)` (`:89`): registry
+ jobs whose `channelSlug` is the slug, running (queued too with `includeQueued`), minus
+ `ignoreKinds`, then every lane's in-flight units for the slug (`getAutoRunnerStatus(lane).inFlight`).
+ A lane download unit is also an `auto-download-unit` job; it is named once, as the lane's.
+ `describeChannelWriter` (`:164`) names the first task's video ("a transcription of v50t5yt is running
+ (Transcribe all, job …)"); `channelWritersRefusal` (`:190`) is the move's sentence. The move's job
+ passes `ignoreKinds: ["relocate-channel-media"]` (it is itself running on the slug; the relocation
+ queue runs one at a time); the preview passes nothing. **The editor's `channelMediaBusyReason`
+ (`editor/app/channels/lib/mediaBusy.ts`) reads the same list with `includeQueued`**, keeps its
+ counts sentence (the e2e regex `running\/queued job\(s\) for this channel` still matches) and
+ appends `Now: <the first writer>.`. Corpus-wide jobs carry no slug and are not writers here:
+ `normalizeAll` and `evictClipWindows` re-ask the hold per channel, and the mirror covers a write that
+ lands mid-copy.
+- **The hold's words: `lib/channelMediaHold.ts`.** `HELD_REASON["in-transition"]` (`:40`) leads with
+ "its media is moving (a move is in progress or was interrupted)" — it also reaches both builds' held
+ lines. `mediaHoldText(status)` (`:52`) is "held: <reason>" or null — the rack's chip
+ (`ChannelTierSelect`, `aria-label="media hold for <slug>"`, visible "media held", the sentence on
+ `title` and in sr-only text), the Storage panel's line (`aria-label="media hold"`), the refused
+ job's sentence, the runners' skip log. It reaches the row as `ChannelRowView.mediaHold`
+ (`views/channelRow.ts:114`) and the panel as a server-computed prop; `lib/storageLocations.ts`
+ (imported by the hold module) is never pulled into a client file by this. The runners skip on
+ `isMediaHeld` (`autoRunner.ts:641`); their skip log is now `[auto] skipping <slug>, held: <reason>:
+ media <status> — <detail>` (the old `media <status> — <detail>` is its tail).
+- **The mirror: `controller/relocateDir.ts`.** `MIRROR_ARGS = ["-a", "--delete", "--info=del"]` —
+ every removal is a `deleting <path>` line in the job log; the verify's dry run is `VERIFY_ARGS`
+ (`--delete` included, so an extra on the destination is a difference, `*deleting`). **`--delete`
+ never reaches the live media**: `assertMirrorDirection` (`:468`, asked by `copyMirrorVerify` before its
+ FIRST rsync, and by `mirrorTree`) takes `live` from the caller — `channels/<slug>/data` for a
+ channel (both directions; resolved through its link, so on the way back it is the relocated
+ target) and `paths.savedVideosDir` for the store — and refuses when the destination's realpath is
+ the live media's, contains it or is inside it; and when the source and destination contain one
+ another. `live` is never derived from `src`/`dest`, so a swapped call is refused with nothing
+ spawned (review M1: the first cut compared `dest` with an `underConstruction` every caller set to
+ `dest`, which could not fail). `copyMirrorVerify` (`:565`) is the copy phase for both
+ movers and both directions: the progress copy (`COPY_ARGS`), the mirror pass, `verifyCopy` in mirror
+ mode (`verifyMirrored`, `:709`): an empty itemized dry run plus equal counts; a difference gets one
+ more mirror pass (`retried`), a second refuses with `CopyVerificationError` (`:430`) carrying the
+ differences by kind (`classifyDrift`, `:391`: `*deleting` → extra on the destination, all-`+`
+ attributes → missing on the destination, anything else → changed). **A swap-phase re-verify
+ mirrors only from the live media** (`data/` still a real directory, `relocateChannelMedia.ts:872`);
+ against a parked `data.relocated-*` it is the strict legacy verify — the target is never mirrored
+ from a copy that is no longer live. Resume takes the copy phase again, so stale files on a
+ destination are removed.
+- **Reconcile and resume** = `relocateChannelMedia({ reconcile: true })`, refused without a marker
+ (`relocateChannelMedia.ts:556`): the copy phase skips the progress copy, logs `diffTrees` by kind
+ ("Reconciling: the destination copy differs from the source — …"), and the mirror pass carries the
+ progress sink. `RelocateChannelMediaResult.reconciled` is true only when the copy phase ran as a
+ reconcile; a reconcile of a marker past the copy phase is a plain resume and says "resumed". Editor: `reconcileRelocationAction` (`storageActions.ts:136`, Resume's guards), the
+ Storage panel's button beside Resume move (log `"Reconcile and resume output"`), the job's done line
+ "(reconciled and resumed an interrupted move)".
+- **The saved-video store shares the pipeline** (`relocateSavedVideos.ts` copy phases →
+ `copyMirrorVerify`; the swap re-verify mirrors while the store is still a real directory) and takes
+ `busy` (`:246`, asked first, `:259`); the editor's job passes the store check with `runningOnly`
+ (`editor/app/storage/lib/storeBusy.ts`). It has no reconcile button.
+- **A cancel is a request, not an exit.** `registry.cancel` on a RUNNING job sets `cancelled` at once
+ (no `endedAt`) and the function winds down; `markTerminal` then skips it, so `registry.finalize`
+ stamps `endedAt` itself for a cancelled record that has none, and `forceRelease` does for one never
+ finalized. `channelWriters` counts `cancelled` with no `endedAt` as a writer, status "stopping"
+ ("— wait for it to stop"). Before this, a job cancelled while running never got an `endedAt` at all
+ (the 2026-09-30 Transcribe all's meta has none).
+- **The race the third ask leaves is closed by ordering** (the review's trace): `registry.enqueue`
+ marks a job running before `start()`, whose guard then reads the marker; a lane unit is in
+ `inFlight` before the pick→run backstop reads the marker. Either the third ask sees the writer or
+ the writer sees the marker.
+- **Not covered:** a corpus-wide
+ writer mid-channel when a move starts; a resumed or reconciled move-out is charged the whole tree
+ against the destination's free space, not the remainder (move-back charges the remainder).
diff --git a/plans/release-16.md b/plans/release-16.md
@@ -24,7 +24,7 @@ slice's prompt carries its ruling, and this record carries what was built. Rules
| CK | `r16/search-in` | A "Search in" row — Transcripts, Posts, Live chat — on the export and hub search, Transcripts and Posts on by default | `common/components/{FiltersPanel,SearchSessionContext,SearchResults,SearchBar,exportFilterStorage}.tsx/.ts`, `common/lib/searchQuery.ts` and `common/lib/search/*` as its prompt names, `export/e2e/search-in.spec.ts` (new) and the specs its prompt names, `export/e2e/helpers.ts`; records: `plans/FACTS.md` |
| DX | `r16/research-setup` | The research-only setup (source → `pnpm install` → `claude mcp add archilyzer` → `/ask`) in one place, the homepage's AI and MCP doc; the sites' and the hub's Use-with-AI page removed and its links pointed at the doc; `README.md` §1/§4 and `mcp/README.md` their own copies (as amended) | `homepage/content/docs/ai-and-mcp.md`, `export/app/use-with-ai/` (removed), the Use with AI links (`export/app/components/{Header,MobileMenu,Footer}.tsx`, `export/app/(workspace)/ask/page.tsx`), `common/lib/{project,corpus}.ts` + `common/bin/compose-site.ts` (what named the page), `mcp/README.md`, `README.md` §1/§4 (wording only), `homepage/e2e/docs.spec.ts`, `export/e2e{,-hub}/use-with-ai-link.spec.ts` and the specs that visited the page |
| FK | `r16/forms-keep-input` | Every editor form keeps what was typed when its action fails: actions return the submitted values with the error, the shared field helpers seed from them | new `editor/app/lib/formState.ts` + test; `editor/app/components/forms/Field.tsx` and the local `Field`s in `SiteForm.tsx`, `ChannelForm.tsx`; every action that returns `{ok:false,error}`/`{error}` (sites, settings, operations/settingsActions, scheduler, storage, homepageActions, cutReleaseAction, channels, videoActions); the 15 forms the ruling lists; e2e `forms-keep-input.spec.ts` (new) + the existing `sites-crud`, `settings`, `channels` specs; records: `plans/FACTS.md`. As shipped, also `editor/app/components/forms/Controlled.tsx` (new) and one tag swap each in `DurationField`, `SocialLinksField`, `SiteMembershipsSection`, `WorkersField` ("Slice FK, as shipped") |
-| RM | `r16/move-holds-writers` | A media move holds the channel's writers and mirrors its copy, so a transcription or download during the move cannot fail it | `common/controller/relocateChannelMedia.ts` (+ the saved-video mover if it shares the code) + tests; the lane runners' per-channel skip (`common/controller/autoRunner.ts`, the backfill/digest/normalize/clip-fetch entry points, `common/lib/channelMediaHold.ts`); `common/jobs/jobKinds.ts` (`needsMedia`); the channel page's Storage panel and the rack's reason text; `editor/e2e/relocate*.spec.ts`; records: `plans/FACTS.md` |
+| RM | `r16/move-holds-writers` | A media move holds the channel's writers and mirrors its copy, so a transcription or download during the move cannot fail it | `common/controller/relocateChannelMedia.ts` (+ the saved-video mover if it shares the code) + tests; the lane runners' per-channel skip (`common/controller/autoRunner.ts`, the backfill/digest/normalize/clip-fetch entry points, `common/lib/channelMediaHold.ts`); `common/jobs/jobKinds.ts` (`needsMedia`); the channel page's Storage panel and the rack's reason text; `editor/e2e/relocate*.spec.ts`; records: `plans/FACTS.md`. As shipped, also `common/controller/{channelWriters,relocateDir}.ts` (+ tests), `common/jobs/streamCommand.ts`, `common/views/channelRow.ts`, `editor/app/channels/lib/{mediaBusy,relocationJob}.ts`, `editor/app/storage/lib/{storeBusy,savedVideosJob}.ts`, `editor/app/api/test/stuck-job/route.ts`; the relocate spec is `editor/e2e/channel-storage.spec.ts` ("Slice RM, as shipped") |
| XL | `r16/x-login` | Connecting an X account works: the Connect window is the operator's own browser without automation signals, and the fetchers can use the operator's browser login directly (`cookiesFromBrowser`) with no window at all | `common/social/{xSessionBroker,xGalleryDlFetcher,fetchers,playwrightRuntime}.ts` + tests; `common/lib/settingsSchema.ts` (one key) + `SETTINGS.md`; `editor/app/settings/{xSessionActions.ts,components/XSessionSection.tsx}`; `editor/e2e/settings*.spec.ts`; records: `plans/FACTS.md`, `ENVIRONMENT.md` if an env var is added |
## Slice CK — the ruling (2026-09-30)
@@ -785,6 +785,215 @@ FK row.
`new-channel-onboarding`, `social-channel`, `sites-crud`, `settings`: **70 passed**, 0 failed,
5.6 min. The full suite was not re-run, as the parent directed.
+### Slice RM, as shipped — a move holds the channel's writers and mirrors its copy (2026-10-01)
+
+Branch `r16/move-holds-writers` off `main` `2e2f6b2a` (the ruling's remediation bullet merged in from
+`9cec5be6`, a fast-forward), worktree `~/Projects/plans-export-header-first-search` (editor 4201,
+test 4211, export 4210 — `pnpm wt list`'s block #12), one Opus implementer. Scratch files `rm-*` in
+the job's `tmp`. The ruling is above ("Slice RM — the ruling").
+
+**Which writer got past the guard on 2026-09-30** (the job metas and logs in `transcripts/.jobs/`
+and `realcandaceo/data/v50t5yt/transcribe-outcome.json`, read only). Not a lane: the lanes skip a
+channel with a marker twice over (the projection's media check and the pick→run backstop), and both
+held. It was a manual **Transcribe all** (`whisper-all`, job `01M3TGYDXA7WCBGN014A90J3JW`, the scan
+over the channel's 434 items):
+
+| When | What |
+|---|---|
+| 21:16:42 | The move (`01M3TGHZSMQHHYXH44ATXXFKSQ`) is queued; the action's busy check finds nothing running |
+| 21:23:30 | Transcribe all is queued and starts. No marker exists yet, so its media guard passes — correctly |
+| ≈21:25:37 | Its transcription of `v50t5yt` begins (749.78 s); parakeet's scratch dir `.audio.mp3.parakeet/` fills window by window |
+| 21:36:28 | The move starts, after twenty minutes behind other moves on the one relocation queue. It asks nothing about the channel's jobs, writes the marker and copies — the scratch dir (meta + four windows) included |
+| 21:38:06–07 | `v50t5yt`'s transcript, cues and outcome land, after rsync had passed the directory; the scratch dir is deleted from the source |
+| 22:02:37 | The verify refuses: four paths differ. The marker stays (phase copy) |
+| 22:40, 23:13, 23:21 | The batch keeps writing with the marker standing — `v511otv`, `v519k6c`, `v51fpcd`'s scratch — until cancelled at 23:29:44 |
+| 2026-10-01 10:17 | The resume copies the new files and refuses on the counts, 1755 against 1750: the five scratch files on the destination, which a copy that never deletes could not settle |
+
+So three things let it through, and the slice closes each: the move asked about the channel's jobs
+only when it was enqueued; a job's media guard ran only at enqueue, never when its queue started it
+(and its per-video loop never re-asks); and seven per-video writers were not in `JOB_KINDS` at all,
+so they were never refused for a moving channel.
+
+**What it does.**
+- **A move refuses to start over a writer, and names it.** `common/controller/channelWriters.ts`
+ (new) reads who is writing into a channel: running registry jobs on its slug (queued ones too for the
+ editor's courtesy check) and every lane's in-flight units for it; a lane's download unit, which is
+ also a registry job, is named once. The refusal: `Cannot move the media of "realcandaceo" now: a
+ transcription of v50t5yt is running (Transcribe all, job 01M3…) — wait for it or cancel it. Nothing
+ has been touched.` (a lane unit: "wait for it, or hold the transcription lane"). Asked by
+ `previewRelocation`, by the job's first step, and a third time the moment the copy phase's marker is
+ written — after which the lanes skip the channel and a media job refuses to start, so that answer
+ cannot go stale; a refusal there removes a marker this run created. A job that has been cancelled
+ but whose function has not returned yet still counts ("is stopping — wait for it to stop"): the
+ registry stamps `endedAt` on a cancelled-while-running job only when it has actually stopped
+ (`registry.finalize`, and `forceRelease` for one never finalized). The editor's
+ `channelMediaBusyReason` (the panel's blocked message, rename, delete, the bulk move's skips) reads
+ the same list and appends `Now: <the first writer>.` to its counts.
+- **The writers are held while the marker stands.** A media job's guard (`needsMedia`) is asked again
+ when its queue starts it (`jobs/streamCommand.ts`), so a job queued before a move and started during
+ it fails before it runs, in the hold's words: `Channel "x" is held: its media is moving (a move is in
+ progress or was interrupted) — …`. `whisper-video`, `transcribe-one`, `download-one-pipeline`,
+ `transcode-audio`, `check-availability`, `quick-availability-check` and `check-maybe-missing` now
+ declare `needsMedia` (labels still absent, so `/jobs` shows them as before). The lanes skip on
+ `isMediaHeld`, logging the hold's sentence. `HELD_REASON["in-transition"]` leads with "its media is
+ moving", and `mediaHoldText` is the one wording: the rack's **media held** chip in the Tier cell
+ (`aria-label="media hold for <slug>"`, the sentence on `title`) and the Storage panel's "Held: its
+ media is moving (…). The lanes skip this channel and its media jobs refuse to start until the hold
+ lifts — when the move completes, or its marker is cleared below." Both go with the marker.
+- **The copy mirrors, toward the destination only.** `relocateDir.ts`'s `copyMirrorVerify` is the copy
+ phase for both movers and both directions: the progress copy, then `rsync -a --delete` source → the
+ copy under construction, then a verify whose dry run carries `--delete` plus equal counts. One
+ change seen between the mirror and the check gets one more mirror pass; a second refuses with the
+ paths by kind — extra on the destination, missing on the destination, changed — and "Something is
+ still writing into <src>: stop it, then Reconcile and resume." The mirror runs with `--info=del`,
+ so every file it removes is a `deleting <path>` line in the job log. **What the direction guard
+ proves** (`assertMirrorDirection`, asked by `copyMirrorVerify` before its first rsync and by
+ `mirrorTree`): the destination's real path is not the live media's, does not contain it and is not
+ inside it — `live` being what the caller knows is live, `channels/<slug>/data` resolved through its
+ link (so the relocated target on the way back) or the saved-video store, never derived from `src`
+ or `dest`; and the source and destination do not contain one another. A call with the two swapped
+ is refused with nothing spawned, in either direction. A resume takes the same path, so the
+ realcandaceo leftovers are mirrored away. A swap-phase re-verify mirrors only while `data/` is still the live directory; against
+ a parked `data.relocated-*` it stays the strict verify, so the target is never mirrored from a copy
+ that is no longer live.
+- **Reconcile and resume** (the remediation bullet): a button beside **Resume move**, with its own log
+ ("Reconcile and resume output"), offered under Resume's condition (a marker, nothing running). One
+ relocation job with `reconcile`: no copy pass; it lists how the destination differs from the source
+ by kind ("Reconciling: the destination copy differs from the source — 2 extra on the destination
+ (…), 0 missing on the destination, 2 changed (…)"), mirrors (the progress bar rides the mirror
+ pass), verifies, swaps and reclaims; the done line says "(reconciled and resumed an interrupted
+ move)" — only when the copy phase ran as a reconcile (`result.reconciled`): a marker already past the
+ copy phase has nothing to reconcile, and that run says "(resumed …)". Without a marker it refuses.
+- **The saved-video store's move shares the pipeline** — `relocateSavedVideos` uses
+ `copyMirrorVerify` both ways and mirrors its swap re-verify while the store is still a real
+ directory — and takes a `busy` first step, which the editor's job fills with the store check
+ (`savedVideosStoreBusyReason(…, { runningOnly: true })`: by then a merely queued download is not
+ writing, and its persist is refused by the marker anyway). It has no Reconcile button; its Resume
+ mirrors.
+
+**Commits**
+
+| Commit | What |
+|---|---|
+| `50fc580c` | `common:` `channelWriters.ts` + test; the mirror, the direction assertion, `copyMirrorVerify`, the verify by kind (`relocateDir.ts` + test); both movers (`reconcile`, `writers`, `busy`) + tests; the start-time media guard + test; seven `needsMedia` kinds; the hold's words, the lanes' skip + test; `ChannelRowView.mediaHold` |
+| `1c5b9739` | `editor:` the Storage panel's hold line and Reconcile and resume; the rack's chip; `mediaBusy` names the writer; the store job's first step |
+| `0ee291f9` | `editor(e2e):` `channel-storage.spec.ts` — three new cases and Resume extended (below); `/api/test/stuck-job` takes `slug` and `task` |
+| `05383fe6` | `plans:` this section, the slices table's RM row; FACTS; the editor changelog |
+| `f9dcf2d4` | `common, editor:` the review's fixes — M1 the guard checks the live media the caller names; L1 `--info=del` and the preview's partial-copy line; L2 a cancelled job counts until it has stopped (`registry.finalize`/`forceRelease` stamp `endedAt`); L3 `reconciled` on the result; tests and one e2e case |
+| `01de2613` | `editor:` the changelog bullet — the fixes, and a made-up video id in its example (N3) |
+| `3cc7a5da` | merge of `main` `a2229d68` (plans only: slice XL's ruling) — the slices table keeps both rows |
+| this commit | `plans:` the review, its fixes and the gates after them |
+
+#### Gates (logs `$T/rm-*.log`)
+
+- **tsc** (all workspaces) clean at `0ee291f9`, after deleting the worktree's stale
+ `export/.next/{dev/,}types` (still naming the removed `/use-with-ai` page, DX's note).
+- **common:** **2,432/2,432**, 78 s (`main`'s 2,404 + 28) — new: `channelWriters.test.ts` 7, `relocateDir.test.ts` +8 (the
+ mirror carries a file added after the copy pass and removes one deleted after it; the realcandaceo
+ stale dir; one change → one more pass; a second → refused by kind; reconcile lists then mirrors;
+ the `--delete` direction refused before rsync spawns; the itemize classifier),
+ `relocateChannelMedia.test.ts` 39 (two old cases re-premised, below; +9: a stray file mirrored
+ away, the strict parked re-verify, the refusal over a running job, the preview's, a writer seen
+ under the marker, a resume over a stale scratch dir, the same coming back, reconcile, reconcile
+ without a marker),
+ `relocateSavedVideos.test.ts` +2, `streamCommand.test.ts` +1 (a media job queued before the marker
+ refuses at its start), `autoRunner.test.ts` +1 (every lane projects nothing for a marked channel and
+ the hold lifts with the marker). **Editor unit:** 109/109. **test:scripts:** 302 passed, 2 skipped
+ (304).
+- **Build:** the capped editor build with the corpus linked (`ln -sT` the primary's `transcripts`,
+ `systemd-run --scope -p MemoryMax=5G`, `pnpm --filter editor exec next build`, the link removed
+ after): exit 0, 47 s, 1.68 GB peak, at `0ee291f9`.
+- **Numbers tool:** none.
+
+ | Run | At | Specs | Result |
+ |---|---|---|---|
+ | 1 | `0ee291f9` | `channel-storage`, `storage-locations`, `channels-storage-columns`, `channels`, `channel-rename`, `ops-api`, `channel-priority`, `channels-rack-layers`, `saved-videos`, `review` (which also matched `site-publish-preview`), `fetch-window` | **88 passed**, 1 failed, 6.9 min — `channel-rename` "rename and delete refuse while a job is writing": **Delete channel** was disabled when clicked (below) |
+ | 2 | `0ee291f9` | `channel-rename` × 4 | **8 passed**, 0 failed, 49 s |
+ | 3 | `0ee291f9` | the full editor suite | **683 passed**, 3 failed, 12 skipped (the rack-audit shots), 51.8 min — `jobs-channel` "auto-refreshes the jobs list": every button on the Playlist stage still disabled after 30 s (the page never hydrated); `perf-budget` "the dashboard renders within budget": `resetData` met `EEXIST` making `test-transcripts/channels` (the fixture race helpers.ts describes); `widget` "+N more expands": `ERR_CONNECTION_REFUSED` — the log's one `[WebServer] ⚠ Server is approaching the used memory threshold, restarting...` landed on it |
+ | 4 | `0ee291f9` | `jobs-channel`, `perf-budget`, `widget` | **34 passed**, 0 failed, 1.4 min (after ~4 min waiting for the queue). None of the three touches what the slice changed beyond the shared job start |
+
+ Three new `channel-storage` cases and one extended: **a move that starts while a job writes into the channel
+ refuses, naming the job** (the 2026-09-30 shape: a stuck job holds the relocation queue for 4 s, the
+ move goes in through `/api/ops/relocate`, a fake running Transcribe all on the channel's video
+ appears, the move starts and fails with the exact sentence; no marker, no copy; the panel's blocked
+ message names the writer); **the rack and the Storage panel show the hold while a marker stands, and
+ it lifts with it** (Clear marker); **Resume move** extended with a planted stale scratch dir on the
+ destination, gone after the resume; **Reconcile and resume settles an extra and a changed file on
+ the destination**.
+
+#### Found and left
+
+- **Run 1's `channel-rename` failure is that spec's race with the page's own refresh**, not this
+ slice: the spec renders the Danger zone while the channel is quiet, starts a sync out of band, and
+ submits Rename then Delete from the stale page. The sync's start moves the pulse token, and when the
+ next pulse tick (every 5 s) lands between the two submits, the refreshed page renders Delete
+ disabled with the busy reason (`error-context.md` shows both forms blocked, the reasons naming the
+ sync). 8/8 on the rerun. Not changed: the spec is not this slice's.
+- **The two existing tests whose premise the ruling reversed**, re-premised: "a verify failure keeps
+ the source…" planted a stray file on the target and expected a refusal — now the mirror removes it
+ (that case is "a stray file on the destination is removed by the mirror pass"), and the refusal case
+ is a source that keeps changing; "out @ swap: a file the target is missing is still a refusal" is now
+ "… is mirrored by one more pass".
+- **Corpus-wide writers carry no slug**, so the move's check cannot name them: a corpus-wide Normalize
+ or Evict already inside a channel when its move starts is not refused. Both re-ask the hold per
+ channel, and a write that lands mid-copy is carried by the mirror (or refused by kind if it keeps
+ landing).
+- **Kinds still outside `JOB_KINDS` with a slug:** `worker-transcribe` / `worker-unit` (this box as a
+ worker, writing into a scratch dir or another machine's mount), `refresh-report` (it asserts
+ reachability itself), and the slug-less `normalize-live-chat` / `archive-*`.
+- ~~A lane tick that inspected the channel just before the marker landed can dispatch a unit after
+ the third ask.~~ **Closed, by ordering** (the review's trace): a unit is in `inFlight` before the
+ pick→run backstop reads the marker, so a unit the third ask missed reads the marker itself; and a
+ job is marked running before its start-time guard reads the marker.
+- **A resumed or reconciled move-out is charged the whole tree against the destination's free space**,
+ not what is still missing (the move back charges the remainder). Unchanged; a near-full destination
+ can refuse a resume it has room for.
+- **A server action that writes into a video's directory without a job** (the video page's edits) is
+ not held by the marker; the mirror carries its write.
+
+#### Decisions the operator could overturn
+
+| What I did | The alternative |
+|---|---|
+| The move's own check refuses over RUNNING jobs and lane units; a QUEUED media job is refused when it starts instead | Refuse the move over queued jobs too (the editor's enqueue-time check does): the operator re-queues the move after a twenty-minute wait for a job that would not have written yet |
+| A third ask right after the marker is written, removing a marker the run created | The first step only; a job that starts while a big tree is measured would be copied over |
+| The swap-phase re-verify mirrors while `data/` is live, never from a parked copy | Mirror from the parked copy too: it would delete from the target what was written through the link after the swap |
+| The Storage panel's busy message keeps its counts and adds `Now: <writer>.` | Replace the counts with the named writer (the `ops-api` spec pins the counts' wording) |
+| The rack's hold is a chip in the Tier cell, shown for every held status (unreachable, stalled and inconsistent too: the lanes skip all of them) | Only for a moving channel, as the ruling's sentence names |
+| Reconcile and resume skips the progress copy and lets the mirror pass carry the bar | Run the copy pass first, as Resume does — the same bytes either way |
+
+#### Review
+
+**Verdict: SHIP AFTER FIXES** (`rm-review.md` in the job's scratch): one Medium, three Lows, three
+notes. Every `--delete` was traced and found pointed the right way; the fixes are what the guard and
+the wording promise.
+
+| Finding | Where |
+|---|---|
+| M1: half the `--delete` direction check could not fail — every caller set `underConstruction` to `dest`, so "dest equals the copy under construction" held by construction, and a swap across two volumes would have passed the nesting check | `f9dcf2d4`: `assertMirrorDirection` takes `live` from the caller (`channels/<slug>/data` through its link, or the store) and refuses a destination whose realpath is, contains or sits inside it; `copyMirrorVerify` asks before its first rsync. Tests: swapped across two roots, out and back, nothing spawned; a destination inside and one containing the live media; the containment check kept |
+| L1: the mirror's removals were not logged, and the preview did not say a pre-existing copy loses files | `f9dcf2d4`: `--info=del` on `MIRROR_ARGS` (a test reads the `deleting` line); the partial-copy line adds "files there that the channel no longer has are removed from that copy (never from the channel)" |
+| L2: a job cancelled but still winding down was not a writer | `f9dcf2d4`: `registry.finalize` stamps `endedAt` on a cancelled-while-running job when it actually stops (`forceRelease` on one never finalized); `channelWriters` counts "cancelled, no `endedAt`" as "stopping — wait for it to stop". Unit (injected and through the live registry) and e2e ("a job cancelled but still stopping holds the move until it has stopped", the stuck-job route cancelled from `/jobs`) |
+| L3: Reconcile on a marker past the copy phase said "reconciled" | `f9dcf2d4`: `RelocateChannelMediaResult.reconciled`; the done line says "resumed" otherwise; a unit test on a swap-phase marker |
+| N1: "four new" cases — three are new, Resume was extended | This commit |
+| N2: any running job on the slug refuses a move, `refresh-report` included | No action (errs safe) |
+| N3: the changelog's example quoted a real video id | `01de2613`: a made-up id |
+
+The review also traced the race "Found and left" listed (a lane tick dispatching after the third
+ask) and found it closed by ordering; that item is struck above.
+
+#### Gates after the review and the merge (at `3cc7a5da`; logs `$T/rm-*2.log`, `rm-tsc3.log`, `rm-e2e5.log`)
+
+- **tsc** (all workspaces) clean, 52 s. **common:** **2,438/2,438**, 94 s (+6: three direction
+ tests in place of one, the mirror's logged removal, two stopping-writer tests, the reconcile past
+ the copy phase). **Editor unit:** 109/109. **test:scripts:** 302 passed, 2 skipped.
+- **Build:** the capped editor build with the corpus linked: exit 0, 63 s, 1.67 GB, the link removed.
+- **e2e** (after ~1.5 min waiting for the queue):
+
+ | Run | At | Specs | Result |
+ |---|---|---|---|
+ | 5 | `3cc7a5da` | `channel-storage` (13, the stopping case included), `storage-locations`, `channels-storage-columns` | **24 passed**, 1 failed, 5.4 min — `storage-locations` "a volume that came up somewhere else is re-pointed": the first step's `clickUntil` on **Refresh** never saw the volume's identity written to settings in 30 s (the page showed the probe's identity; the action's write did not land). Nothing in that step is the slice's |
+ | 6 | `3cc7a5da` | `storage-locations` × 2 | **18 passed**, 0 failed, 2.0 min |
+
## Rollout
Both slices are export- and homepage-side; the editor and umtool are not rebuilt for this release.