Archilyzer · Source

archilyzer

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

commit f698dc0f88b04633856e3d694b88054ee89a6a8b
parent 1eb37bde03bab3949985c10cf14e6cf7047e91ea
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Fri,  2 Oct 2026 02:10:35 -0400

common: archilyzer storage migrate-tier — the one-off migration off the retired whole-directory layout

A legacy channel (data/ an absolute link to <root>/<slug>/data, config.dataDir)
gets its text copied home to channels/<slug>/data.incoming from a NUL list the
classifier wrote (text, clips/, scratch and source-media.*; never a glob),
verified by an itemized dry run and per-kind counts and bytes; every tierable
file gets its relative link carrying the file's times; the swap renames the
platter's data to media, links channels/<slug>/media, drops the old data link,
renames data.incoming to data and records mediaDir with dataDir unset, each
step looking at the disk first. Resumable from the tier-migration marker's
phase; idempotent; a dry run writes nothing; --reclaim deletes the platter's
text copies. A real run refuses while the editor's port answers or .jobs/
shows a job in a live process; --all goes smallest text first and stops
before the three big-text channels, printing the free space and the command
to continue (--include-large, or per slug). Fixture tests over the real rsync.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

Diffstat:
Mcommon/bin/archilyzer.ts | 2++
Acommon/bin/migrate-media-tier.test.ts | 567+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/bin/migrate-media-tier.ts | 1183+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
3 files changed, 1752 insertions(+), 0 deletions(-)

