Archilyzer · Source

archilyzer

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

commit 513f3466c9ca304217c6e01bcdfe79ec4e658864
parent 1058f75cbcc62ea6c3afce5278d2bf6d0eee1387
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Fri, 11 Sep 2026 11:13:07 -0400

common: the move itself, both directions, with the source held until verified

`relocateChannelMedia({ paths, slug, direction, root })` copies
`channels/<slug>/data` out to `<root>/<slug>/data`, verifies it, swaps in an
absolute symlink and records the target in `config.json`. Direction "back"
undoes it. Same shape as the saved-video backup job: real rsync from
`paths.rsyncBin` through execa with `cancelSignal`, output streamed to the
caller's onLog.

THE SOURCE IS NEVER TOUCHED UNTIL THE COPY IS VERIFIED. An abort, an unmounted
root, a full target, an rsync failure and a verify failure all leave `data/`
exactly where it was; the only destructive step is the reclaim, which runs after
the link is in place and the config is written. `rsync -a` preserves mtimes, and
the LMDB index stores ids and mtimes with no paths, so a relocated channel needs
no reindex — the test asserts the mtime rather than trusting the flag.

The verify is two checks, not one. `rsync --dry-run --itemize-changes` must
print nothing AND both trees must measure the same file count and byte sum:
rsync runs without `--delete`, so a stray file already at the target is invisible
to the dry run and is caught only by the measurement. That is the case the test
seeds.

`.relocating.json` in the channel dir carries the target, direction and phase.
It makes the channel "in-transition" to every guard, and it is what lets a
killed job resume from copy / swap / reclaim instead of restarting. The rerun is
the ONE caller allowed to look past its own marker — a marker naming a different
target is refused.

`config.dataDir` is written only on success, from the config re-read at swap
time, so it is a record of what is on disk and never an intention.

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

Diffstat:
Acommon/controller/relocateChannelMedia.test.ts | 387+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/controller/relocateChannelMedia.ts | 589+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
2 files changed, 976 insertions(+), 0 deletions(-)

