#!/usr/bin/env tsx // THE ONE-OFF MEDIA-TIER MIGRATION (release 17, slice T3). // // archilyzer storage migrate-tier … [--dry-run] [--reclaim] // archilyzer storage migrate-tier --all [--order smallest] [--include-large] [--dry-run] [--reclaim] // // Before release 17 a relocation moved a channel's WHOLE `data/` to another // drive: `channels//data` an absolute symlink to `//data`, // `config.dataDir` recording it. Such a channel is `legacy` // (lib/channelMedia.ts): its text is on the far drive, so it is held by every // guard until this runs. This brings the TEXT home and leaves the big files // where they are, in the release-17 layout (plans/release-17.md, "The model"): // // channels//data/ a REAL directory on the corpus disk again // /audio.mp3 -> ../../media//audio.mp3 (relative, the file's mtime) // channels//media -> //media (ONE absolute link) // //media//… (the old `//data`, renamed) // config.json: mediaDir = //media, dataDir removed // // THE PHASES, per channel (plans/release-17.md, "The one-off migration"): // // preflight the channel is on the retired layout and its drive answers; // `assertRelocationRootPresent(root)`; the old tree is walked and // classified BY NAME (lib/mediaTier.ts); the space rule // `copy bytes + resume margin <= free(channelsDir) - floor`. // copy marker `{target: //media, direction: "out", // phase: "copy", scope: "tier-migration"}`; every entry that is // NOT tierable (the text, `clips/`, scratch, `source-media.*`) — // except a postprocessor's dead `*.temp.*`, left on the platter — // is copied to `channels//data.incoming/` by rsync from a // NUL-separated list the classifier wrote (never a glob); the // copy is verified by an itemized dry run AND per-kind counts and // bytes; then every tierable file gets its relative link, carrying // the file's times (`lutimes` — without them every migrated live // chat's cues would read as stale). // swap platter `rename(//data -> //media)`; // `symlink(//media, channels//media)`; // `unlink(channels//data)`; `rename(data.incoming -> data)`; // `patchChannelConfig(slug, { mediaDir }, { unset: ["dataDir"] })`. // Each step looks at the disk first, so a rerun after a crash // between any two of them finishes the rest. The reclaim note // (`.tier-migration.reclaim`) is left beside the config. // reclaim only with --reclaim, and only on a channel carrying the note: // the platter's copies of the text (an entry of // `//media//` that is not tierable, and whose twin // on the corpus disk `lstat`s the same) and the dead temps are // deleted, empty video dirs dropped; anything else is kept and // reported. // done the marker is cleared; bytes by kind and links made are printed. // // THE EDITOR MUST BE STOPPED. It holds its job registry in memory, so nothing // on disk can say a job of its is about to write into the channel; a real run // refuses while the editor's port answers or `.jobs/` shows a job in a live // process (`editorRunningReason`). A DRY RUN WRITES NOTHING — no marker, no // list, no directory — and only warns. // // IDEMPOTENT AND RESUMABLE. A channel already on the new layout is a no-op; a // `tier-migration` marker is resumed from its phase. Any OTHER marker (a media // move's) is refused: that is the Storage panel's to finish. // // No index rebuild follows: rsync -a keeps every text file's mtime, the links // carry the media files' mtimes, and the index stats only sidecars. import os from "node:os"; import path from "node:path"; import { lstat, lutimes, mkdir, readdir, readlink, rename, rm, rmdir, stat, symlink, unlink, writeFile, } from "node:fs/promises"; import type { Stats } from "node:fs"; import type { Paths } from "../lib/paths"; import type { SiteSettings } from "../lib/settings"; import { formatBytes } from "../lib/format"; import { getFreeBytes } from "../lib/diskSpace"; import { classifyEntry, CLIPS_DIR, isTierable } from "../lib/mediaTier"; import { channelMediaLink, relocatedMediaDir, tierLinkTarget, } from "../lib/mediaTier-server"; import { forgetChannelMedia, inspectChannelMedia, readRelocationMarker, relocatedDataDir, relocationMarkerPath, type RelocationMarker, } from "../lib/channelMedia"; import { clearDirMarker, linkOrDirState, makeProgressSink, rsyncTree, writeDirMarker, } from "../controller/relocateDir"; import { assertRelocationRootPresent, rootOfRelocatedMediaDir, } from "../controller/relocateChannelMedia"; import { patchChannelConfig, readChannelConfig } from "../controller/channels"; import { readJobMeta } from "../jobs/jobMeta"; import { processIsAlive } from "../jobs/bootQueuedJobs"; import { parseArgv } from "./_parseFlags"; import { runIfEntryPoint } from "./_cli"; const GB = 1024 ** 3; // THE THREE BIG-TEXT CHANNELS (plans/release-17.md, Rulings: "a stop before // the three big-text channels … with the free space reported; the operator // decides"). `--all` migrates everything else, smallest text first, and stops // before these; `--include-large` (or naming one) goes on. export const LARGE_TEXT_CHANNELS: readonly string[] = [ "omnibased", "rekietalaw", "the-quartering-rumble", ]; // The rebuilt text tree, beside `data` in the channel dir until the swap. export const INCOMING_NAME = "data.incoming"; // The NUL-separated rsync list, beside the marker; removed when the channel is done. export const LIST_NAME = ".tier-migration.files"; // THE RECLAIM NOTE, beside the config from the swap until a reclaim: "this // tool migrated this channel, and the drive still holds its text copies". // `--all --reclaim` takes only channels carrying it, and a named `--reclaim` // refuses one without it — a channel the Storage panel moved has nothing of // the retired layout to take, and walking its whole tree to find that out // costs minutes on a platter. export const RECLAIM_NOTE = ".tier-migration.reclaim"; // --------------------------------------------------------------------------- // The inventory: the old tree, by kind // --------------------------------------------------------------------------- // What the copy carries, by kind. `source` is a media file that is NOT // tierable — `source-media.*` stays a real file in `data//` (the // saved-video store renames it), so it is copied with the text rather than // left on the platter with no link (and then deleted by --reclaim). export type CopyKind = "text" | "clips" | "scratch" | "source"; export const COPY_KINDS: readonly CopyKind[] = ["text", "clips", "scratch", "source"]; export type Tally = { files: number; bytes: number }; export type KindTallies = Record; function emptyTallies(): KindTallies { return { text: { files: 0, bytes: 0 }, clips: { files: 0, bytes: 0 }, scratch: { files: 0, bytes: 0 }, source: { files: 0, bytes: 0 }, }; } export type TierableFile = { id: string; name: string; bytes: number; atime: Date; mtime: Date; }; export type Inventory = { // Every top-level directory of the old tree (a video dir), in order. videoIds: string[]; // Paths relative to the old tree, for rsync's --files-from. copyList: string[]; copy: KindTallies; copyBytes: number; tierable: TierableFile[]; tierableBytes: number; // Dead postprocessor temps (`*.temp.*`) left on the platter, neither copied // nor linked; `--reclaim` deletes them. left: Tally; }; // Which copy kind a video-dir entry is. `stats` decides only one thing: a // tierable NAME that is not a regular file (a directory, a link) is not // something the link step can stand in for, so it is copied as it is. function copyKindOf(name: string): CopyKind { if (name === CLIPS_DIR) return "clips"; const kind = classifyEntry(name); if (kind === "scratch") return "scratch"; if (kind === "media") return "source"; return "text"; } // Files and bytes under one entry: a file or link counts itself (`lstat` // size — rsync -a copies a link as a link), a directory everything below it. async function measureEntry(p: string, st?: Stats): Promise { const s = st ?? (await lstat(p)); if (!s.isDirectory()) return { files: 1, bytes: s.size }; const out = { files: 0, bytes: 0 }; for (const name of await readdir(p)) { const t = await measureEntry(path.join(p, name)); out.files += t.files; out.bytes += t.bytes; } return out; } function add(t: KindTallies, kind: CopyKind, by: Tally): void { t[kind].files += by.files; t[kind].bytes += by.bytes; } // A POSTPROCESSOR'S DEAD TEMP — `source-media.temp.mp4`, `audio.temp.mp3`: // yt-dlp's postprocessor writes it and renames it over the final name, so one // still there after the run is a leftover (15 GB of them on one channel, // measured 2026-10-02). It is scratch by the classifier; the migration leaves // it on the platter, unlinked, rather than carry it to the corpus disk, and // `--reclaim` deletes it with the text copies. export function isLeftOnPlatter(name: string): boolean { return classifyEntry(name) === "scratch" && /\.temp\./.test(name); } // ONE WALK OF THE OLD TREE: each video dir's entries classified by name. A // tierable regular file is left where it is (a link stands in for it); every // other entry is listed for the copy and measured. export async function inventoryTree(dir: string): Promise { const inv: Inventory = { videoIds: [], copyList: [], copy: emptyTallies(), copyBytes: 0, tierable: [], tierableBytes: 0, left: { files: 0, bytes: 0 }, }; const top = (await readdir(dir, { withFileTypes: true })).sort((a, b) => a.name.localeCompare(b.name), ); for (const e of top) { const p = path.join(dir, e.name); if (!e.isDirectory()) { // A stray file at the top of `data/`: text, carried. inv.copyList.push(e.name); add(inv.copy, copyKindOf(e.name), await measureEntry(p)); continue; } inv.videoIds.push(e.name); const names = (await readdir(p)).sort(); for (const name of names) { const fp = path.join(p, name); const st = await lstat(fp); if (isTierable(name) && st.isFile()) { inv.tierable.push({ id: e.name, name, bytes: st.size, atime: st.atime, mtime: st.mtime, }); inv.tierableBytes += st.size; continue; } if (isLeftOnPlatter(name)) { const t = await measureEntry(fp, st); inv.left.files += t.files; inv.left.bytes += t.bytes; continue; } inv.copyList.push(`${e.name}/${name}`); add(inv.copy, copyKindOf(name), await measureEntry(fp, st)); } } inv.copyBytes = COPY_KINDS.reduce((n, k) => n + inv.copy[k].bytes, 0); return inv; } // THE REBUILT TREE, measured the same way, minus the links the migration made // (a link whose name is tierable and whose target is the tier link). A resumed // copy's tree already holds them. export async function measureIncoming(dir: string): Promise { const out = emptyTallies(); let top; try { top = await readdir(dir, { withFileTypes: true }); } catch { return out; } for (const e of top) { const p = path.join(dir, e.name); if (!e.isDirectory()) { add(out, copyKindOf(e.name), await measureEntry(p)); continue; } for (const name of await readdir(p)) { const fp = path.join(p, name); const st = await lstat(fp); if (st.isSymbolicLink() && isTierable(name)) { const target = await readlink(fp).catch(() => ""); if (target === tierLinkTarget(e.name, name)) continue; } add(out, copyKindOf(name), await measureEntry(fp, st)); } } return out; } function tallyText(t: Tally): string { return `${t.files} file(s), ${formatBytes(t.bytes)}`; } function kindsLine(t: KindTallies): string { return COPY_KINDS.map((k) => `${k} ${tallyText(t[k])}`).join("; "); } // --------------------------------------------------------------------------- // The editor must be stopped // --------------------------------------------------------------------------- export const DEFAULT_EDITOR_URL = "http://localhost:3001"; export type EditorCheckOptions = { url?: string; fetchImpl?: typeof fetch; isAlive?: (pid: number) => boolean; // Ms since epoch the machine booted (a pid-less `running` meta from before // it is a ghost). machineBootMs?: number; timeoutMs?: number; }; // WHY THE EDITOR IS (OR MAY BE) RUNNING, or null. Two questions, either // enough: // 1. Does its port answer? `GET /api/pulse`: any HTTP // answer is an editor; a refused connection is none; a connection that // does not answer within the timeout is treated as one (the editor's // main thread can be busy for minutes — release 17 step 0). // 2. Does `.jobs/` show a job in a live process? A `running` or `queued` // meta whose writer (`pid`) is alive and is not this process — the editor, // or an `archilyzer run`. A `running` meta with no pid (written before // release 17) counts when it was written after the machine booted: the // editor's boot pass closes those, so one left is a process that has not // been through it. export async function editorRunningReason( paths: Pick, opts: EditorCheckOptions = {}, ): Promise { const base = (opts.url ?? process.env.ARCHILYZER_EDITOR_URL ?? DEFAULT_EDITOR_URL) .trim() .replace(/\/+$/, ""); const url = `${base}/api/pulse`; const doFetch = opts.fetchImpl ?? fetch; const timeoutMs = opts.timeoutMs ?? 3000; try { const res = await doFetch(url, { signal: AbortSignal.timeout(timeoutMs) }); return `the editor answers at ${url} (HTTP ${res.status})`; } catch (err) { const name = (err as Error | null)?.name; if (name === "TimeoutError" || name === "AbortError") { return ( `something is listening at ${base} but did not answer ${url} within ` + `${Math.round(timeoutMs / 1000)} s — if it is the editor, stop it` ); } /* refused, unresolvable: nothing is listening there */ } const isAlive = opts.isAlive ?? processIsAlive; const bootMs = opts.machineBootMs ?? Date.now() - os.uptime() * 1000; let names: string[]; try { names = await readdir(paths.jobsDir); } catch { return null; } const live: string[] = []; for (const file of names) { if (!file.endsWith(".meta.json")) continue; // A meta a live process wrote was written after the machine booted; an // older file cannot be one, and is not read. try { if ((await stat(path.join(paths.jobsDir, file))).mtimeMs < bootMs) continue; } catch { continue; } const meta = await readJobMeta( paths as Paths, file.slice(0, -".meta.json".length), ); if (!meta || (meta.status !== "running" && meta.status !== "queued")) continue; const what = `job ${meta.id} (${meta.kind}${meta.channelSlug ? `, ${meta.channelSlug}` : ""}) is ${meta.status}`; if (typeof meta.pid === "number") { if (meta.pid !== process.pid && isAlive(meta.pid)) { live.push(`${what} in a live process (pid ${meta.pid})`); } } else if (meta.status === "running") { const at = meta.startedAt ?? meta.queuedAt; if (typeof at === "number" && at >= bootMs) { live.push(`${what} with no recorded process, since after this machine booted`); } } } if (live.length === 0) return null; const more = live.length > 3 ? ` (and ${live.length - 3} more)` : ""; return `${live.slice(0, 3).join("; ")}${more}`; } // --------------------------------------------------------------------------- // One channel // --------------------------------------------------------------------------- export type MigrationStep = | "copied" | "linked" | "marker-swap" | "platter-renamed" | "media-linked" | "data-unlinked" | "data-renamed" | "config-written" | "reclaimed"; export type MigrateDeps = { // Why the editor is running, or null. Default: `editorRunningReason`. editorProbe?: (paths: Paths) => Promise; freeBytes?: (dir: string) => Promise; // Test seam: called after each named step; a throw is a kill there. checkpoint?: (slug: string, step: MigrationStep) => void | Promise; }; export type MigrateSettings = Pick< SiteSettings, "minFreeDiskGB" | "resumeMarginGB" | "storage" >; export type MigrateOptions = { paths: Paths; settings: MigrateSettings; // Slugs named on the command line, or every legacy channel. slugs: string[] | "all"; dryRun?: boolean; reclaim?: boolean; // `--all` goes on past the three big-text channels. includeLarge?: boolean; log?: (line: string) => void; deps?: MigrateDeps; }; export type ChannelOutcome = | "migrated" | "already" | "not-legacy" | "would-migrate" | "no-space" | "refused" | "failed"; export type ChannelReport = { slug: string; outcome: ChannelOutcome; detail?: string; copy?: KindTallies; copyBytes?: number; tierableFiles?: number; tierableBytes?: number; // Dead `*.temp.*` left on the platter (not copied, not linked). leftOnPlatter?: Tally; linksMade?: number; linksAlready?: number; reclaimedFiles?: number; reclaimedBytes?: number; // Platter entries `--reclaim` kept: no same-size copy on the corpus disk. reclaimKept?: string[]; // The phase a marker was resumed from. resumedFrom?: string; }; export type MigrateResult = { exitCode: number; channels: ChannelReport[]; // `--all` stopped before these (the big three), with this much free. stoppedBefore?: string[]; freeBytes?: number; editorRunning?: string | null; }; class Refusal extends Error { constructor( message: string, readonly outcome: "refused" | "no-space" = "refused", ) { super(message); this.name = "Refusal"; } } type Ctx = { paths: Paths; settings: MigrateSettings; dryRun: boolean; reclaim: boolean; log: (line: string) => void; freeBytes: (dir: string) => Promise; checkpoint: (slug: string, step: MigrationStep) => Promise; // A dry run of several channels: the space the earlier ones would take. projectedUse: number; }; // Where a channel stands, read from the corpus disk alone. type Plan = | { kind: "migrate"; D: string; root: string; target: string; marker: RelocationMarker | null } | { kind: "migrated"; mediaDir: string } | { kind: "not-legacy" }; function sameDir(a: string, b: string): boolean { return path.resolve(a) === path.resolve(b); } async function planChannel(ctx: Ctx, slug: string): Promise { const { paths } = ctx; const channelDir = path.join(paths.channelsDir, slug); const config = await readChannelConfig(paths, slug); if (!config) throw new Refusal(`no channel "${slug}" (no readable config.json)`); const dataLink = path.join(channelDir, "data"); const dataState = await linkOrDirState(dataLink); const marker = await readRelocationMarker(paths, slug); if (marker && marker.scope !== "tier-migration") { throw new Refusal( `a media move (${marker.direction}, phase "${marker.phase}", target ` + `${marker.target}) is in flight or was interrupted — finish it from the ` + `channel's Storage panel first; this tool resumes only its own marker`, ); } if (marker) { // RESUMING: the marker is the source of truth for the target — the swap // unsets `dataDir` before the marker is cleared. const target = marker.target; if (path.basename(target) !== "media" || path.basename(path.dirname(target)) !== slug) { throw new Refusal(`its tier-migration marker names ${target}, not /${slug}/media`); } const root = rootOfRelocatedMediaDir(target, slug); const D = relocatedDataDir(root, slug); const recorded = config.dataDir?.trim(); if (recorded && !sameDir(recorded, D)) { throw new Refusal( `its tier-migration marker names ${target} but config.json records dataDir ${recorded}`, ); } return { kind: "migrate", D, root, target, marker }; } const recorded = config.dataDir?.trim(); const linked = dataState.kind === "link" ? path.resolve(channelDir, dataState.linkTarget) : undefined; if (!recorded && !linked) { const mediaDir = config.mediaDir?.trim(); if (mediaDir && dataState.kind !== "other") return { kind: "migrated", mediaDir }; return { kind: "not-legacy" }; } const D = recorded ?? (linked as string); if (recorded && linked && !sameDir(recorded, linked)) { throw new Refusal( `config.json records dataDir ${recorded} but ${dataLink} points at ${linked} — ` + `inconsistent; fix one by hand first`, ); } if (dataState.kind !== "link") { throw new Refusal( `config.json records dataDir ${recorded} but ${dataLink} is ` + `${dataState.kind === "real-dir" ? "a real directory" : dataState.kind} — ` + `inconsistent; fix one by hand first`, ); } if (path.basename(D) !== "data" || path.basename(path.dirname(D)) !== slug) { throw new Refusal(`its dataDir ${D} is not /${slug}/data`); } if (config.mediaDir?.trim()) { throw new Refusal( `config.json records both dataDir ${D} and mediaDir ${config.mediaDir.trim()} — ` + `inconsistent; fix one by hand first`, ); } const root = rootOfRelocatedMediaDir(D, slug); return { kind: "migrate", D, root, target: relocatedMediaDir(root, slug), marker: null }; } async function isRealDir(p: string): Promise { try { return (await stat(p)).isDirectory(); } catch { return false; } } async function writeMarker(ctx: Ctx, slug: string, marker: RelocationMarker): Promise { await writeDirMarker(relocationMarkerPath(ctx.paths, slug), marker); forgetChannelMedia(slug); } // THE VERIFY'S DRY RUN: what rsync would still send, by line. A directory // whose only difference is its time (`.d..t`) is not content. function contentDrift(output: string): string[] { return output .split(/[\r\n]+/) .map((l) => l.trim()) .filter((l) => l.length > 0) .filter((l) => !l.startsWith("sending incremental")) .filter((l) => !/^(sent |total size |building file list|done$)/.test(l)) .filter((l) => !/^\.d\.\.t/.test(l)); } function rsyncListArgs(listFile: string): string[] { // -r EXPLICITLY: --files-from turns off -a's implied recursion, and // `clips/` and a scratch dir are directories to carry whole. return ["-a", "-r", "--from0", `--files-from=${listFile}`]; } async function migrateOne(ctx: Ctx, slug: string): Promise { const { paths, log } = ctx; const plan = await planChannel(ctx, slug); if (plan.kind === "not-legacy") { return { slug, outcome: "not-legacy", detail: "not on the retired layout — nothing to migrate" }; } if (plan.kind === "migrated") { const report: ChannelReport = { slug, outcome: "already", detail: "already on the media tier" }; if (ctx.reclaim) Object.assign(report, await reclaimChannel(ctx, slug, plan.mediaDir)); return report; } const { D, root, target, marker } = plan; const channelDir = path.join(paths.channelsDir, slug); const dataLink = path.join(channelDir, "data"); const incoming = path.join(channelDir, INCOMING_NAME); const listFile = path.join(channelDir, LIST_NAME); const mediaLink = channelMediaLink(paths, slug); let phase = marker?.phase ?? "copy"; const report: ChannelReport = { slug, outcome: "migrated" }; if (marker) { report.resumedFrom = phase; log(`${slug}: resuming a tier migration from phase "${phase}"`); } if (phase === "copy") { // PREFLIGHT. Everything that can refuse does so here, before the marker. if (!(await isRealDir(D))) { throw new Refusal(`${D} is not reachable (is the drive mounted?)`); } await assertRelocationRootPresent(root, ctx.settings.storage, paths); const targetState = await linkOrDirState(target); if (targetState.kind !== "missing") { throw new Refusal( `${target} already exists — a media move or an earlier attempt left it; ` + `look at it before migrating`, ); } const mediaState = await linkOrDirState(mediaLink); if (mediaState.kind !== "missing") { throw new Refusal(`${mediaLink} already exists — look at it before migrating`); } if (!marker && (await linkOrDirState(incoming)).kind !== "missing") { throw new Refusal( `${incoming} exists with no tier-migration marker — a leftover; remove it ` + `by hand after checking it`, ); } log(`${slug}: walking ${D}…`); const inv = await inventoryTree(D); report.copy = inv.copy; report.copyBytes = inv.copyBytes; report.tierableFiles = inv.tierable.length; report.tierableBytes = inv.tierableBytes; report.leftOnPlatter = inv.left; log( `${slug}: ${inv.videoIds.length} video dir(s); to the corpus disk: ` + `${formatBytes(inv.copyBytes)} (${kindsLine(inv.copy)}); staying on the ` + `media drive: ${inv.tierable.length} file(s), ${formatBytes(inv.tierableBytes)}` + (inv.left.files > 0 ? `; dead postprocessor temps left there unlinked (--reclaim deletes them): ${tallyText(inv.left)}` : ""), ); // THE SPACE RULE: the copy plus the resume margin must fit above the disk // gate's floor on the corpus disk — landing the text with nothing to spare // would stop the next download at the gate. The margin is zero when the // gate is off (the mover's rule). A resumed copy needs only what it lacks. const floorGB = ctx.settings.minFreeDiskGB; const marginGB = floorGB > 0 ? ctx.settings.resumeMarginGB : 0; const already = marker ? sumBytes(await measureIncoming(incoming)) : 0; const needed = Math.max(0, inv.copyBytes - already) + marginGB * GB; const free = (await ctx.freeBytes(paths.channelsDir)) - ctx.projectedUse; const usable = free - floorGB * GB; if (needed > usable) { throw new Refusal( `not enough space on the corpus disk: ${formatBytes(Math.max(0, free))} free` + (floorGB > 0 ? ` (${formatBytes(Math.max(0, usable))} above the ${floorGB} GB floor)` : "") + `; ${formatBytes(Math.max(0, inv.copyBytes - already))} to copy` + (marginGB > 0 ? ` plus a ${marginGB} GB resume margin` : ""), "no-space", ); } if (ctx.dryRun) { ctx.projectedUse += Math.max(0, inv.copyBytes - already); report.outcome = "would-migrate"; report.detail = `would copy ${formatBytes(inv.copyBytes - already)} to ${incoming}, link ` + `${inv.tierable.length} file(s), rename ${D} to ${target}` + (ctx.reclaim ? `, then reclaim the platter's text copies` : ""); return report; } // COPY. await writeMarker(ctx, slug, { target, direction: "out", startedAt: marker?.startedAt || new Date().toISOString(), phase: "copy", scope: "tier-migration", }); await writeFile(listFile, inv.copyList.map((p) => `${p}\0`).join("")); await mkdir(incoming).catch((err: NodeJS.ErrnoException) => { if (err.code !== "EEXIST") throw err; }); const sink = makeProgressSink({ totalBytes: inv.copyBytes, alreadyBytes: already, log }); const copied = await rsyncTree({ rsyncBin: paths.rsyncBin, src: D, dest: incoming, args: [...rsyncListArgs(listFile), "--partial", "--info=progress2"], log, progress: sink, }); if (copied.exitCode !== 0) { throw new Error( `rsync failed (exit ${copied.exitCode}); nothing on ${D} was touched — run again to resume`, ); } // VERIFY: rsync has nothing left to send… const dry = await rsyncTree({ rsyncBin: paths.rsyncBin, src: D, dest: incoming, args: [...rsyncListArgs(listFile), "--dry-run", "--itemize-changes"], log: (m) => { if (m.startsWith("$ ")) log(m); }, }); if (dry.exitCode !== 0) throw new Error(`the verification rsync failed (exit ${dry.exitCode})`); const drift = contentDrift(dry.output); if (drift.length > 0) { throw new Error( `verification failed: ${drift.length} item(s) still differ (first: ${drift[0]}); ` + `nothing on ${D} was touched`, ); } // …and the rebuilt tree holds what the old one does, kind by kind. const got = await measureIncoming(incoming); const off = COPY_KINDS.filter( (k) => got[k].files !== inv.copy[k].files || got[k].bytes !== inv.copy[k].bytes, ); if (off.length > 0) { throw new Error( `verification failed: ${off .map((k) => `${k} ${tallyText(got[k])} copied, ${tallyText(inv.copy[k])} on ${D}`) .join("; ")}; nothing on ${D} was touched`, ); } log(`${slug}: copy verified (${kindsLine(got)})`); await ctx.checkpoint(slug, "copied"); // THE LINKS, each carrying its file's times. let made = 0; let had = 0; for (const id of inv.videoIds) { await mkdir(path.join(incoming, id)).catch((err: NodeJS.ErrnoException) => { if (err.code !== "EEXIST") throw err; }); } for (const f of inv.tierable) { const p = path.join(incoming, f.id, f.name); const want = tierLinkTarget(f.id, f.name); try { await symlink(want, p); made++; } catch (err) { if ((err as NodeJS.ErrnoException).code !== "EEXIST") throw err; const st = await linkOrDirState(p); if (st.kind !== "link" || st.linkTarget !== want) { throw new Error(`${p} exists and is not the link ${want}`); } had++; } await lutimes(p, f.atime, f.mtime); } report.linksMade = made; report.linksAlready = had; log(`${slug}: ${made} link(s) made${had ? `, ${had} already there` : ""}`); await ctx.checkpoint(slug, "linked"); await writeMarker(ctx, slug, { target, direction: "out", startedAt: marker?.startedAt || new Date().toISOString(), phase: "swap", scope: "tier-migration", }); phase = "swap"; await ctx.checkpoint(slug, "marker-swap"); } else if (ctx.dryRun) { report.outcome = "would-migrate"; report.detail = `would resume its tier migration from phase "${phase}"`; return report; } if (phase === "swap") { await swap(ctx, slug, { D, target, dataLink, incoming, mediaLink }); await writeFile( path.join(channelDir, RECLAIM_NOTE), JSON.stringify({ target, migratedAt: new Date().toISOString() }) + "\n", ); await rm(listFile, { force: true }); phase = "reclaim"; } if (ctx.reclaim || marker?.phase === "reclaim") { await writeMarker(ctx, slug, { target, direction: "out", startedAt: marker?.startedAt || new Date().toISOString(), phase: "reclaim", scope: "tier-migration", }); Object.assign(report, await reclaimTree(ctx, slug, target)); await ctx.checkpoint(slug, "reclaimed"); await rm(path.join(channelDir, RECLAIM_NOTE), { force: true }); } await clearDirMarker(relocationMarkerPath(paths, slug)); forgetChannelMedia(slug); const after = await inspectChannelMedia(paths, slug, undefined, { fresh: true }); if (after.status !== "ok") { report.outcome = "failed"; report.detail = `migrated, but its media reads ${after.status}: ${after.detail ?? ""}`; } return report; } function sumBytes(t: KindTallies): number { return COPY_KINDS.reduce((n, k) => n + t[k].bytes, 0); } // THE SWAP, every step asking the disk what is there first, so a rerun after a // crash between any two of them does the rest and nothing twice. async function swap( ctx: Ctx, slug: string, p: { D: string; target: string; dataLink: string; incoming: string; mediaLink: string }, ): Promise { const { log } = ctx; // 1. The platter: //data -> //media. const d = await linkOrDirState(p.D); const t = await linkOrDirState(p.target); if (d.kind === "real-dir" && t.kind === "missing") { await rename(p.D, p.target); log(`${slug}: renamed ${p.D} -> ${p.target}`); } else if (!(d.kind === "missing" && t.kind === "real-dir")) { throw new Error( `cannot swap: ${p.D} is ${d.kind} and ${p.target} is ${t.kind} — the drive may ` + `not be mounted; nothing was changed by this step`, ); } await ctx.checkpoint(slug, "platter-renamed"); // 2. The corpus disk: channels//media -> //media. const m = await linkOrDirState(p.mediaLink); if (m.kind === "missing") { await symlink(p.target, p.mediaLink); } else if (!(m.kind === "link" && sameDir(path.resolve(path.dirname(p.mediaLink), m.linkTarget), p.target))) { throw new Error(`cannot swap: ${p.mediaLink} is ${m.kind}, not the link to ${p.target}`); } await ctx.checkpoint(slug, "media-linked"); // 3. The old whole-directory link goes. const dl = await linkOrDirState(p.dataLink); if (dl.kind === "link") { if (!sameDir(path.resolve(path.dirname(p.dataLink), dl.linkTarget), p.D)) { throw new Error(`cannot swap: ${p.dataLink} points at ${dl.linkTarget}, not ${p.D}`); } await unlink(p.dataLink); } else if (dl.kind === "other") { throw new Error(`cannot swap: ${p.dataLink} is neither a link nor a directory`); } await ctx.checkpoint(slug, "data-unlinked"); // 4. The rebuilt text tree takes its name. const inc = await linkOrDirState(p.incoming); const now = await linkOrDirState(p.dataLink); if (inc.kind === "real-dir" && now.kind === "missing") { await rename(p.incoming, p.dataLink); } else if (!(inc.kind === "missing" && now.kind === "real-dir")) { throw new Error(`cannot swap: ${p.incoming} is ${inc.kind} and ${p.dataLink} is ${now.kind}`); } await ctx.checkpoint(slug, "data-renamed"); // 5. The record. await patchChannelConfig(ctx.paths, slug, { mediaDir: p.target }, { unset: ["dataDir"] }); forgetChannelMedia(slug); log(`${slug}: config.json records mediaDir ${p.target}; dataDir removed`); await ctx.checkpoint(slug, "config-written"); } // RECLAIM: the platter's copies of the text. In `//media//` // a tierable regular file is the media and stays; a dead `*.temp.*` (left // there by the migration, never copied) goes; every other entry goes ONLY when // one `lstat` of its twin on the corpus disk — `channels//data//` // — finds it there, a directory for a directory or a file of the same size. // Anything else is KEPT and reported: the platter copy may be the only one (a // reclaim long after the migration, a file changed or removed since). A video // dir left empty goes (non-recursive `rmdir`). A top-level non-directory entry // is a stray file copied from the top of the old tree, checked the same way. async function reclaimTree( ctx: Ctx, slug: string, mediaDir: string, ): Promise> { const ssd = path.join(ctx.paths.channelsDir, slug, "data"); let files = 0; let bytes = 0; const kept: string[] = []; const hasTwin = async (rel: string, st: Stats): Promise => { let twin: Stats; try { twin = await lstat(path.join(ssd, rel)); } catch { return false; } if (st.isDirectory()) return twin.isDirectory(); return !twin.isDirectory() && twin.size === st.size; }; const take = async (p: string, st: Stats): Promise => { const t = await measureEntry(p, st); files += t.files; bytes += t.bytes; if (!ctx.dryRun) await rm(p, { recursive: true, force: true }); }; const top = await readdir(mediaDir, { withFileTypes: true }); for (const e of top) { const p = path.join(mediaDir, e.name); if (!e.isDirectory()) { const st = await lstat(p); if (await hasTwin(e.name, st)) await take(p, st); else kept.push(e.name); continue; } for (const name of await readdir(p)) { const fp = path.join(p, name); const st = await lstat(fp); if (isTierable(name) && st.isFile()) continue; const rel = `${e.name}/${name}`; if (isLeftOnPlatter(name) || (await hasTwin(rel, st))) await take(fp, st); else kept.push(rel); } if (!ctx.dryRun) await rmdir(p).catch(() => {}); } ctx.log( `${slug}: ${ctx.dryRun ? "would reclaim" : "reclaimed"} ${files} file(s), ` + `${formatBytes(bytes)} of text copies and dead temps from ${mediaDir}`, ); if (kept.length > 0) { ctx.log( `${slug}: kept ${kept.length} entr${kept.length === 1 ? "y" : "ies"} on the drive with no ` + `same-size copy on the corpus disk (${kept.slice(0, 5).join(", ")}` + `${kept.length > 5 ? `, … ${kept.length - 5} more` : ""})`, ); } return { reclaimedFiles: files, reclaimedBytes: bytes, reclaimKept: kept }; } // --reclaim on a channel already on the media tier: only one this tool // migrated (its reclaim note stands), whose media drive answers and whose // `data/` is a real directory. async function reclaimChannel( ctx: Ctx, slug: string, mediaDir: string, ): Promise> { const note = path.join(ctx.paths.channelsDir, slug, RECLAIM_NOTE); if ((await linkOrDirState(note)).kind === "missing") { throw new Refusal( `cannot reclaim: no ${RECLAIM_NOTE} beside its config — migrate-tier did not ` + `migrate it, or its reclaim is done; its drive holds no text copies of the ` + `retired layout to take`, ); } const loc = await inspectChannelMedia(ctx.paths, slug, undefined, { fresh: true }); if (loc.status !== "ok") { throw new Refusal( `cannot reclaim: its media reads ${loc.status}${loc.detail ? ` (${loc.detail})` : ""}`, ); } if (!(await isRealDir(path.join(ctx.paths.channelsDir, slug, "data")))) { throw new Refusal(`cannot reclaim: its data/ is not a real directory`); } if (ctx.dryRun) return reclaimTree(ctx, slug, mediaDir); await writeMarker(ctx, slug, { target: mediaDir, direction: "out", startedAt: new Date().toISOString(), phase: "reclaim", scope: "tier-migration", }); const out = await reclaimTree(ctx, slug, mediaDir); await ctx.checkpoint(slug, "reclaimed"); await rm(note, { force: true }); await clearDirMarker(relocationMarkerPath(ctx.paths, slug)); forgetChannelMedia(slug); return out; } // --------------------------------------------------------------------------- // The run // --------------------------------------------------------------------------- // Every channel this tool has work on: on the retired layout (a `data` link or // a recorded `dataDir`), or carrying its own marker. export async function legacyChannels(paths: Paths): Promise { const entries = await readdir(paths.channelsDir, { withFileTypes: true }).catch(() => []); const out: string[] = []; for (const e of entries) { if (!e.isDirectory()) continue; const slug = e.name; const marker = await readRelocationMarker(paths, slug); if (marker?.scope === "tier-migration") { out.push(slug); continue; } const config = await readChannelConfig(paths, slug).catch(() => null); if (!config) continue; const data = await linkOrDirState(path.join(paths.channelsDir, slug, "data")); if (config.dataDir?.trim() || data.kind === "link") out.push(slug); } return out.sort(); } // Channels this tool migrated whose drive still holds the text copies (their // reclaim note stands), for `--all --reclaim`. async function reclaimableChannels(paths: Paths): Promise { const entries = await readdir(paths.channelsDir, { withFileTypes: true }).catch(() => []); const out: string[] = []; for (const e of entries) { if (!e.isDirectory()) continue; const note = await linkOrDirState(path.join(paths.channelsDir, e.name, RECLAIM_NOTE)); if (note.kind !== "missing") out.push(e.name); } return out.sort(); } // The text bytes a legacy channel would copy, for `--order smallest`; a // channel whose tree cannot be walked sorts last (and is refused in its turn). async function copyBytesOf(paths: Paths, slug: string): Promise { const marker = await readRelocationMarker(paths, slug); if (marker?.scope === "tier-migration") return -1; // resume these first const config = await readChannelConfig(paths, slug).catch(() => null); const data = path.join(paths.channelsDir, slug, "data"); const D = config?.dataDir?.trim() || (await readlink(data).then((t) => path.resolve(path.dirname(data), t)).catch(() => "")); if (!D) return Number.MAX_SAFE_INTEGER; try { return (await inventoryTree(D)).copyBytes; } catch { return Number.MAX_SAFE_INTEGER; } } function describe(r: ChannelReport): string { const parts: string[] = [`${r.slug}: ${r.outcome}`]; if (r.detail) parts.push(`— ${r.detail}`); return parts.join(" "); } export async function migrateMediaTier(opts: MigrateOptions): Promise { const log = opts.log ?? ((l: string) => console.log(l)); const deps = opts.deps ?? {}; const ctx: Ctx = { paths: opts.paths, settings: opts.settings, dryRun: opts.dryRun === true, reclaim: opts.reclaim === true, log, freeBytes: deps.freeBytes ?? getFreeBytes, checkpoint: async (slug, step) => { await deps.checkpoint?.(slug, step); }, projectedUse: 0, }; const result: MigrateResult = { exitCode: 0, channels: [] }; const probe = deps.editorProbe ?? ((p: Paths) => editorRunningReason(p)); const running = await probe(opts.paths); result.editorRunning = running; if (running) { if (!ctx.dryRun) { log( `Refusing: the editor must be stopped for the migration — ${running}. ` + `It holds its job registry in memory, so nothing on disk can say a job of ` + `its is about to write into a channel. Stop it and run this again.`, ); result.exitCode = 2; return result; } log(`Note: ${running}. A dry run reads only; a real run would refuse until it is stopped.`); } if (ctx.dryRun) log("DRY RUN — nothing will be written."); let order: string[]; let deferred: string[] = []; // Every legacy channel's copy bytes, from ONE walk each: the order, and the // stop's report on the big three. const sizes = new Map(); if (opts.slugs === "all") { const legacy = await legacyChannels(opts.paths); const large = legacy.filter((s) => LARGE_TEXT_CHANNELS.includes(s)); const rest = legacy.filter((s) => !LARGE_TEXT_CHANNELS.includes(s)); const measure = async (slugs: string[]) => { const sized: { slug: string; bytes: number }[] = []; for (const slug of slugs) { const bytes = await copyBytesOf(opts.paths, slug); sizes.set(slug, bytes); sized.push({ slug, bytes }); } return sized .sort((a, b) => a.bytes - b.bytes || a.slug.localeCompare(b.slug)) .map((s) => s.slug); }; log(`${legacy.length} channel(s) on the retired layout: ${legacy.join(", ") || "none"}`); order = await measure(rest); const largeOrder = await measure(large); if (opts.includeLarge) order.push(...largeOrder); else deferred = largeOrder; if (ctx.reclaim) { for (const slug of await reclaimableChannels(opts.paths)) { if (!order.includes(slug)) order.push(slug); } } } else { order = opts.slugs; } for (const slug of order) { let report: ChannelReport; try { report = await migrateOne(ctx, slug); } catch (err) { const refusal = err instanceof Refusal; report = { slug, outcome: refusal ? (err as Refusal).outcome : "failed", detail: (err as Error).message, }; } result.channels.push(report); log(describe(report)); if (report.outcome === "migrated" || report.outcome === "already") { if (report.copy) { log( ` ${slug}: copied ${formatBytes(report.copyBytes ?? 0)} (${kindsLine(report.copy)}); ` + `${report.linksMade ?? 0} link(s) made; ${report.tierableFiles ?? 0} media file(s), ` + `${formatBytes(report.tierableBytes ?? 0)} stay on the media drive` + (report.leftOnPlatter?.files ? `; ${tallyText(report.leftOnPlatter)} of dead temps left there for --reclaim` : ""), ); } } const bad = report.outcome === "refused" || report.outcome === "failed" || report.outcome === "no-space"; if (bad) { result.exitCode = 1; // A real run stops at the first channel it could not finish: the // operator decides. A dry run reports every channel. if (!ctx.dryRun) { if (opts.slugs === "all") { log( `Stopped at ${slug}. Channels done so far are migrated; run ` + `\`archilyzer storage migrate-tier --all\` again after fixing it (finished ` + `channels are a no-op).`, ); } return result; } } } if (deferred.length > 0) { result.stoppedBefore = deferred; const free = await ctx.freeBytes(opts.paths.channelsDir); result.freeBytes = free; log(""); log( `STOPPED before the big-text channels: ${deferred.join(", ")}. The corpus disk has ` + `${formatBytes(free - (ctx.dryRun ? ctx.projectedUse : 0))} free${ctx.dryRun ? " (after the channels above)" : ""}; ` + `the gate floor is ${opts.settings.minFreeDiskGB} GB.`, ); for (const slug of deferred) { const bytes = sizes.get(slug) ?? (await copyBytesOf(opts.paths, slug)); log( ` ${slug}: ${ bytes < 0 ? "a tier migration is in flight — resume it by name" : bytes === Number.MAX_SAFE_INTEGER ? "its old tree could not be walked" : `${formatBytes(bytes)} to copy` }`, ); } log( `Continue with \`archilyzer storage migrate-tier \` one at a time, or ` + `\`archilyzer storage migrate-tier --all --include-large\`.`, ); } return result; } // --------------------------------------------------------------------------- // The command line // --------------------------------------------------------------------------- const USAGE = "usage: archilyzer storage migrate-tier … | --all [--order smallest] " + "[--include-large] [--dry-run] [--reclaim]"; export async function main(argv: string[]): Promise { const BOOLS = ["all", "dry-run", "reclaim", "include-large", "help"]; const { flags, positionals } = parseArgv(argv, BOOLS); if (flags.help) { console.log(USAGE); return 0; } const known = new Set([...BOOLS, "order"]); const unknown = Object.keys(flags).filter((k) => !known.has(k)); if (unknown.length > 0) { console.error(`migrate-tier: unknown flag(s) ${unknown.map((k) => `--${k}`).join(", ")}\n${USAGE}`); return 2; } const all = flags.all === true; if (all === positionals.length > 0) { console.error(`migrate-tier: name channel slugs or --all, not both and not neither\n${USAGE}`); return 2; } if (flags.order !== undefined && (flags.order !== "smallest" || !all)) { console.error(`migrate-tier: --order takes "smallest", with --all\n${USAGE}`); return 2; } if (flags["include-large"] && !all) { console.error(`migrate-tier: --include-large goes with --all (a named channel always runs)\n${USAGE}`); return 2; } const { getPaths } = await import("../lib/paths"); const { getSettings } = await import("../lib/settings"); const paths = getPaths(); const settings = getSettings(); console.log(`channels: ${paths.channelsDir}`); const res = await migrateMediaTier({ paths, settings, slugs: all ? "all" : positionals, dryRun: flags["dry-run"] === true, reclaim: flags.reclaim === true, includeLarge: flags["include-large"] === true, }); return res.exitCode; } runIfEntryPoint(import.meta.url, () => main(process.argv.slice(2)));