commit 8fe456e4351de308e7dfea3373e1b036d338215d
parent 0f8a0a1b38f9f64c2fc3277951efffc8005aaa7c
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Thu, 1 Oct 2026 12:25:58 -0400
common, editor: review fixes — the --delete guard checks the live media the caller names; every removal logged; a cancelled job counts until it has stopped; Reconcile past the copy phase says "resumed"
- M1: assertMirrorDirection takes `live` from the caller — the channel's
`channels/<slug>/data` (through its link) or the saved-video store — and
refuses a destination whose real path is, contains or sits inside it, plus
the source/destination containment check as before. copyMirrorVerify asks
before its first rsync, so a swapped call copies nothing either. Tests: src
and dest swapped across two roots, out and back (nothing spawned); a
destination inside / containing the live media; the containment check.
- L1: the mirror pass runs with --info=del, so every removal is in the job
log; the preview's partial-copy line says extras are removed from that copy.
- L2: registry.finalize stamps endedAt on a job cancelled while running (and
forceRelease on one cancelled and never finalized); channelWriters counts
"cancelled with no endedAt" as a writer, "stopping — wait for it to stop".
Unit tests (injected and through the live registry) and an e2e case through
the stuck-job route cancelled from /jobs.
- L3: RelocateChannelMediaResult.reconciled — true only when the copy phase
ran as a reconcile; the done line says "resumed" otherwise.
- N3: the changelog's example uses a made-up video id.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Diffstat:
11 files changed, 378 insertions(+), 65 deletions(-)
diff --git a/common/controller/channelWriters.test.ts b/common/controller/channelWriters.test.ts
@@ -1,6 +1,6 @@
import { test } from "node:test";
import assert from "node:assert/strict";
-import type { JobRecord } from "../jobs/registry";
+import { getRegistry, newJobId, type JobRecord } from "../jobs/registry";
import type { AutoQueueKind } from "../lib/autoQueueTypes";
import type { AutoRunnerInFlight } from "./autoRunner";
import {
@@ -170,3 +170,47 @@ test("more than one writer: the first named, the rest counted; none is no refusa
);
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
@@ -36,7 +36,10 @@ export type ChannelWriter =
kind: string;
// The kind's label (jobKindLabel), or the raw kind when it has none.
label: string;
- status: Extract<JobStatus, "running" | "queued">;
+ // "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;
@@ -93,16 +96,26 @@ export function channelWriters(
const queued: ChannelWriter[] = [];
for (const j of source.jobs()) {
if (j.channelSlug !== slug || ignore.has(j.kind)) continue;
- if (j.status !== "running" && !(opts.includeQueued && j.status === "queued")) {
+ // 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 === "running" ? running : queued).push({
+ (j.status === "queued" ? queued : running).push({
source: "job",
jobId: j.id,
kind: j.kind,
label: jobKindLabel(j.kind),
- status: j.status,
+ status: stopping ? "stopping" : (j.status as "running" | "queued"),
...(task ? { videoId: task.id, taskKind: task.kind } : {}),
...(!task && j.videoId ? { videoId: j.videoId } : {}),
});
@@ -164,11 +177,11 @@ export function describeChannelWriter(w: ChannelWriter): string {
}
// What the operator does about the first writer: a job is waited for or
-// cancelled on /jobs; a lane's unit is waited for, or its lane held.
+// 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 {
- return w.source === "lane"
- ? `wait for it, or hold the ${w.lane} lane`
- : "wait for it or cancel it";
+ 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
diff --git a/common/controller/relocateChannelMedia.test.ts b/common/controller/relocateChannelMedia.test.ts
@@ -1336,6 +1336,7 @@ test("reconcile: an extra and a changed file on the destination are settled, and
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\)/);
@@ -1352,6 +1353,33 @@ test("reconcile: an extra and a changed file on the destination are settled, and
});
});
+// 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" } });
diff --git a/common/controller/relocateChannelMedia.ts b/common/controller/relocateChannelMedia.ts
@@ -101,6 +101,11 @@ export type RelocateChannelMediaResult = {
// 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
@@ -741,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}`,
@@ -803,6 +809,8 @@ async function moveOut(args: {
rsyncBin: paths.rsyncBin,
src: dataDir,
dest: target,
+ // The channel's media, which --delete may never reach.
+ live: dataDir,
log,
progress: makeProgressSink({
totalBytes: measured.bytes,
@@ -819,6 +827,7 @@ async function moveOut(args: {
),
})
).retried;
+ reconciled = args.reconcile;
phase = "swap";
}
@@ -860,7 +869,7 @@ async function moveOut(args: {
log,
signal,
...(verifySrc === dataDir
- ? { mirror: { underConstruction: target } }
+ ? { mirror: { live: dataDir } }
: {}),
})
).retried;
@@ -938,6 +947,7 @@ async function moveOut(args: {
files: measured.files,
resumed: args.resumed,
retried: verifyRetried,
+ reconciled,
};
}
@@ -1000,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))
@@ -1062,11 +1073,14 @@ async function moveBack(args: {
});
// 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 copyMirrorVerify({
rsyncBin: paths.rsyncBin,
src: target,
dest: incoming,
+ live: dataDir,
log,
progress: makeProgressSink({
totalBytes: measured.bytes,
@@ -1085,6 +1099,7 @@ async function moveBack(args: {
),
})
).retried;
+ reconciled = args.reconcile;
phase = "swap";
}
@@ -1159,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
@@ -8,6 +8,7 @@ import {
readdir,
rm,
stat,
+ symlink,
utimes,
writeFile,
} from "node:fs/promises";
@@ -174,6 +175,7 @@ test("a file added to the source after the copy pass is carried by the mirror pa
rsyncBin: "rsync",
src,
dest,
+ live: src,
cancelled,
log: (m) => {
// The transcription that finished after rsync had passed v1/.
@@ -200,6 +202,7 @@ test("a file removed from the source after the copy pass is removed from the des
rsyncBin: "rsync",
src,
dest,
+ live: src,
cancelled,
log: (m) => {
// The transcriber cleaning up its scratch dir once it finished.
@@ -231,6 +234,7 @@ test("a stale extra dir on the destination (the realcandaceo case) is removed, t
rsyncBin: "rsync",
src,
dest,
+ live: src,
cancelled,
log: () => {},
});
@@ -249,6 +253,7 @@ test("one change between the mirror and the check gets one more pass", async ()
rsyncBin: "rsync",
src,
dest,
+ live: src,
cancelled,
log: (m) => {
lines.push(m);
@@ -275,6 +280,7 @@ test("a source that changes again after the second pass refuses, naming the path
rsyncBin: "rsync",
src,
dest,
+ live: src,
cancelled,
log: (m) => {
// Something that keeps writing: a new file before every check.
@@ -311,6 +317,7 @@ test("a reconcile lists the differences by kind first, then mirrors without a co
rsyncBin: "rsync",
src,
dest,
+ live: src,
cancelled,
reconcile: true,
log: (m) => lines.push(m),
@@ -329,51 +336,154 @@ test("a reconcile lists the differences by kind first, then mirrors without a co
});
});
-// THE ONE DIRECTION --delete MAY RUN IN, asserted before rsync is spawned. The
-// fake rsync here records that it ran; it must never be reached.
-test("--delete is refused toward anything but the copy under construction", async () => {
+// 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");
- await put(dest, "v1/other.mp3", "two");
- const ran = path.join(dir, "rsync-ran");
- const fake = path.join(dir, "fake-rsync.sh");
- await writeFile(fake, `#!/bin/sh\ntouch "${ran}"\nexit 0\n`);
- await chmod(fake, 0o755);
-
- // SWAPPED: the source named as the destination — the live media as the
- // target of --delete.
+ const { bin, ran } = await fakeRsync(dir);
+ // Inside: a copy built under the channel's own data/.
await assert.rejects(
() =>
mirrorTree({
- rsyncBin: fake,
+ rsyncBin: bin,
src: dest,
- dest: src,
- underConstruction: dest,
+ dest: path.join(src, "v1"),
+ live: src,
log: () => {},
}),
- MirrorDirectionError,
+ /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/,
);
- // One inside the other, either way round.
+ 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(
() =>
- assertMirrorDirection({
+ mirrorTree({
+ rsyncBin: bin,
src,
dest: path.join(src, "v1"),
- underConstruction: path.join(src, "v1"),
+ live: elsewhere,
+ log: () => {},
}),
/one is inside the other/,
);
await assert.rejects(
- () => assertMirrorDirection({ src: path.join(dest, "x"), dest, underConstruction: dest }),
+ () =>
+ 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");
- // Both trees as they were.
assert.deepEqual(await tree(src), ["v1/audio.mp3"]);
assert.deepEqual(await tree(dest), ["v1/other.mp3"]);
+ await assertMirrorDirection({ src, dest, live: src });
+ });
+});
- // The right way round is allowed, and does what it says.
- await assertMirrorDirection({ src, dest, underConstruction: dest });
+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)}`,
+ );
});
});
diff --git a/common/controller/relocateDir.ts b/common/controller/relocateDir.ts
@@ -337,11 +337,16 @@ export const COPY_ARGS = ["-a", "--partial", "--info=progress2"];
// 1755 files against 1750, with the source untouched — and nothing in the move
// could ever settle it. `--delete` is what can.
//
-// `--delete` IS ONLY EVER POINTED AT THE COPY UNDER CONSTRUCTION, never at the
-// source: run the other way round it would delete from the channel's live
-// media every file the partial copy has not reached yet. `mirrorTree` refuses
-// any other direction before rsync is spawned (assertMirrorDirection).
-export const MIRROR_ARGS = ["-a", "--delete"];
+// `--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.
@@ -452,43 +457,53 @@ export class MirrorDirectionError extends Error {
}
// THE ONE RULE `--delete` LIVES UNDER, asserted before rsync is spawned: the
-// destination is the copy under construction — the directory the move itself
-// created to copy into (`<root>/<slug>/data` out, `data.incoming` back,
-// `<root>/saved-videos` for the store) — and it is neither the source nor
-// inside it, nor the other way round. A caller that swaps `src` and `dest`
-// names the live media as the destination, which is not the copy it is
-// building, and is refused here with nothing deleted.
+// 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;
- underConstruction: string;
+ live: string;
}): Promise<void> {
- if (path.resolve(opts.dest) !== path.resolve(opts.underConstruction)) {
- throw new MirrorDirectionError(
- `Refusing to mirror into ${opts.dest}: --delete may only target the copy ` +
- `under construction (${opts.underConstruction}), never the source.`,
- );
- }
- const [src, dest] = await Promise.all([
+ 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.`,
+ `other, so --delete would reach the source. Nothing has been touched.`,
);
}
}
-// One `rsync -a --delete` pass, source → the copy under construction. With a
-// progress sink it reports like the copy pass (Reconcile and resume's pass is
-// the copy).
+// 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;
- underConstruction: string;
+ live: string;
log: (m: string) => void;
progress?: (line: string) => void;
signal?: AbortSignal;
@@ -552,6 +567,9 @@ export async function copyMirrorVerify(opts: {
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;
@@ -561,7 +579,8 @@ export async function copyMirrorVerify(opts: {
// The caller's sentence for a cancel, which names its own two directories.
cancelled: () => Error;
}): Promise<{ bytes: number; files: number; retried: boolean }> {
- const { rsyncBin, src, dest, log, signal } = opts;
+ 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(
@@ -589,7 +608,7 @@ export async function copyMirrorVerify(opts: {
rsyncBin,
src,
dest,
- underConstruction: dest,
+ live,
log,
progress: opts.reconcile ? opts.progress : undefined,
signal,
@@ -608,7 +627,7 @@ export async function copyMirrorVerify(opts: {
dest,
log,
signal,
- mirror: { underConstruction: dest },
+ mirror: { live },
});
}
@@ -631,7 +650,7 @@ export async function verifyCopy(opts: {
dest: string;
log: (m: string) => void;
signal?: AbortSignal;
- mirror?: { underConstruction: string };
+ mirror?: { live: string };
}): Promise<{ bytes: number; files: number; retried: boolean }> {
if (opts.mirror) return verifyMirrored({ ...opts, mirror: opts.mirror });
let retried = false;
@@ -693,7 +712,7 @@ async function verifyMirrored(opts: {
dest: string;
log: (m: string) => void;
signal?: AbortSignal;
- mirror: { underConstruction: string };
+ mirror: { live: string };
}): Promise<{ bytes: number; files: number; retried: boolean }> {
const { rsyncBin, src, dest, log, signal } = opts;
let retried = false;
@@ -728,7 +747,7 @@ async function verifyMirrored(opts: {
rsyncBin,
src,
dest,
- underConstruction: opts.mirror.underConstruction,
+ live: opts.mirror.live,
log,
signal,
});
diff --git a/common/controller/relocateSavedVideos.ts b/common/controller/relocateSavedVideos.ts
@@ -420,6 +420,7 @@ async function moveStoreOut(a: Inner): Promise<SavedVideosRelocateResult> {
rsyncBin: paths.rsyncBin,
src: store,
dest: target,
+ live: store,
log,
progress: makeProgressSink({
totalBytes: measured.bytes,
@@ -460,7 +461,7 @@ async function moveStoreOut(a: Inner): Promise<SavedVideosRelocateResult> {
signal: a.signal,
// The store is still the live directory here, so a difference gets
// the mirror pass a copy phase would give it.
- mirror: { underConstruction: target },
+ mirror: { live: store },
})
).retried;
await rename(store, parked);
@@ -630,6 +631,9 @@ async function moveStoreBack(a: Inner): Promise<SavedVideosRelocateResult> {
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,
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/editor/app/channels/[slug]/components/stages/StorageStage.tsx b/editor/app/channels/[slug]/components/stages/StorageStage.tsx
@@ -538,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 && (
diff --git a/editor/app/channels/lib/relocationJob.ts b/editor/app/channels/lib/relocationJob.ts
@@ -97,8 +97,10 @@ export async function enqueueRelocation(opts: {
onLog(
`${direction === "out" ? "Moved" : "Moved back"} ${result.files} file(s) / ` +
`${formatBytes(result.bytes)} — ${result.target}` +
+ // "Reconciled" only when the copy phase ran as a reconcile: a
+ // marker already past the copy phase had nothing to reconcile.
(result.resumed
- ? reconcile
+ ? result.reconciled
? " (reconciled and resumed an interrupted move)"
: " (resumed an interrupted move)"
: "") +
diff --git a/editor/e2e/channel-storage.spec.ts b/editor/e2e/channel-storage.spec.ts
@@ -741,6 +741,68 @@ test("a move that starts while a job writes into the channel refuses, naming the
).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