diff --git a/common/controller/relocateChannelMedia.test.ts b/common/controller/relocateChannelMedia.test.ts @@ -0,0 +1,387 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { + lstat, + mkdir, + mkdtemp, + readdir, + readFile, + readlink, + rm, + stat, + utimes, + writeFile, +} from "node:fs/promises"; +import { tmpdir } from "node:os"; +import path from "node:path"; +import type { Paths } from "../lib/paths"; +import { + inspectChannelMedia, + readRelocationMarker, + relocatedDataDir, +} from "../lib/channelMedia"; +import { previewRelocation, relocateChannelMedia } from "./relocateChannelMedia"; +import { readChannelConfig } from "./channels"; + +// Run with: +// pnpm --filter yt-dlp-transcript-common exec tsx --test controller/relocateChannelMedia.test.ts +// +// These exercise the REAL rsync binary (paths.rsyncBin -> "rsync"), like +// controller/backupSavedVideos.test.ts. Everything happens inside one mkdtemp. + +const MTIME = new Date("2021-03-04T05:06:07.000Z"); + +async function withTmp( + fn: (paths: Paths, root: string, dir: string) => Promise<void>, +): Promise<void> { + const dir = await mkdtemp(path.join(tmpdir(), "ttb-relocate-")); + const paths = { + transcriptsDir: dir, + channelsDir: path.join(dir, "channels"), + rsyncBin: "rsync", + } as Paths; + const root = path.join(dir, "platter"); + await mkdir(root, { recursive: true }); + try { + await fn(paths, root, dir); + } finally { + await rm(dir, { recursive: true, force: true }); + } +} + +async function seed( + paths: Paths, + slug: string, + videos: Record<string, Record<string, string>>, + config: Record<string, unknown> = {}, +): Promise<string> { + const channelDir = path.join(paths.channelsDir, slug); + await mkdir(channelDir, { recursive: true }); + await writeFile( + path.join(channelDir, "config.json"), + JSON.stringify( + { handling: "transcribe", url: "https://example.com/c", ...config }, + null, + 2, + ) + "\n", + ); + // These live in the CHANNEL dir, not data/ — they must not move. + await writeFile(path.join(channelDir, "archive"), "youtube v1\n"); + await writeFile(path.join(channelDir, "playlist"), "https://x/v1\n"); + for (const [id, files] of Object.entries(videos)) { + const videoDir = path.join(channelDir, "data", id); + await mkdir(videoDir, { recursive: true }); + for (const [name, content] of Object.entries(files)) { + const file = path.join(videoDir, name); + await writeFile(file, content); + // A known mtime, so "rsync -a preserved it" is an assertion and not a + // coincidence. The LMDB index keys on mtimes, so this is what makes a + // relocation free of a reindex. + await utimes(file, MTIME, MTIME); + } + } + return channelDir; +} + +test("out: copies, links, records the target, keeps mtimes and reclaims the source", async () => { + await withTmp(async (paths, root) => { + const channelDir = await seed(paths, "alpha", { + v1: { "audio.m4a": "one".repeat(500), "transcript.json": "{}" }, + v2: { "audio.m4a": "two".repeat(500) }, + }); + const target = relocatedDataDir(root, "alpha"); + + const res = await relocateChannelMedia({ + paths, + slug: "alpha", + direction: "out", + root, + onLog: () => {}, + }); + assert.equal(res.direction, "out"); + assert.equal(res.target, target); + assert.equal(res.files, 3); + assert.equal(res.resumed, false); + + // data/ is an absolute symlink at the target... + const link = await lstat(path.join(channelDir, "data")); + assert.ok(link.isSymbolicLink()); + assert.equal(await readlink(path.join(channelDir, "data")), target); + + // ...so the on-disk contract channelDir/data/<id>/<file> still resolves. + assert.equal( + await readFile( + path.join(channelDir, "data", "v1", "transcript.json"), + "utf8", + ), + "{}", + ); + assert.equal( + (await stat(path.join(target, "v1", "audio.m4a"))).mtime.getTime(), + MTIME.getTime(), + ); + + // The record is in config.json, written only now, on success. + assert.equal((await readChannelConfig(paths, "alpha"))?.dataDir, target); + assert.equal((await inspectChannelMedia(paths, "alpha")).status, "ok"); + + // The parked copy is gone and the marker with it — this is the step that + // actually frees the source volume. + const left = await readdir(channelDir); + assert.deepEqual( + left.filter((n) => n.startsWith("data.")), + [], + ); + assert.equal(await readRelocationMarker(paths, "alpha"), null); + // Channel-level files never moved. + assert.ok(left.includes("archive")); + assert.ok(left.includes("playlist")); + }); +}); + +test("abort from an onLog hook leaves the source intact, and the rerun completes", async () => { + await withTmp(async (paths, root) => { + const channelDir = await seed(paths, "alpha", { + v1: { "audio.m4a": "x".repeat(20000) }, + }); + const target = relocatedDataDir(root, "alpha"); + + // Aborting from the log hook is how a Cancel button reaches this code: the + // job's onLog and its AbortSignal belong to the same run. Firing on the + // first line is deterministic — what it pins is the contract (the source is + // untouched, the marker is kept, a rerun resumes), not rsync's own + // --partial mechanics, which are rsync's to test. + const controller = new AbortController(); + await assert.rejects( + () => + relocateChannelMedia({ + paths, + slug: "alpha", + direction: "out", + root, + signal: controller.signal, + onLog: () => controller.abort(), + }), + /Cancelled/, + ); + + // The source is a REAL directory still, with its file. + assert.ok((await lstat(path.join(channelDir, "data"))).isDirectory()); + assert.equal( + (await readFile(path.join(channelDir, "data", "v1", "audio.m4a"), "utf8")) + .length, + 20000, + ); + // Nothing was recorded, and the marker says where it was going. + assert.equal((await readChannelConfig(paths, "alpha"))?.dataDir, undefined); + const marker = await readRelocationMarker(paths, "alpha"); + assert.equal(marker?.target, target); + assert.equal(marker?.direction, "out"); + assert.equal(marker?.phase, "copy"); + // And every guard now reads the channel as in transition. + assert.equal( + (await inspectChannelMedia(paths, "alpha")).status, + "in-transition", + ); + + // The rerun finishes it. It is the ONE caller allowed to look past its own + // marker. + const res = await relocateChannelMedia({ + paths, + slug: "alpha", + direction: "out", + root, + onLog: () => {}, + }); + assert.equal(res.resumed, true); + assert.equal((await inspectChannelMedia(paths, "alpha")).status, "ok"); + assert.equal((await readChannelConfig(paths, "alpha"))?.dataDir, target); + }); +}); + +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" }, + }); + 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"); + + await assert.rejects( + () => + relocateChannelMedia({ + paths, + slug: "alpha", + direction: "out", + root, + onLog: () => {}, + }), + /Verification failed/, + ); + + assert.ok((await lstat(path.join(channelDir, "data"))).isDirectory()); + assert.equal( + await readFile(path.join(channelDir, "data", "v1", "audio.m4a"), "utf8"), + "one", + ); + assert.equal((await readChannelConfig(paths, "alpha"))?.dataDir, undefined); + }); +}); + +test("back: restores a real directory, clears the config and reclaims the target", async () => { + await withTmp(async (paths, root) => { + const channelDir = await seed(paths, "alpha", { + v1: { "audio.m4a": "one", "transcript.json": "{}" }, + }); + const target = relocatedDataDir(root, "alpha"); + await relocateChannelMedia({ + paths, + slug: "alpha", + direction: "out", + root, + onLog: () => {}, + }); + + const res = await relocateChannelMedia({ + paths, + slug: "alpha", + direction: "back", + onLog: () => {}, + }); + assert.equal(res.direction, "back"); + assert.equal(res.files, 2); + + const back = await lstat(path.join(channelDir, "data")); + assert.ok(back.isDirectory()); + assert.ok(!back.isSymbolicLink()); + assert.equal( + await readFile(path.join(channelDir, "data", "v1", "audio.m4a"), "utf8"), + "one", + ); + assert.equal( + (await stat(path.join(channelDir, "data", "v1", "audio.m4a"))).mtime.getTime(), + MTIME.getTime(), + ); + assert.equal((await readChannelConfig(paths, "alpha"))?.dataDir, undefined); + assert.equal((await inspectChannelMedia(paths, "alpha")).status, "in-place"); + // The target is reclaimed, and so is <root>/<slug> once it is empty. + assert.equal(await pathThere(target), false); + assert.equal(await pathThere(path.dirname(target)), false); + assert.equal(await readRelocationMarker(paths, "alpha"), null); + assert.equal( + (await readdir(channelDir)).includes("data.incoming"), + false, + ); + }); +}); + +test("moving back a channel that was never moved is refused", async () => { + await withTmp(async (paths) => { + await seed(paths, "alpha", { v1: { "audio.m4a": "one" } }); + await assert.rejects( + () => + relocateChannelMedia({ + paths, + slug: "alpha", + direction: "back", + onLog: () => {}, + }), + /already in place/, + ); + }); +}); + +test("a relocated channel is refused a second move without a move back", async () => { + await withTmp(async (paths, root, dir) => { + await seed(paths, "alpha", { v1: { "audio.m4a": "one" } }); + await relocateChannelMedia({ + paths, + slug: "alpha", + direction: "out", + root, + onLog: () => {}, + }); + const other = path.join(dir, "platter2"); + await mkdir(other, { recursive: true }); + await assert.rejects( + () => + relocateChannelMedia({ + paths, + slug: "alpha", + direction: "out", + root: other, + onLog: () => {}, + }), + /already relocated/, + ); + }); +}); + +test("an unmounted root is refused before anything is written", async () => { + await withTmp(async (paths, _root, dir) => { + const channelDir = await seed(paths, "alpha", { v1: { "audio.m4a": "one" } }); + await assert.rejects( + () => + relocateChannelMedia({ + paths, + slug: "alpha", + direction: "out", + root: path.join(dir, "not-mounted"), + onLog: () => {}, + }), + /does not exist or is not a directory/, + ); + assert.equal(await readRelocationMarker(paths, "alpha"), null); + assert.ok((await lstat(path.join(channelDir, "data"))).isDirectory()); + }); +}); + +test("a social channel has no media to relocate", async () => { + await withTmp(async (paths, root) => { + await seed(paths, "poster", {}, { sourceKind: "social" }); + await assert.rejects( + () => + relocateChannelMedia({ + paths, + slug: "poster", + direction: "out", + root, + onLog: () => {}, + }), + /social channel/, + ); + }); +}); + +test("preview measures the tree and both volumes without moving anything", async () => { + await withTmp(async (paths, root) => { + const channelDir = await seed(paths, "alpha", { + v1: { "audio.m4a": "abcde", "transcript.json": "{}" }, + v2: { "audio.m4a": "fgh" }, + }); + const preview = await previewRelocation({ paths, slug: "alpha", root }); + assert.equal(preview.target, relocatedDataDir(root, "alpha")); + assert.equal(preview.files, 3); + assert.equal(preview.bytesToMove, 5 + 2 + 3); + assert.equal(preview.existingPartial, false); + assert.ok(preview.freeOnRoot > 0); + assert.ok(preview.freeOnSource > 0); + // Both dirs are under one tmpdir, so this really is the same volume. + assert.equal(preview.sameDevice, true); + assert.ok((await lstat(path.join(channelDir, "data"))).isDirectory()); + }); +}); + +async function pathThere(p: string): Promise<boolean> { + try { + await stat(p); + return true; + } catch { + return false; + } +} diff --git a/common/controller/relocateChannelMedia.ts b/common/controller/relocateChannelMedia.ts @@ -0,0 +1,589 @@ +import path from "node:path"; +import { + mkdir, + readdir, + readFile, + rename, + rm, + stat, + symlink, + unlink, + writeFile, +} from "node:fs/promises"; +import { execa } from "execa"; +import type { Paths } from "../lib/paths"; +import { isSocialChannel } from "../lib/channelConfig"; +import { getFreeBytes } from "../lib/diskSpace"; +import { formatBytes } from "../lib/format"; +import { + inspectChannelMedia, + relocatedDataDir, + relocationMarkerPath, + type RelocationDirection, + type RelocationMarker, + type RelocationPhase, +} from "../lib/channelMedia"; +import { readChannelConfig, writeChannelConfig } from "./channels"; + +// MOVE A CHANNEL'S MEDIA TO ANOTHER DRIVE, AND BACK. +// +// The mechanism is a symlink (see common/lib/channelMedia.ts for why): the media +// is copied to `<root>/<slug>/data`, verified, and `channels/<slug>/data` becomes +// an absolute link to it while `config.json` records the target in `dataDir`. +// Nothing that reads a channel changes, because the on-disk contract +// `channelDir/data/<id>/…` is preserved exactly. +// +// THE SOURCE IS NEVER TOUCHED UNTIL THE COPY IS VERIFIED. An abort, a full +// target, a crash, an rsync failure — all of them leave `data/` exactly where it +// was, and leave the partial copy resumable. The only destructive step is the +// reclaim at the very end, after the link is in place and the config is written. +// +// `rsync -a` preserves mtimes, which is what makes this free for the LMDB index: +// it stores ids and mtimes, no paths, so a relocated channel needs no reindex. + +export type RelocateChannelMediaResult = { + slug: string; + direction: RelocationDirection; + // Absolute path of the relocated data dir. For "back" this is what was + // reclaimed, not where the media now lives. + target: string; + bytes: number; + files: number; + // True when the run picked up an interrupted one from its marker rather than + // starting from scratch. + resumed: boolean; +}; + +type RelocateOpts = { + paths: Paths; + slug: string; + direction: RelocationDirection; + // Required for "out"; ignored for "back", which reads the target from config. + root?: string; + onLog?: (line: string) => void; + signal?: AbortSignal; +}; + +export type RelocationPreview = { + slug: string; + target: string; + bytesToMove: number; + files: number; + freeOnRoot: number; + freeOnSource: number; + // The move is pointless (and rename-unsafe assumptions change) when the root + // is the volume the corpus is already on. + sameDevice: boolean; + // A resumable partial copy from an earlier attempt is already at the target. + existingPartial: boolean; +}; + +// One walk, used by the preview, the space check and the verify. Follows no +// symlinks (a relocated channel is never the SOURCE of another relocation). +async function measureTree( + dir: string, +): Promise<{ bytes: number; files: number }> { + let bytes = 0; + let files = 0; + const stack = [dir]; + while (stack.length > 0) { + const cur = stack.pop() as string; + let entries; + try { + entries = await readdir(cur, { withFileTypes: true }); + } catch { + continue; + } + for (const e of entries) { + const p = path.join(cur, e.name); + if (e.isDirectory()) { + stack.push(p); + } else if (e.isFile()) { + try { + const st = await stat(p); + bytes += st.size; + files++; + } catch { + /* vanished mid-walk; the verify is what catches real drift */ + } + } + } + } + return { bytes, files }; +} + +async function pathExists(p: string): Promise<boolean> { + try { + await stat(p); + return true; + } catch { + return false; + } +} + +async function isDirectory(p: string): Promise<boolean> { + try { + return (await stat(p)).isDirectory(); + } catch { + return false; + } +} + +async function writeMarker( + paths: Paths, + slug: string, + marker: RelocationMarker, +): Promise<void> { + const file = relocationMarkerPath(paths, slug); + const tmp = `${file}.tmp-${process.pid}`; + await writeFile(tmp, JSON.stringify(marker, null, 2) + "\n"); + await rename(tmp, file); +} + +async function clearMarker(paths: Paths, slug: string): Promise<void> { + await rm(relocationMarkerPath(paths, slug), { force: true }); +} + +async function readMarkerRaw( + paths: Paths, + slug: string, +): Promise<RelocationMarker | null> { + try { + const raw = await readFile(relocationMarkerPath(paths, slug), "utf8"); + const r = JSON.parse(raw) as Partial<RelocationMarker>; + if (typeof r.target !== "string") return null; + return { + target: r.target, + direction: r.direction === "back" ? "back" : "out", + startedAt: typeof r.startedAt === "string" ? r.startedAt : "", + phase: + r.phase === "swap" || r.phase === "reclaim" + ? r.phase + : ("copy" as RelocationPhase), + }; + } catch { + return null; + } +} + +// rsync, exactly as backupSavedVideos does it: the real binary from +// paths.rsyncBin, cancellable, streaming its output into the job log. +async function rsyncTree(opts: { + paths: Paths; + src: string; + dest: string; + args: string[]; + log: (m: string) => void; + signal?: AbortSignal; +}): Promise<{ exitCode: number; output: string }> { + // Trailing slash on src: copy the CONTENTS, so <src>/ -> <dest>/ and not + // <dest>/data/. Getting this wrong is a silently nested corpus. + const args = [...opts.args, `${opts.src}/`, `${opts.dest}/`]; + opts.log(`$ ${opts.paths.rsyncBin} ${args.join(" ")}`); + const child = execa(opts.paths.rsyncBin, args, { + cancelSignal: opts.signal, + all: true, + buffer: false, + reject: false, + }); + let output = ""; + child.all?.on("data", (c: Buffer) => { + const text = c.toString("utf8"); + output += text; + opts.log(text); + }); + const result = await child; + return { exitCode: result.exitCode ?? 1, output }; +} + +// 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. +async function verifyCopy(opts: { + paths: Paths; + src: string; + dest: string; + log: (m: string) => void; + signal?: AbortSignal; +}): Promise<{ bytes: number; files: number }> { + const { exitCode, output } = await rsyncTree({ + ...opts, + args: ["-a", "--dry-run", "--itemize-changes"], + }); + if (exitCode !== 0) { + throw new Error(`Verification rsync failed (exit ${exitCode})`); + } + const drift = output + .split("\n") + .map((l) => l.trim()) + .filter((l) => l.length > 0 && !l.startsWith("sending incremental")) + .filter((l) => !/^(sent|total size|$)/.test(l)); + if (drift.length > 0) { + throw new Error( + `Verification failed: ${drift.length} file(s) still differ ` + + `(first: ${drift[0]}). The source has NOT been touched.`, + ); + } + const [a, b] = await Promise.all([measureTree(opts.src), measureTree(opts.dest)]); + if (a.files !== b.files || a.bytes !== b.bytes) { + throw new Error( + `Verification failed: source has ${a.files} file(s)/${formatBytes(a.bytes)}, ` + + `target has ${b.files}/${formatBytes(b.bytes)}. The source has NOT been touched.`, + ); + } + return a; +} + +// What the operator sees before committing to a move. Cheap enough to run on a +// form keystroke debounce: one tree walk of the channel plus two statfs calls. +export async function previewRelocation({ + paths, + slug, + root, +}: { + paths: Paths; + slug: string; + root: string; +}): Promise<RelocationPreview> { + const source = path.join(paths.channelsDir, slug, "data"); + const target = relocatedDataDir(root, slug); + const [measured, freeOnRoot, freeOnSource, existingPartial] = await Promise.all([ + measureTree(source), + getFreeBytes(root), + getFreeBytes(paths.channelsDir), + pathExists(target), + ]); + let sameDevice = false; + try { + const [a, b] = await Promise.all([stat(source), stat(root)]); + sameDevice = a.dev === b.dev; + } catch { + /* an unmounted or absent root is not "same device" */ + } + return { + slug, + target, + bytesToMove: measured.bytes, + files: measured.files, + freeOnRoot, + freeOnSource, + sameDevice, + existingPartial, + }; +} + +export async function relocateChannelMedia( + opts: RelocateOpts, +): Promise<RelocateChannelMediaResult> { + const { paths, slug, direction, signal } = opts; + const log = opts.onLog ?? ((m: string) => console.log(m)); + const config = await readChannelConfig(paths, slug); + if (!config) throw new Error(`Channel "${slug}" not found`); + if (isSocialChannel(config)) { + throw new Error( + `Channel "${slug}" is a social channel — it has no downloaded media to relocate`, + ); + } + + const channelDir = path.join(paths.channelsDir, slug); + const dataDir = path.join(channelDir, "data"); + const existingMarker = await readMarkerRaw(paths, slug); + + if (direction === "back") { + const target = config.dataDir?.trim(); + if (!target) { + throw new Error( + `Channel "${slug}" is not relocated — its media is already in place`, + ); + } + if (existingMarker && existingMarker.target !== target) { + throw new Error( + `A relocation to ${existingMarker.target} is already in progress for "${slug}"`, + ); + } + return moveBack({ + paths, + slug, + channelDir, + dataDir, + target, + log, + signal, + resumed: Boolean(existingMarker), + phase: existingMarker?.direction === "back" ? existingMarker.phase : "copy", + }); + } + + const root = opts.root?.trim(); + if (!root) throw new Error("No destination root given"); + if (!path.isAbsolute(root)) { + throw new Error(`The destination root must be an absolute path (got "${root}")`); + } + const target = relocatedDataDir(root, slug); + if (existingMarker && existingMarker.target !== target) { + throw new Error( + `A relocation to ${existingMarker.target} is already in progress for "${slug}" — ` + + `finish or clear it before moving to ${target}`, + ); + } + if (config.dataDir?.trim() && config.dataDir.trim() !== target) { + throw new Error( + `Channel "${slug}" is already relocated to ${config.dataDir.trim()}. ` + + `Move it back in place first.`, + ); + } + + return moveOut({ + paths, + slug, + config, + channelDir, + dataDir, + root, + target, + log, + signal, + resumed: Boolean(existingMarker), + phase: existingMarker?.direction === "out" ? existingMarker.phase : "copy", + }); +} + +async function moveOut(args: { + paths: Paths; + slug: string; + config: NonNullable<Awaited<ReturnType<typeof readChannelConfig>>>; + channelDir: string; + dataDir: string; + root: string; + target: string; + log: (m: string) => void; + signal?: AbortSignal; + resumed: boolean; + phase: RelocationPhase; +}): Promise<RelocateChannelMediaResult> { + const { paths, slug, channelDir, dataDir, root, target, log, signal } = args; + + // PREFLIGHT. Everything that can refuse does so here, before a single byte is + // written and before the marker exists. + if (!(await isDirectory(root))) { + throw new Error( + `The destination root ${root} does not exist or is not a directory ` + + `(is the drive mounted?)`, + ); + } + // A RESUMING RUN MUST TOLERATE ITS OWN MARKER. `inspectChannelMedia` reports + // "in-transition" for any channel carrying one, which is exactly right for + // every guard and exactly wrong here: the rerun that finishes an interrupted + // move is the one caller allowed to see it. A marker for a DIFFERENT target + // was already refused above. + const location = await inspectChannelMedia(paths, slug, args.config); + if (!args.resumed && location.status !== "in-place") { + throw new Error( + `Channel "${slug}" is not in a movable state: ${ + location.detail ?? location.status + }`, + ); + } + + const measured = await measureTree(dataDir); + log( + `Relocating ${slug}: ${measured.files} file(s), ${formatBytes(measured.bytes)} ` + + `-> ${target}`, + ); + + let phase = args.phase; + if (phase === "copy") { + const free = await getFreeBytes(root); + if (free < measured.bytes) { + throw new Error( + `Not enough space on ${root}: ${formatBytes(free)} free, ` + + `${formatBytes(measured.bytes)} to move`, + ); + } + await mkdir(target, { recursive: true }); + await writeMarker(paths, slug, { + target, + direction: "out", + 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({ + paths, + src: dataDir, + dest: target, + args: ["-a", "--partial", "--info=progress2"], + log, + signal, + }); + 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…"); + await verifyCopy({ paths, src: dataDir, dest: target, log, signal }); + phase = "swap"; + } + + if (phase === "swap") { + await writeMarker(paths, slug, { + target, + direction: "out", + startedAt: new Date().toISOString(), + phase: "swap", + }); + // Same device, so the rename is atomic: `data/` is a real dir one instant + // and the parked copy the next, never half of each. + const parked = path.join(channelDir, `data.relocated-${Date.now()}`); + if (await pathExists(dataDir)) { + await rename(dataDir, parked); + } + await symlink(target, dataDir); + // Written only now, on success: config.dataDir is a record of what is on + // disk, never an intention. + const fresh = (await readChannelConfig(paths, slug)) ?? args.config; + await writeChannelConfig(paths, slug, { ...fresh, dataDir: target }); + log(`Swapped: ${dataDir} -> ${target}`); + await writeMarker(paths, slug, { + target, + direction: "out", + startedAt: new Date().toISOString(), + phase: "reclaim", + }); + if (await pathExists(parked)) { + await rm(parked, { recursive: true, force: true }); + } + } else { + // Resuming at "reclaim": the swap already happened, so any parked copy left + // behind is the only thing outstanding. + for (const name of await readdir(channelDir).catch(() => [] as string[])) { + if (name.startsWith("data.relocated-")) { + await rm(path.join(channelDir, name), { recursive: true, force: true }); + } + } + } + + await clearMarker(paths, slug); + const freeNow = await getFreeBytes(paths.channelsDir); + log( + `Done. ${formatBytes(measured.bytes)} now on ${root}; ` + + `${formatBytes(freeNow)} free on the source volume.`, + ); + return { + slug, + direction: "out", + target, + bytes: measured.bytes, + files: measured.files, + resumed: args.resumed, + }; +} + +async function moveBack(args: { + paths: Paths; + slug: string; + channelDir: string; + dataDir: string; + target: string; + log: (m: string) => void; + signal?: AbortSignal; + resumed: boolean; + phase: RelocationPhase; +}): Promise<RelocateChannelMediaResult> { + const { paths, slug, channelDir, dataDir, target, log, signal } = args; + if (!(await isDirectory(target))) { + throw new Error( + `The relocated media at ${target} is not reachable (is the drive mounted?)`, + ); + } + const measured = await measureTree(target); + const incoming = path.join(channelDir, "data.incoming"); + log( + `Moving ${slug} back in place: ${measured.files} file(s), ` + + `${formatBytes(measured.bytes)} <- ${target}`, + ); + + let phase = args.phase; + if (phase === "copy") { + const free = await getFreeBytes(paths.channelsDir); + if (free < measured.bytes) { + throw new Error( + `Not enough space on the corpus volume: ${formatBytes(free)} free, ` + + `${formatBytes(measured.bytes)} to move back`, + ); + } + await mkdir(incoming, { recursive: true }); + await writeMarker(paths, slug, { + target, + direction: "back", + startedAt: new Date().toISOString(), + phase: "copy", + }); + const { exitCode } = await rsyncTree({ + paths, + src: target, + dest: incoming, + args: ["-a", "--partial", "--info=progress2"], + log, + signal, + }); + 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…"); + await verifyCopy({ paths, src: target, dest: incoming, log, signal }); + phase = "swap"; + } + + if (phase === "swap") { + await writeMarker(paths, slug, { + target, + direction: "back", + startedAt: new Date().toISOString(), + phase: "swap", + }); + // unlink, not rm -r: `data` is the LINK here, and removing it recursively + // would be the one way this whole design eats the media. + await unlink(dataDir).catch(() => {}); + await rename(incoming, dataDir); + const fresh = await readChannelConfig(paths, slug); + if (fresh) { + const { dataDir: _dropped, ...rest } = fresh; + await writeChannelConfig(paths, slug, rest); + } + log(`Swapped: ${dataDir} is a real directory again`); + await writeMarker(paths, slug, { + target, + direction: "back", + startedAt: new Date().toISOString(), + phase: "reclaim", + }); + } + + await rm(target, { recursive: true, force: true }); + // Leave <root>/<slug> behind only if something else is in it. + const slugRoot = path.dirname(target); + if ((await readdir(slugRoot).catch(() => ["keep"])).length === 0) { + await rm(slugRoot, { recursive: true, force: true }); + } + await clearMarker(paths, slug); + log(`Done. ${formatBytes(measured.bytes)} back in place.`); + return { + slug, + direction: "back", + target, + bytes: measured.bytes, + files: measured.files, + resumed: args.resumed, + }; +}