diff --git a/common/bin/archilyzer.ts b/common/bin/archilyzer.ts @@ -267,6 +267,8 @@ export const COMMANDS: Command[] = [ "--channel <slug> list duplicate and missing transcripts"), script(["migrate", "channel-priority"], "migrate-channel-priority.ts", "[--dry-run] the one-shot channel-priority migration (plans/channel-priority.md, S5)"), + script(["storage", "migrate-tier"], "migrate-media-tier.ts", + "<slug>…|--all [--order smallest] [--include-large] [--dry-run] [--reclaim] bring a channel off the retired whole-directory layout onto the media tier, its text home to the corpus disk (editor stopped; --all stops before the three big-text channels)"), { path: ["brand", "media"], usage: "[--out <dir>] [--video-kit] render the Archilyzer Media channel's assets", diff --git a/common/bin/migrate-media-tier.test.ts b/common/bin/migrate-media-tier.test.ts @@ -0,0 +1,567 @@ +import { after, test } from "node:test"; +import assert from "node:assert/strict"; +import { + lstatSync, + mkdirSync, + mkdtempSync, + readdirSync, + readFileSync, + readlinkSync, + renameSync, + rmSync, + statSync, + symlinkSync, + utimesSync, + writeFileSync, +} from "node:fs"; +import { tmpdir } from "node:os"; +import path from "node:path"; + +// Run with: +// pnpm --filter yt-dlp-transcript-common exec tsx --test bin/migrate-media-tier.test.ts +// +// The one-off media-tier migration over a fixture: a "platter" dir beside a +// corpus, a legacy channel (its `data/` an absolute link to +// `<platter>/<slug>/data`, `config.dataDir`) holding text, audio, a raw live +// chat, a clip window, a source container and scratch. The REAL rsync runs. +// The environment is pinned before any module reads it (getPaths caches), as +// buildIndex.test.ts does, so the index case can build over the same corpus. + +const ROOT = mkdtempSync(path.join(tmpdir(), "migrate-tier-")); +const PINNED: Record<string, string> = { + TRANSCRIPTS_DIR: path.join(ROOT, "transcripts"), + SAVED_VIDEOS_DIR: path.join(ROOT, "saved-videos"), + SITES_DIR: path.join(ROOT, "transcripts", "sites"), + SETTINGS_FILE: path.join(ROOT, "settings.json"), + EXPORT_PUBLIC_DIR: path.join(ROOT, "public"), + EXPORT_INDEX_DIR: path.join(ROOT, ".export-index"), + EXPORT_BUILDS_DIR: path.join(ROOT, ".export-builds"), + EDITOR_CHANGELOG_FILE: path.join(ROOT, "editor-CHANGELOG.md"), + EXPORT_CHANGELOG_FILE: path.join(ROOT, "export-CHANGELOG.md"), + CHARTS_CONFIG_FILE: path.join(ROOT, "chart-templates.json"), + SEARCH_ALIASES_FILE: path.join(ROOT, "transcripts", "search-aliases.json"), + CURATED_TAGS_FILE: path.join(ROOT, "transcripts", "tags.json"), + ARCHILYZER_CONFIG_DIR: path.join(ROOT, "config"), + ARCHILYZER_SOURCE_SCRATCH: path.join(ROOT, "source-scratch"), +}; +Object.assign(process.env, PINNED); +delete process.env.ARCHILYZER_INDEX_ALLOW_HELD; +after(() => rmSync(ROOT, { recursive: true, force: true })); + +const { getPaths } = await import("../lib/paths"); +const { + editorRunningReason, + legacyChannels, + migrateMediaTier, + INCOMING_NAME, + LIST_NAME, +} = await import("./migrate-media-tier"); +const { inspectChannelMedia, readRelocationMarker } = await import("../lib/channelMedia"); +const { classifyEntry } = await import("../lib/mediaTier"); +const { readChannelConfig } = await import("../controller/channels"); +const { buildIndex } = await import("../controller/buildIndex"); +const { isLiveChatCuesFresh, normalizeLiveChat } = await import( + "../controller/normalizeLiveChat" +); + +type MigrateOpts = Parameters<typeof migrateMediaTier>[0]; +type MigrationStep = NonNullable< + NonNullable<MigrateOpts["deps"]>["checkpoint"] +> extends (slug: string, step: infer S) => unknown + ? S + : never; + +const paths = getPaths(); +const PLATTER = path.join(ROOT, "platter"); +const SLUG = "legacy-ch"; +const SITE = "testsite"; + +const OLD = new Date("2021-03-04T05:06:07.000Z"); +const CUES_AT = new Date("2021-03-05T05:06:07.000Z"); + +// Settings with the disk gate off, so the tmpfs's free space is not the +// subject (relocateChannelMedia.test.ts says why at length); the free-space +// cases inject `freeBytes` and turn the floor on. +const SETTINGS = { + minFreeDiskGB: 0, + resumeMarginGB: 2, + storage: { locations: [], defaultLocationId: "" }, +} as unknown as MigrateOpts["settings"]; + +const CHAT_LINE = + JSON.stringify({ + replayChatItemAction: { + actions: [ + { + addChatItemAction: { + item: { + liveChatTextMessageRenderer: { + message: { runs: [{ text: "hello" }] }, + authorName: { simpleText: "a" }, + timestampUsec: "1000000", + }, + }, + }, + }, + ], + videoOffsetTimeMsec: "1000", + }, + }) + "\n"; + +const VTT = + "WEBVTT\nKind: captions\nLanguage: en\n\n" + + "00:00:00.000 --> 00:00:05.000\nFirst caption line.\n\n" + + "00:01:00.000 --> 00:01:50.000\nSecond caption line.\n"; + +// One video dir's files: name -> content. `clips/` is a directory. +const VIDEOS: Record<string, Record<string, string>> = { + v1: { + "transcript.en.vtt": VTT, + "audio.mp3": "A".repeat(4096), + "transcript.live_chat.json": CHAT_LINE.repeat(20), + "source-media.mp4": "S".repeat(2048), + "audio.m4a.part": "P".repeat(512), + "audio.tmp-1234.mp3": "T".repeat(256), + "clips/w-0001.mp4": "C".repeat(1024), + }, + v2: { + "transcript.json": JSON.stringify({ segments: [{ start: 0, end: 1, text: "hi there" }] }), + "audio.opus": "O".repeat(3000), + }, + v3: {}, +}; + +function writeJson(file: string, value: unknown): void { + mkdirSync(path.dirname(file), { recursive: true }); + writeFileSync(file, JSON.stringify(value, null, 2)); +} + +function channelDir(slug = SLUG): string { + return path.join(paths.channelsDir, slug); +} +function oldTree(slug = SLUG): string { + return path.join(PLATTER, slug, "data"); +} +function mediaTarget(slug = SLUG): string { + return path.join(PLATTER, slug, "media"); +} + +function writeConfig(slug: string, extra: Record<string, unknown> = {}): void { + writeJson(path.join(channelDir(slug), "config.json"), { + handling: "youtube", + name: slug, + url: `https://www.youtube.com/@${slug}/videos`, + ...extra, + }); +} + +function resetCorpus(channels: string[] = [SLUG]): void { + for (const p of [paths.transcriptsDir, PINNED.EXPORT_INDEX_DIR, PLATTER]) { + rmSync(p, { recursive: true, force: true }); + } + mkdirSync(paths.channelsDir, { recursive: true }); + mkdirSync(PLATTER, { recursive: true }); + writeFileSync(paths.settingsFile, JSON.stringify({})); + writeJson(path.join(paths.sitesDir, SITE, "site.json"), { + siteId: SITE, + siteTitle: "Test Site", + siteDescription: "fixture", + headerTitle: "Test Site", + homeTagline: "", + socialLinks: [], + groups: [{ id: "default", name: "All channels", selectedByDefault: true }], + defaultGroupId: "default", + channels: channels.map((slug) => ({ slug, groupId: "default" })), + }); +} + +// Write the videos into `dataDir` (every file at OLD, the cues a day later). +function writeVideos(dataDir: string, slug: string, scale = 1): void { + for (const [id, files] of Object.entries(VIDEOS)) { + const dir = path.join(dataDir, id); + mkdirSync(dir, { recursive: true }); + writeJson(path.join(dir, "metadata.info.json"), { + id, + title: `Video ${id}`, + channel: slug, + upload_date: "20260601", + duration: 120, + description: "fixture", + webpage_url: `https://www.youtube.com/watch?v=${id}`, + extractor_key: "Youtube", + }); + utimesSync(path.join(dir, "metadata.info.json"), OLD, OLD); + for (const [name, content] of Object.entries(files)) { + const file = path.join(dir, name); + mkdirSync(path.dirname(file), { recursive: true }); + writeFileSync(file, content.repeat(scale)); + utimesSync(file, OLD, OLD); + } + } +} + +// A channel on the RETIRED layout, as the pre-release-17 mover left it. +function seedLegacy(slug = SLUG, scale = 1): void { + writeVideos(oldTree(slug), slug, scale); + writeConfig(slug, { dataDir: oldTree(slug) }); + symlinkSync(oldTree(slug), path.join(channelDir(slug), "data")); +} + +// What is on disk under ROOT's corpus and platter: type, size, mtime, link +// target, content of small files. Equal snapshots = nothing was written. +function snapshot(): Record<string, string> { + const out: Record<string, string> = {}; + const walk = (p: string) => { + const st = lstatSync(p); + const rel = path.relative(ROOT, p); + if (st.isSymbolicLink()) { + out[rel] = `link ${readlinkSync(p)} ${st.mtimeMs}`; + return; + } + if (st.isDirectory()) { + out[rel] = `dir`; + for (const n of readdirSync(p)) walk(path.join(p, n)); + return; + } + out[rel] = `file ${st.size} ${st.mtimeMs} ${readFileSync(p, "utf8").slice(0, 64)}`; + }; + walk(paths.transcriptsDir); + walk(PLATTER); + return out; +} + +function run(over: Partial<MigrateOpts> = {}, lines: string[] = []) { + return migrateMediaTier({ + paths, + settings: SETTINGS, + slugs: [SLUG], + log: (l) => lines.push(l), + ...over, + deps: { + editorProbe: async () => null, + ...over.deps, + }, + }); +} + +// THE MIGRATED STATE, every rule the slice names. +async function assertMigrated(slug = SLUG, scale = 1): Promise<void> { + const ch = channelDir(slug); + const config = await readChannelConfig(paths, slug); + assert.equal(config?.mediaDir, mediaTarget(slug)); + assert.equal(config?.dataDir, undefined, "dataDir unset"); + assert.ok(lstatSync(path.join(ch, "data")).isDirectory(), "data/ is a real dir"); + assert.equal(readlinkSync(path.join(ch, "media")), mediaTarget(slug)); + assert.throws(() => lstatSync(oldTree(slug)), "the old tree is renamed"); + assert.throws(() => lstatSync(path.join(ch, INCOMING_NAME))); + assert.throws(() => lstatSync(path.join(ch, LIST_NAME))); + assert.equal(await readRelocationMarker(paths, slug), null, "no marker"); + + for (const [id, files] of Object.entries(VIDEOS)) { + const dir = path.join(ch, "data", id); + const meta = lstatSync(path.join(dir, "metadata.info.json")); + assert.ok(meta.isFile()); + assert.equal(meta.mtimeMs, OLD.getTime(), `${id} metadata keeps its mtime`); + for (const name of Object.keys(files)) { + const top = name.split("/")[0]; + const p = path.join(dir, name); + const st = lstatSync(p); + if (name === "audio.mp3" || name === "audio.opus" || name === "transcript.live_chat.json") { + assert.ok(st.isSymbolicLink(), `${id}/${name} is a link`); + assert.equal(readlinkSync(p), `../../media/${id}/${name}`, "relative"); + assert.equal(st.mtimeMs, OLD.getTime(), `${id}/${name}: the link carries the file's mtime`); + assert.equal(readFileSync(p, "utf8"), files[name].repeat(scale), "readable through the link"); + assert.ok(lstatSync(path.join(mediaTarget(slug), id, name)).isFile(), "bytes on the platter"); + } else { + assert.ok(st.isFile(), `${id}/${name} stays real (${top})`); + assert.equal(st.mtimeMs, OLD.getTime(), `${id}/${name} keeps its mtime`); + assert.equal(readFileSync(p, "utf8"), files[name].repeat(scale)); + } + } + } + const loc = await inspectChannelMedia(paths, slug, undefined, { fresh: true }); + assert.equal(loc.status, "ok", loc.detail); + assert.equal(loc.text.readable, true); +} + +test("clips/ is text by the classifier, and so is carried to the corpus disk", () => { + assert.equal(classifyEntry("clips"), "text"); +}); + +test("a dry run writes nothing, and says what it would do", async () => { + resetCorpus(); + seedLegacy(); + const before = snapshot(); + const lines: string[] = []; + const res = await run({ dryRun: true }, lines); + assert.equal(res.exitCode, 0); + const [r] = res.channels; + assert.equal(r.outcome, "would-migrate", r.detail); + assert.equal(r.tierableFiles, 3); + assert.equal(r.copy?.clips.files, 1); + assert.equal(r.copy?.source.files, 1, "source-media is carried, not tiered"); + assert.equal(r.copy?.scratch.files, 2); + assert.deepEqual(snapshot(), before); + assert.ok(lines.some((l) => l.includes("DRY RUN"))); +}); + +test("a run migrates the channel; a rerun is a no-op", async () => { + resetCorpus(); + seedLegacy(); + const lines: string[] = []; + const res = await run({}, lines); + assert.equal(res.exitCode, 0, lines.join("\n")); + const [r] = res.channels; + assert.equal(r.outcome, "migrated", r.detail); + assert.equal(r.linksMade, 3); + assert.equal(r.copy?.clips.files, 1); + await assertMigrated(); + // The live chat's cues (none here) aside, the raw replay's link answers the + // freshness check from the corpus disk with the file's time. + assert.equal(lstatSync(path.join(channelDir(), "data", "v1", "transcript.live_chat.json")).mtimeMs, OLD.getTime()); + // Without --reclaim the platter keeps its text copies. + assert.ok(lstatSync(path.join(mediaTarget(), "v1", "transcript.en.vtt")).isFile()); + + const before = snapshot(); + const again = await run(); + assert.equal(again.exitCode, 0); + assert.equal(again.channels[0].outcome, "already"); + assert.deepEqual(snapshot(), before, "the rerun writes nothing"); +}); + +test("--reclaim deletes the platter's text copies and keeps the media", async () => { + resetCorpus(); + seedLegacy(); + const res = await run({ reclaim: true }); + assert.equal(res.exitCode, 0, res.channels[0].detail); + const r = res.channels[0]; + assert.equal(r.outcome, "migrated"); + assert.ok((r.reclaimedFiles ?? 0) > 0); + await assertMigrated(); + assert.deepEqual(readdirSync(path.join(mediaTarget(), "v1")).sort(), [ + "audio.mp3", + "transcript.live_chat.json", + ]); + assert.deepEqual(readdirSync(path.join(mediaTarget(), "v2")), ["audio.opus"]); + assert.throws(() => lstatSync(path.join(mediaTarget(), "v3")), "an emptied dir goes"); + + // Reclaim on its own, later: nothing left to take, nothing touched. + const before = snapshot(); + const again = await run({ reclaim: true }); + assert.equal(again.channels[0].outcome, "already"); + assert.equal(again.channels[0].reclaimedFiles, 0); + assert.deepEqual(snapshot(), before); +}); + +test("--reclaim after a plain run takes the copies a first run left", async () => { + resetCorpus(); + seedLegacy(); + await run(); + const res = await run({ reclaim: true, dryRun: true }); + assert.equal(res.channels[0].outcome, "already"); + const wouldTake = res.channels[0].reclaimedFiles ?? 0; + assert.ok(wouldTake > 0); + assert.ok(lstatSync(path.join(mediaTarget(), "v1", "transcript.en.vtt")).isFile(), "dry: kept"); + const real = await run({ reclaim: true }); + assert.equal(real.channels[0].reclaimedFiles, wouldTake); + assert.throws(() => lstatSync(path.join(mediaTarget(), "v1", "transcript.en.vtt"))); + await assertMigrated(); +}); + +test("a media move's marker is refused, with or without a scope; nothing is touched", async () => { + for (const marker of [ + { target: mediaTarget(), direction: "out", startedAt: "", phase: "copy", scope: "media" }, + { target: path.join(PLATTER, "elsewhere", SLUG, "data"), direction: "out", startedAt: "", phase: "swap" }, + ]) { + resetCorpus(); + seedLegacy(); + writeJson(path.join(channelDir(), ".relocating.json"), marker); + const before = snapshot(); + const res = await run(); + assert.equal(res.exitCode, 1); + assert.equal(res.channels[0].outcome, "refused"); + assert.match(res.channels[0].detail ?? "", /media move .* Storage panel/); + assert.deepEqual(snapshot(), before); + } +}); + +test("a running editor refuses a real run; a dry run only notes it", async () => { + resetCorpus(); + seedLegacy(); + const before = snapshot(); + const lines: string[] = []; + const editorProbe = async () => "the editor answers at http://localhost:3001/api/pulse (HTTP 200)"; + const res = await run({ deps: { editorProbe } }, lines); + assert.equal(res.exitCode, 2); + assert.equal(res.channels.length, 0); + assert.ok(lines.some((l) => /Refusing: the editor must be stopped/.test(l))); + assert.deepEqual(snapshot(), before); + const dry = await run({ dryRun: true, deps: { editorProbe } }); + assert.equal(dry.exitCode, 0); + assert.equal(dry.channels[0].outcome, "would-migrate"); + assert.deepEqual(snapshot(), before); +}); + +test("editorRunningReason: the port, then a live job meta", async () => { + const jobsDir = path.join(ROOT, "jobs-probe"); + rmSync(jobsDir, { recursive: true, force: true }); + mkdirSync(jobsDir, { recursive: true }); + const refused = async () => { + throw Object.assign(new TypeError("fetch failed"), { cause: { code: "ECONNREFUSED" } }); + }; + const timedOut = async () => { + throw Object.assign(new Error("The operation was aborted due to timeout"), { name: "TimeoutError" }); + }; + const answers = async () => new Response("ok", { status: 200 }); + const opts = { url: "http://localhost:3999", machineBootMs: 0 }; + + assert.equal(await editorRunningReason({ jobsDir }, { ...opts, fetchImpl: refused as typeof fetch }), null); + assert.match( + (await editorRunningReason({ jobsDir }, { ...opts, fetchImpl: answers as typeof fetch })) ?? "", + /answers at http:\/\/localhost:3999\/api\/pulse \(HTTP 200\)/, + ); + assert.match( + (await editorRunningReason({ jobsDir }, { ...opts, fetchImpl: timedOut as typeof fetch })) ?? "", + /did not answer/, + ); + + const meta = { id: "j1", kind: "whisper-video", queueKey: "q", channelSlug: "c", status: "running", queuedAt: Date.now(), pid: 424242 }; + writeFileSync(path.join(jobsDir, "j1.meta.json"), JSON.stringify(meta)); + assert.match( + (await editorRunningReason({ jobsDir }, { ...opts, fetchImpl: refused as typeof fetch, isAlive: () => true })) ?? "", + /job j1 \(whisper-video, c\) is running in a live process \(pid 424242\)/, + ); + assert.equal( + await editorRunningReason({ jobsDir }, { ...opts, fetchImpl: refused as typeof fetch, isAlive: () => false }), + null, + "a dead writer's meta is a ghost", + ); + writeFileSync(path.join(jobsDir, "j1.meta.json"), JSON.stringify({ ...meta, pid: undefined })); + assert.match( + (await editorRunningReason({ jobsDir }, { ...opts, fetchImpl: refused as typeof fetch })) ?? "", + /no recorded process/, + ); + assert.equal( + await editorRunningReason({ jobsDir }, { ...opts, machineBootMs: Date.now() + 60_000, fetchImpl: refused as typeof fetch }), + null, + "a meta from before the machine booted is not read", + ); +}); + +test("killed at any step between copy and swap, a rerun resumes and finishes", async () => { + const steps: MigrationStep[] = [ + "copied", + "linked", + "marker-swap", + "platter-renamed", + "media-linked", + "data-unlinked", + "data-renamed", + "config-written", + ]; + for (const step of steps) { + resetCorpus(); + seedLegacy(); + const killed = await run({ + deps: { + checkpoint: (_slug, s) => { + if (s === step) throw new Error(`killed at ${s}`); + }, + }, + }); + assert.equal(killed.channels[0].outcome, "failed", step); + assert.equal(killed.exitCode, 1); + const marker = await readRelocationMarker(paths, SLUG); + assert.equal(marker?.scope, "tier-migration", step); + const loc = await inspectChannelMedia(paths, SLUG, undefined, { fresh: true }); + assert.equal(loc.status, "in-transition", step); + assert.equal(loc.text.readable, false, `${step}: the text is held mid-migration`); + assert.deepEqual(await legacyChannels(paths), [SLUG], `${step}: --all would resume it`); + + const resumed = await run(); + assert.equal(resumed.exitCode, 0, `${step}: ${resumed.channels[0].detail}`); + assert.equal(resumed.channels[0].outcome, "migrated", step); + assert.equal( + resumed.channels[0].resumedFrom, + ["copied", "linked"].includes(step) ? "copy" : "swap", + step, + ); + await assertMigrated(); + } +}); + +test("the free-space stop: refused before anything is written", async () => { + resetCorpus(); + seedLegacy(); + const before = snapshot(); + const settings = { ...SETTINGS, minFreeDiskGB: 5 } as MigrateOpts["settings"]; + // 6 GB free, a 5 GB floor and a 2 GB margin: the copy cannot fit. + const res = await run({ settings, deps: { freeBytes: async () => 6 * 1024 ** 3 } }); + assert.equal(res.exitCode, 1); + assert.equal(res.channels[0].outcome, "no-space"); + assert.match(res.channels[0].detail ?? "", /not enough space on the corpus disk: 6\.00 GB free .*5 GB floor.*2 GB resume margin/); + assert.deepEqual(snapshot(), before); + // 8 GB free clears it. + const ok = await run({ settings, deps: { freeBytes: async () => 8 * 1024 ** 3 } }); + assert.equal(ok.channels[0].outcome, "migrated", ok.channels[0].detail); +}); + +test("--all: smallest text first, and a stop before the big three", async () => { + const big = "omnibased"; + resetCorpus(["bigger", "small", big]); + seedLegacy("bigger", 4); + seedLegacy("small", 1); + seedLegacy(big, 1); + // A classic channel and a social one are not candidates. + writeConfig("classic"); + mkdirSync(path.join(channelDir("classic"), "data", "x1"), { recursive: true }); + + const lines: string[] = []; + const res = await run({ slugs: "all" }, lines); + assert.equal(res.exitCode, 0, lines.join("\n")); + assert.deepEqual(res.channels.map((c) => c.slug), ["small", "bigger"]); + assert.deepEqual(res.stoppedBefore, [big]); + assert.ok(lines.some((l) => l.includes("STOPPED before the big-text channels: omnibased"))); + assert.ok(lines.some((l) => l.includes("--all --include-large"))); + await assertMigrated("small"); + await assertMigrated("bigger", 4); + assert.equal(lstatSync(path.join(channelDir(big), "data")).isSymbolicLink(), true, "untouched"); + + const rest = await run({ slugs: "all", includeLarge: true }); + assert.deepEqual(rest.channels.map((c) => c.slug), [big]); + await assertMigrated(big); + assert.deepEqual(await legacyChannels(paths), []); +}); + +test("the index build after the migration rewrites no record", async () => { + resetCorpus(); + // Built while the channel's text was readable (before release 17 the index + // read a legacy channel through its link) — here, a classic layout with the + // same files and times. + writeConfig(SLUG); + writeVideos(path.join(channelDir(), "data"), SLUG); + const v1 = path.join(channelDir(), "data", "v1"); + const outcome = await normalizeLiveChat({ videoDir: v1, channelSlug: SLUG, tier: false }); + assert.equal(outcome.status, "wrote"); + utimesSync(path.join(v1, "live_chat.cues.json"), CUES_AT, CUES_AT); + const first = await buildIndex({ paths, onLog: () => {} }); + assert.equal(first.added, 3); + const subsBefore = statSync(path.join(v1, "transcript.live_chat.json")).mtimeMs; + + // Onto the retired layout, as the old mover left it (rename keeps mtimes). + mkdirSync(path.dirname(oldTree()), { recursive: true }); + renameSync(path.join(channelDir(), "data"), oldTree()); + symlinkSync(oldTree(), path.join(channelDir(), "data")); + writeConfig(SLUG, { dataDir: oldTree() }); + + const res = await run(); + assert.equal(res.channels[0].outcome, "migrated", res.channels[0].detail); + assert.equal(lstatSync(path.join(v1, "transcript.live_chat.json")).mtimeMs, subsBefore); + assert.equal((await isLiveChatCuesFresh(v1)).fresh, true, "the cues are not stale"); + + const second = await buildIndex({ paths, onLog: () => {} }); + assert.deepEqual(second.heldChannels, []); + assert.equal(second.added, 0); + assert.equal(second.changed, 0, "no record rewritten: every mtime the index keys on is unchanged"); + assert.equal(second.removed, 0); +}); diff --git a/common/bin/migrate-media-tier.ts b/common/bin/migrate-media-tier.ts @@ -0,0 +1,1183 @@ +#!/usr/bin/env tsx +// THE ONE-OFF MEDIA-TIER MIGRATION (release 17, slice T3). +// +// archilyzer storage migrate-tier <slug>… [--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/<slug>/data` an absolute symlink to `<root>/<slug>/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/<slug>/data/ a REAL directory on the corpus disk again +// <id>/audio.mp3 -> ../../media/<id>/audio.mp3 (relative, the file's mtime) +// channels/<slug>/media -> <root>/<slug>/media (ONE absolute link) +// <root>/<slug>/media/<id>/… (the old `<root>/<slug>/data`, renamed) +// config.json: mediaDir = <root>/<slug>/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: <root>/<slug>/media, direction: "out", +// phase: "copy", scope: "tier-migration"}`; every entry that is +// NOT tierable (the text, `clips/`, scratch, `source-media.*`) +// is copied to `channels/<slug>/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(<root>/<slug>/data -> <root>/<slug>/media)`; +// `symlink(<root>/<slug>/media, channels/<slug>/media)`; +// `unlink(channels/<slug>/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. +// reclaim only with --reclaim: the platter's copies of the text (every +// entry of `<root>/<slug>/media/<id>/` that is not tierable) are +// deleted, empty video dirs dropped. +// 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 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/<id>/` (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<CopyKind, Tally>; + +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; +}; + +// 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<Tally> { + 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; +} + +// 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<Inventory> { + const inv: Inventory = { + videoIds: [], + copyList: [], + copy: emptyTallies(), + copyBytes: 0, + tierable: [], + tierableBytes: 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; + } + 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<KindTallies> { + 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 <ARCHILYZER_EDITOR_URL>/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<Paths, "jobsDir">, + opts: EditorCheckOptions = {}, +): Promise<string | null> { + 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<string | null>; + freeBytes?: (dir: string) => Promise<number>; + // Test seam: called after each named step; a throw is a kill there. + checkpoint?: (slug: string, step: MigrationStep) => void | Promise<void>; +}; + +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; + linksMade?: number; + linksAlready?: number; + reclaimedFiles?: number; + reclaimedBytes?: number; + // 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<number>; + checkpoint: (slug: string, step: MigrationStep) => Promise<void>; + // 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<Plan> { + 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 <root>/${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 <root>/${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<boolean> { + try { + return (await stat(p)).isDirectory(); + } catch { + return false; + } +} + +async function writeMarker(ctx: Ctx, slug: string, marker: RelocationMarker): Promise<void> { + 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<ChannelReport> { + 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; + 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)}`, + ); + + // 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 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 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<void> { + const { log } = ctx; + // 1. The platter: <root>/<slug>/data -> <root>/<slug>/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/<slug>/media -> <root>/<slug>/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. Everything in +// `<root>/<slug>/media/<id>/` that is not a tierable regular file was copied +// to the corpus disk and verified; it goes, and a video dir left empty goes +// with it. A top-level non-directory entry was a stray text file, copied too. +async function reclaimTree( + ctx: Ctx, + slug: string, + mediaDir: string, +): Promise<Pick<ChannelReport, "reclaimedFiles" | "reclaimedBytes">> { + let files = 0; + let bytes = 0; + const top = await readdir(mediaDir, { withFileTypes: true }); + for (const e of top) { + const p = path.join(mediaDir, e.name); + if (!e.isDirectory()) { + const t = await measureEntry(p); + files += t.files; + bytes += t.bytes; + if (!ctx.dryRun) await rm(p, { force: true }); + 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 t = await measureEntry(fp, st); + files += t.files; + bytes += t.bytes; + if (!ctx.dryRun) await rm(fp, { recursive: true, force: true }); + } + if (!ctx.dryRun) await rmdir(p).catch(() => {}); + } + ctx.log( + `${slug}: ${ctx.dryRun ? "would reclaim" : "reclaimed"} ${files} file(s), ` + + `${formatBytes(bytes)} of text copies from ${mediaDir}`, + ); + return { reclaimedFiles: files, reclaimedBytes: bytes }; +} + +// --reclaim on a channel already on the media tier (this run's earlier +// migration, or a previous one): its media drive must answer, and nothing may +// be moving it. +async function reclaimChannel( + ctx: Ctx, + slug: string, + mediaDir: string, +): Promise<Pick<ChannelReport, "reclaimedFiles" | "reclaimedBytes">> { + 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 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<string[]> { + 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 already on the media tier (for `--all --reclaim`). +async function migratedChannels(paths: Paths): Promise<string[]> { + const entries = await readdir(paths.channelsDir, { withFileTypes: true }).catch(() => []); + const out: string[] = []; + for (const e of entries) { + if (!e.isDirectory()) continue; + const config = await readChannelConfig(paths, e.name).catch(() => null); + if (config?.mediaDir?.trim() && !config.dataDir?.trim()) 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<number> { + 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<MigrateResult> { + 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[] = []; + 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) sized.push({ slug, bytes: await copyBytesOf(opts.paths, slug) }); + 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); + if (opts.includeLarge) order.push(...(await measure(large))); + else deferred = large; + if (ctx.reclaim) { + for (const slug of await migratedChannels(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`, + ); + } + } + 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 = 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 <slug>\` one at a time, or ` + + `\`archilyzer storage migrate-tier --all --include-large\`.`, + ); + } + return result; +} + +// --------------------------------------------------------------------------- +// The command line +// --------------------------------------------------------------------------- + +const USAGE = + "usage: archilyzer storage migrate-tier <slug>… | --all [--order smallest] " + + "[--include-large] [--dry-run] [--reclaim]"; + +export async function main(argv: string[]): Promise<number> { + 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)));