Archilyzer · Source

archilyzer

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

commit de5b77d53c3e70b344864cb33ffbe84e01256605
parent 1aac49840c13bbb31e960ac20198be30af7067de
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Fri,  2 Oct 2026 00:19:31 -0400

Merge r17/media-tier-model (release 17 slice T1) — the media tier's classifier, hook and guards: files classified by name (media/text/scratch), every media finalisation tiers audio and the raw live chat into channels/<slug>/media through relative links and every deleter follows the link; config mediaDir, the retired whole-data/ layout reads legacy and is held; a text guard for the index, stats, snapshot, normalize and shards so the media tier never holds them; lanes by tier; channelWriters mediaOnly; the cues twin; reviewed SHIP

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

Diffstat:
MCHANNEL.md | 7++++---
Mcommon/bin/doctor.test.ts | 15++++++++++++---
Mcommon/bin/run-operation.test.ts | 7++++---
Mcommon/controller/archiveLiveChat.ts | 2++
Mcommon/controller/autoRunner.test.ts | 71+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Mcommon/controller/autoRunner.ts | 33++++++++++++++++++++++++++++++---
Mcommon/controller/backfillReacquire.ts | 6++++--
Mcommon/controller/buildIndex.test.ts | 216+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------
Mcommon/controller/buildIndex.ts | 109++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------------
Mcommon/controller/buildStats.test.ts | 68++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------
Mcommon/controller/buildStats.ts | 11++++++-----
Mcommon/controller/channelSnapshot.test.ts | 59+++++++++++++++++++++++++++++++++++++++++++++++++++--------
Mcommon/controller/channelSnapshot.ts | 149++++++++++++++++++++++++++++++++++++++++++++++++++++---------------------------
Mcommon/controller/channelWriters.test.ts | 19+++++++++++++++++++
Mcommon/controller/channelWriters.ts | 9++++++++-
Mcommon/controller/channels.ts | 14+++++++++-----
Mcommon/controller/cleanAudioFromTranscribed.ts | 6++++--
Mcommon/controller/cleanExtraAudioFormats.ts | 3++-
Mcommon/controller/evictClipWindows.test.ts | 53++++++++++++++++++++++++++++++++++++++++++++++-------
Mcommon/controller/evictClipWindows.ts | 30+++++++++++++++---------------
Acommon/controller/mediaTierHooks.test.ts | 145+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/controller/normalizeAll.ts | 6++++--
Mcommon/controller/normalizeAllLiveChat.ts | 13+++++++++++++
Mcommon/controller/normalizeLiveChat.ts | 30+++++++++++++++++++++++++++---
Mcommon/controller/operationBatch.ts | 12+++++++++++-
Mcommon/controller/purgeSupersededAutoSubs.ts | 5+++--
Mcommon/controller/recencyIndex.ts | 11++++++++---
Mcommon/controller/relocateChannelMedia.test.ts | 32+++++++++++++++++++-------------
Mcommon/controller/removeWrongFormatAudio.ts | 5+++--
Mcommon/controller/renameChannel.test.ts | 8+++++++-
Mcommon/controller/storageLocations.test.ts | 20+++++++++++++-------
Mcommon/controller/storageStall.test.ts | 211+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------------
Mcommon/controller/storageWatch.test.ts | 20+++++++++++++-------
Mcommon/controller/transcode.ts | 8++++++++
Mcommon/jobs/jobKinds.test.ts | 77++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcommon/jobs/jobKinds.ts | 120+++++++++++++++++++++++++++++++++++++++++++++++++++----------------------------
Mcommon/jobs/streamCommand.test.ts | 64+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcommon/jobs/streamCommand.ts | 20++++++++++++++++++--
Mcommon/lib/channelConfig.ts | 6+++++-
Mcommon/lib/channelConfigSchema.test.ts | 4+++-
Mcommon/lib/channelConfigSchema.ts | 1+
Mcommon/lib/channelMedia.test.ts | 259+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------
Mcommon/lib/channelMedia.ts | 450+++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------------------
Mcommon/lib/channelMediaHold.ts | 40++++++++++++++++++++++++++++------------
Mcommon/lib/fileSchemaDocs.ts | 4++--
Acommon/lib/mediaTier-server.test.ts | 340+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/lib/mediaTier-server.ts | 423+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/lib/mediaTier.test.ts | 88+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/lib/mediaTier.ts | 114+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/lib/storageHealth.test.ts | 8++++++++
Mcommon/lib/storageHealth.ts | 10++++++----
Mcommon/ytdlp/audioCheckedDownload.ts | 3++-
Mcommon/ytdlp/downloadOneManaged.ts | 13++++++++++++-
Mcommon/ytdlp/runYtdlp.ts | 39+++++++++++++++++++++++++++++++++++++++
Meditor/CHANGELOG.md | 1+
Meditor/app/channels/[slug]/lib/fixIncompleteTranscript.ts | 8+++++---
Meditor/app/channels/[slug]/shardActions.ts | 8++++++--
Meditor/app/channels/[slug]/videos/[id]/videoActions.ts | 42+++++++++++++++++++++++++++++++++++++++---
Meditor/app/components/MediaLocationBadge.tsx | 4++++
Mplans/FACTS.md | 50++++++++++++++++++++++++++++++++++++++++++++++++++
Mplans/release-17.md | 229+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mumtool/report-to-video/cues.mjs | 124+++++++++++++++++++++++++++++++++++++------------------------------------------
Mumtool/report-to-video/cues.test.mjs | 138++++++++++++++++++++++++++++++++++++++++++++++---------------------------------
63 files changed, 3463 insertions(+), 637 deletions(-)

diff --git a/CHANNEL.md b/CHANNEL.md @@ -6,9 +6,9 @@ One channel of the corpus, persisted to `transcripts/channels/<slug>/config.json `handling` is the one required key: a file without a valid one is not a channel. The smallest channel is `{ "handling": "youtube", "url": "https://www.youtube.com/@example" }`. -Every other key is optional and has NO default of its own: an absent key means whatever its description says — for the per-channel overrides, inherit the global setting of the same name; for `name`, `url`, `dataDir`, `subLangs` and the sync-state stamps, simply unset. So an ill-typed or out-of-range value is not coerced — it is DROPPED, as if the file did not spell it. Unknown keys (including the retired `excludeFromSync`, now a paused `sync` tier in the channel-priority document) are dropped by every read and every write. +Every other key is optional and has NO default of its own: an absent key means whatever its description says — for the per-channel overrides, inherit the global setting of the same name; for `name`, `url`, `mediaDir`, `subLangs` and the sync-state stamps, simply unset. So an ill-typed or out-of-range value is not coerced — it is DROPPED, as if the file did not spell it. Unknown keys (including the retired `excludeFromSync`, now a paused `sync` tier in the channel-priority document) are dropped by every read and every write. -The three **sync state** keys are not configuration: the sync, sweep and download passes stamp them, the channel form never does, and they live in the same file on purpose. Writers after creation PATCH (`patchChannelConfig`): each re-reads the file at the moment it writes and changes only its own keys, so a stamp and a form save made at once in the editor both land. The two exceptions write a whole config, and only when there is no readable file to patch: a media move and a channel rename record `dataDir` from their own copy of the config. +The three **sync state** keys are not configuration: the sync, sweep and download passes stamp them, the channel form never does, and they live in the same file on purpose. Writers after creation PATCH (`patchChannelConfig`): each re-reads the file at the moment it writes and changes only its own keys, so a stamp and a form save made at once in the editor both land. The two exceptions write a whole config, and only when there is no readable file to patch: a media move and a channel rename record `mediaDir` from their own copy of the config. Regenerate this file with `pnpm --filter yt-dlp-transcript-common exec tsx bin/file-schemas-docs.ts`. @@ -27,7 +27,8 @@ Regenerate this file with `pnpm --filter yt-dlp-transcript-common exec tsx bin/f | `keepLatest` | config | Keep-latest window: the newest N videos (by upload date) are protected from the Clean-audio sweep AND have their source video persisted to the saved-video store. 0 or absent = disabled; positives clamp to [1, 100000]. A kept video later found deleted at the source is pinned permanently via the do-not-clean marker. | | `extractionMode` | config | `"ytdlp"` (default — yt-dlp's own `-x --audio-format` postprocessor, no source container kept) or `"app"` (yt-dlp downloads the source container and the app runs ffmpeg). The keep-latest persistence rule forces `"app"` for the videos it persists. | | `savedVideosDir` | config | Per-channel override for the saved-video store root: this channel's persisted source videos live under `<savedVideosDir>/<slug>/<videoId>/`. Trimmed; blank = the global store. | -| `dataDir` | config | Where this channel's media ACTUALLY lives when relocated to another drive: the absolute path `channels/<slug>/data` is a symlink to. Absent = in place. Written ONLY by the relocate / re-point jobs on success — a record of what is on disk, never free text, because a value that disagrees with the link is an "inconsistent" channel every guard refuses. | +| `dataDir` | config | RETIRED (release 17). The whole-directory layout's record: the absolute path `channels/<slug>/data` was a symlink to. Still parsed for one release so a write never erases it: a channel that carries it — or whose `data/` is a link — is `legacy`, and every media job, lane and build holds it until `archilyzer storage migrate-tier <slug>` moves its text back and its media into `mediaDir`. Never written by anything but that migration, which removes it. | +| `mediaDir` | config | Where this channel's big files live when relocated: `channels/<slug>/media` is a symlink to it, `<root>/<slug>/media`. Absent = in place. Written only by relocate / re-point / the tier migration — a record of what is on disk, never free text, because a value that disagrees with the link is an "inconsistent" channel every media guard refuses. The text (`data/`) never moves. | | `ytdlpExtraArgs` | config | Extra yt-dlp arguments, appended verbatim. Must be an array of strings or it is dropped. | | `subLangs` | config | yt-dlp `--sub-langs` value for caption downloads. | | `lastSyncedAt` | sync state | SYNC STATE. When the channel last synced (ISO time). Stamped by every sync, and by a social fetch; read by the scheduler's cadence gate. | diff --git a/common/bin/doctor.test.ts b/common/bin/doctor.test.ts @@ -237,9 +237,16 @@ test("an unmounted drive is a warning, not an empty channel and not a failure", const c = checkout(); const chan = path.join(c.paths.channelsDir, "moved"); mkdirSync(chan, { recursive: true }); - const target = path.join(c.root, "elsewhere", "moved", "data"); - writeFileSync(path.join(chan, "config.json"), JSON.stringify({ dataDir: target })); - symlinkSync(target, path.join(chan, "data")); // the target does not exist + const target = path.join(c.root, "elsewhere", "moved", "media"); + writeFileSync(path.join(chan, "config.json"), JSON.stringify({ mediaDir: target })); + mkdirSync(path.join(chan, "data")); + symlinkSync(target, path.join(chan, "media")); // the target does not exist + // A channel on the retired whole-directory layout is named with its way out. + const old = path.join(c.paths.channelsDir, "retired"); + mkdirSync(old, { recursive: true }); + const oldTarget = path.join(c.root, "elsewhere", "retired", "data"); + writeFileSync(path.join(old, "config.json"), JSON.stringify({ dataDir: oldTarget })); + symlinkSync(oldTarget, path.join(old, "data")); // No workers: a `{}` file would synthesize whisper workers, whose engine // this machine does not have — a real failure, and not this test's. writeFileSync(c.paths.settingsFile, JSON.stringify({ workers: [] })); @@ -250,6 +257,8 @@ test("an unmounted drive is a warning, not an empty channel and not a failure", const media = r.checks.find((x) => x.id === "media")!; assert.equal(media.status, "warn"); assert.match(media.detail, /moved: unreachable/); + const retired = r.checks.filter((x) => x.id === "media").map((x) => x.detail).join("\n"); + assert.match(retired, /retired: legacy — .*archilyzer storage migrate-tier retired/); assert.equal(r.ok, true, renderDoctorReport(r)); assert.deepEqual(tree(c.root), before); }); diff --git a/common/bin/run-operation.test.ts b/common/bin/run-operation.test.ts @@ -140,12 +140,13 @@ test("the media guard refuses an unmounted channel before any job exists", { tim const slug = "moved"; const channelDir = path.join(getPaths().channelsDir, slug); await mkdir(channelDir, { recursive: true }); - const target = path.join(ROOT, "platter", slug, "data"); // never created + const target = path.join(ROOT, "platter", slug, "media"); // never created await writeFile( path.join(channelDir, "config.json"), - JSON.stringify({ handling: "transcribe", url: "https://example.com/m", dataDir: target }), + JSON.stringify({ handling: "transcribe", url: "https://example.com/m", mediaDir: target }), ); - await symlink(target, path.join(channelDir, "data")); + await mkdir(path.join(channelDir, "data")); + await symlink(target, path.join(channelDir, "media")); const before = jobRecords(); const { o, out } = capture(); const code = await runOperation({ operation: "diarization", channel: slug, ids: [] }, { out }); diff --git a/common/controller/archiveLiveChat.ts b/common/controller/archiveLiveChat.ts @@ -161,6 +161,8 @@ async function stageChannel( videoDir, channelSlug: ch.slug, configName: ch.config.name, + // A build moves no media (release 17): the jobs tier, this reads. + tier: false, }); } catch (err) { normalizeFailed++; diff --git a/common/controller/autoRunner.test.ts b/common/controller/autoRunner.test.ts @@ -1,6 +1,6 @@ import { test } from "node:test"; import assert from "node:assert/strict"; -import { mkdirSync, mkdtempSync, rmSync, writeFileSync } from "node:fs"; +import { mkdirSync, mkdtempSync, rmSync, symlinkSync, writeFileSync } from "node:fs"; import { tmpdir } from "node:os"; import path from "node:path"; import { @@ -12,6 +12,7 @@ import { priorityContextFor, resetPriorityContextForTest, sharedAutoQueueState, + isChannelHeldForLane, } from "./autoRunner"; import { defaultChannelPriority, @@ -33,7 +34,7 @@ import type { } from "../jobs/autoQueuePolicy"; import type { Paths } from "../lib/paths"; import type { SiteSettings } from "../lib/settings"; -import { forgetChannelMedia } from "../lib/channelMedia"; +import { forgetChannelMedia, inspectChannelMedia } from "../lib/channelMedia"; import { bucketLaneOperationId } from "../lib/operations"; // THE RUNNER'S HALF OF S1 (plans/channel-priority.md), and it is here rather @@ -711,6 +712,9 @@ test("a channel with a move marker projects no work on any lane, and the hold li ); forgetChannelMedia(slug); for (const lane of LANES) { + // (Release 17: the digest lane is held only by the text hold, and a media + // move leaves the text readable — `isChannelHeldForLane` below. This + // fixture has no digest work, so its total is 0 either way.) assert.equal( await laneTotal(lane, paths, configs), 0, @@ -723,3 +727,66 @@ test("a channel with a move marker projects no work on any lane, and the hold li forgetChannelMedia(slug); assert.ok((await laneTotal("transcription", paths, configs)) > 0); }); + + +// --- Release 17: the lanes by tier ------------------------------------------ +// +// A stalled, unmounted or moving MEDIA tier holds the lanes that open the big +// files (transcription, download, backfill) and lets the digest lane run: a +// digest reads and writes text, which stays on the corpus disk. Only the text +// hold — a `legacy` channel, a tier migration in flight — keeps the digest lane +// off a channel. Decided from the inspector's real answers over fixtures. + +test("release 17: an unmounted media drive holds transcription/download/backfill, lets digest run; legacy holds all four", async () => { + const paths = pendingFixture(); + const slug = "tiered"; + const channelDir = path.join(paths.channelsDir, slug); + mkdirSync(path.join(channelDir, "data", "v1"), { recursive: true }); + const target = path.join(path.dirname(paths.channelsDir), "platter", slug, "media"); + const writeCfg = (extra: Record<string, unknown>) => + writeFileSync( + path.join(channelDir, "config.json"), + JSON.stringify({ url: "https://www.youtube.com/@t", handling: "transcribe", ...extra }), + ); + writeCfg({ mediaDir: target }); + symlinkSync(target, path.join(channelDir, "media")); // dangling: not mounted + const unmounted = await inspectChannelMedia(paths, slug, undefined, { fresh: true }); + assert.equal(unmounted.status, "unreachable"); + const heldBy = (loc: Parameters<typeof isChannelHeldForLane>[1]) => + Object.fromEntries(LANES.map((lane) => [lane, isChannelHeldForLane(lane, loc)])); + assert.deepEqual(heldBy(unmounted), { + transcription: true, + download: true, + digest: false, + backfill: true, + }); + + // A media move in flight: the same. + writeFileSync( + path.join(channelDir, ".relocating.json"), + JSON.stringify({ target, direction: "out", startedAt: "", phase: "copy", scope: "media" }), + ); + const moving = await inspectChannelMedia(paths, slug, undefined, { fresh: true }); + assert.equal(moving.status, "in-transition"); + assert.equal(heldBy(moving).digest, false); + assert.equal(heldBy(moving).transcription, true); + // A tier migration rebuilds data/ itself: the digest lane is held too. + writeFileSync( + path.join(channelDir, ".relocating.json"), + JSON.stringify({ target, direction: "out", startedAt: "", phase: "copy", scope: "tier-migration" }), + ); + const migrating = await inspectChannelMedia(paths, slug, undefined, { fresh: true }); + assert.equal(heldBy(migrating).digest, true); + rmSync(path.join(channelDir, ".relocating.json")); + + // The retired layout holds every lane. + writeCfg({ dataDir: path.join(path.dirname(paths.channelsDir), "platter", slug, "data") }); + const legacy = await inspectChannelMedia(paths, slug, undefined, { fresh: true }); + assert.equal(legacy.status, "legacy"); + assert.deepEqual(heldBy(legacy), { + transcription: true, + download: true, + digest: true, + backfill: true, + }); +}); diff --git a/common/controller/autoRunner.ts b/common/controller/autoRunner.ts @@ -86,10 +86,12 @@ import { downloadQueueKey } from "../lib/queueKeys"; import { isGateHeld } from "../lib/pauseGates"; import { inspectChannelMedia, + markerHoldsText, readRelocationMarker, + type ChannelMediaLocation, type ChannelMediaStatus, } from "../lib/channelMedia"; -import { isMediaHeld, mediaHoldText } from "../lib/channelMediaHold"; +import { isMediaHeld, isTextHeld, mediaHoldText } from "../lib/channelMediaHold"; import { type ChannelPriority, type FocusSummary, @@ -557,6 +559,21 @@ export function makeFocusHoldReporter( // too — an operator watching the log sees the channel leave and come back. const mediaSkipLogged = new Map<string, ChannelMediaStatus>(); +// WHETHER A LANE SKIPS A CHANNEL, from its inspect answer (release 17). The +// digest lane writes only text, so only the text hold — a `legacy` channel, a +// `data/` it cannot read, a tier migration in flight — keeps it off a channel; +// it runs while the media is moving, stalled or unmounted. The transcription, +// download and backfill lanes open the big files and keep the media hold. +export function isChannelHeldForLane( + kind: AutoQueueKind, + location: Pick<ChannelMediaLocation, "status" | "text">, +): boolean { + if (kind === "digest") { + return isTextHeld(location.status) || !location.text.readable; + } + return isMediaHeld(location.status); +} + function noteSkippedForMedia( slug: string, status: ChannelMediaStatus, @@ -637,8 +654,14 @@ async function buildChannelWork( // drive, a link and a config that disagree. It lifts by itself — the next // tick after the marker is removed (a move completed or abandoned; the // movers forget the memo) or the drive is back reads the channel again. + // + // THE DIGEST LANE READS TEXT (release 17): it is held only by the text + // hold — a `legacy` channel, or a `data/` it cannot read — so a digest + // runs on a channel whose media is moving, stalled or unmounted. The + // transcription, download and backfill lanes open the big files and keep + // the media hold. const location = media[i]; - if (location && isMediaHeld(location.status)) { + if (location && isChannelHeldForLane(kind, location)) { noteSkippedForMedia(slug, location.status, location.detail); continue; } @@ -1894,8 +1917,12 @@ async function runLoop( // lane it retires nothing — runOperationPick owns the (operation, video) // keys and never ran — which is equally fine, because guard 2 has taken // the channel off the list before the next tick could offer it again. + // + // NOT THE DIGEST LANE (release 17): a media move leaves the text where it + // is, and a digest writes only text. A tier migration rebuilds `data/` + // itself, and holds it too. const marker = await readRelocationMarker(paths, channelSlug); - if (marker) { + if (marker && (kind !== "digest" || markerHoldsText(marker))) { onLog( `Auto-${kind}: skipping ${channelSlug}/${pick.videoId}, ` + `${mediaHoldText("in-transition")} — a relocation ` + diff --git a/common/controller/backfillReacquire.ts b/common/controller/backfillReacquire.ts @@ -59,8 +59,9 @@ // deletes audio after a transcription. The outcome is the good one: a video // whose only transcript was YouTube ASR gets our own. +import { removeMediaFile } from "../lib/mediaTier-server"; import path from "node:path"; -import { readdir, rm } from "node:fs/promises"; +import { readdir } from "node:fs/promises"; import type { Paths } from "../lib/paths"; import { getSettings } from "../lib/settings"; import { diskGate, evaluateDiskGate, getFreeBytes } from "../lib/diskSpace"; @@ -503,7 +504,8 @@ function buildCleanup( } for (const name of added) { - await rm(path.join(videoDir, name), { force: true }).catch((err) => { + // Through its link when the hook tiered it (release 17). + await removeMediaFile(videoDir, name).catch((err) => { // Reported, never thrown: this runs in a `finally`, and throwing here // would replace the real outcome of the item with a cleanup error. log(`Could not remove ${name} for ${videoId}: ${(err as Error).message}`); diff --git a/common/controller/buildIndex.test.ts b/common/controller/buildIndex.test.ts @@ -1,15 +1,23 @@ // Integration: the index build's HOLD, through the REAL buildIndex, over a temp // corpus. // -// A channel's `data/` may be an absolute symlink to another drive (AGENTS.md, -// "A channel's `data/` may live on another drive"). With that drive unmounted -// the link dangles, and the index build used to read the channel as having no -// videos: it removed every record the channel had, and the site built next -// published the channel as gone. These cases pin the hold that replaced it: -// the channel is not rescanned, and its records and shared pages are kept as -// they are; a FULL rebuild with a channel held refuses unless -// ARCHILYZER_INDEX_ALLOW_HELD is set; and all of it undoes itself when the -// drive is back. The stats build's twin is buildStats.test.ts case (i). +// A channel's `data/` used to be an absolute symlink to another drive (the +// retired whole-directory layout). With that drive unmounted the link +// dangled, and the index build used to read the channel as having no videos: +// it removed every record the channel had, and the site built next published +// the channel as gone. These cases pin the hold that replaced it: the channel +// is not rescanned, and its records and shared pages are kept as they are; a +// FULL rebuild with a channel held refuses unless ARCHILYZER_INDEX_ALLOW_HELD +// is set; and all of it undoes itself when the text is readable again. The +// stats build's twin is buildStats.test.ts case (i). +// +// RELEASE 17: only the TEXT holds the index. A relocated channel's text stays +// on the corpus disk and only its media (`channels/<slug>/media`) is on the +// drive, so an unmounted MEDIA drive holds nothing here (case (j)). What holds +// is a text tier that cannot be read — here, the channel put back on the +// retired layout with its drive away (`unmount`), which reads `legacy` — and +// what lifts it is the text home again (`remount`, which is what +// `archilyzer storage migrate-tier` does). // // Run with: node_modules/.bin/tsx --test common/controller/buildIndex.test.ts @@ -78,12 +86,19 @@ const COMMON = fileURLToPath(new URL("..", import.meta.url)); // holds. The last case reads it. const writes: string[] = []; let afterStat: ((p: string) => void) | null = null; +// A STALLED MEDIA DRIVE (release 17): while set, a promise-API stat, readFile +// or open whose path is a symlink resolving under it never settles — what a +// call blocked on a stalled drive looks like from here — and is recorded. +let hangUnder: string | null = null; +const driveCalls: string[] = []; { const req = createRequire(import.meta.url); const fsCjs = req("node:fs") as Record<string, unknown>; const fspCjs = req("node:fs/promises") as Record<string, unknown>; const WRITES = ["writeFile", "appendFile", "rename", "mkdir", "rm", "rmdir", "unlink", "copyFile", "cp", "symlink", "link", "utimes", "truncate", "mkdtemp", "chmod"]; - const TWO_PATHS = new Set(["rename", "copyFile", "cp", "symlink", "link"]); + // A symlink's first argument is its TARGET (relative to the link, and never + // written); only the link itself, the second, is. + const TWO_PATHS = new Set(["rename", "copyFile", "cp", "link"]); const opensForWrite = (flags: unknown) => (typeof flags === "string" && /[wa+]/.test(flags)) || (typeof flags === "number" && (flags & 3) !== 0); @@ -94,7 +109,8 @@ let afterStat: ((p: string) => void) | null = null; if (typeof fn !== "function") return; mod[name] = function (this: unknown, ...args: unknown[]) { if (mode === "write" || opensForWrite(args[1])) { - const ps = TWO_PATHS.has(name.replace(/Sync$/, "")) ? [args[0], args[1]] : [args[0]]; + const base = name.replace(/Sync$/, ""); + const ps = base === "symlink" ? [args[1]] : TWO_PATHS.has(base) ? [args[0], args[1]] : [args[0]]; for (const a of ps) { const p = asPath(a); if (p !== null) writes.push(path.resolve(p)); @@ -121,6 +137,31 @@ let afterStat: ((p: string) => void) | null = null; if (p !== null) afterStat?.(path.resolve(p)); return result; }; + const throughLinkUnder = (p: string, root: string): boolean => { + try { + if (!fsCjsLstat(p).isSymbolicLink()) return false; + // Through every link on the way (`media` is itself a link to the drive). + return fsCjsRealpath(p).startsWith(root + path.sep); + } catch { + return false; + } + }; + const fsCjsLstat = (fsCjs.lstatSync as (p: string) => { isSymbolicLink(): boolean }).bind(fsCjs); + const fsCjsRealpath = (fsCjs.realpathSync as (p: string) => string).bind(fsCjs); + for (const n of ["stat", "readFile", "open"]) { + const orig = fspCjs[n] as (...a: unknown[]) => Promise<unknown>; + fspCjs[n] = function (this: unknown, ...args: unknown[]) { + const p = asPath(args[0]); + if (hangUnder && p !== null) { + const abs = path.resolve(p); + if (abs.startsWith(hangUnder + path.sep) || throughLinkUnder(abs, hangUnder)) { + driveCalls.push(`${n} ${abs}`); + return new Promise(() => {}); + } + } + return orig.apply(this, args); + }; + } syncBuiltinESMExports(); } @@ -200,21 +241,40 @@ function resetCorpus(channels: string[] = [CHANNEL, DRIVE_CHANNEL]): void { }); } -// A channel whose data/ is a relocated symlink, the way the editor's Storage -// panel leaves it: channels/<slug>/data -> <MEDIA>/<slug>/data, with -// config.dataDir recording the target. d1 carries a subtitle track. +// A channel whose MEDIA is relocated (release 17): its text in a real data/ +// on the corpus disk, channels/<slug>/media -> <MEDIA>/<slug>/media, with +// config.mediaDir recording the target. d1 carries a subtitle track. +const DRIVE_MEDIA = () => path.join(MEDIA, DRIVE_CHANNEL, "media"); +const driveDataLink = () => path.join(paths.channelsDir, DRIVE_CHANNEL, "data"); +const AWAY_TEXT = () => path.join(AWAY, DRIVE_CHANNEL, "data"); function seedDriveChannel(titles: Record<string, string> = {}): void { - const target = path.join(MEDIA, DRIVE_CHANNEL, "data"); - mkdirSync(target, { recursive: true }); - writeChannel(DRIVE_CHANNEL, { dataDir: target }); - symlinkSync(target, path.join(paths.channelsDir, DRIVE_CHANNEL, "data")); + mkdirSync(DRIVE_MEDIA(), { recursive: true }); + writeChannel(DRIVE_CHANNEL, { mediaDir: DRIVE_MEDIA() }); + symlinkSync(DRIVE_MEDIA(), path.join(paths.channelsDir, DRIVE_CHANNEL, "media")); seedVideo("d1", DRIVE_CHANNEL, { subs: true, title: titles.d1 }); seedVideo("d2", DRIVE_CHANNEL, { title: titles.d2 }); } -// Unmount: the link now dangles, exactly as an absent USB drive leaves it. -const unmount = () => renameSync(MEDIA, AWAY); -const remount = () => renameSync(AWAY, MEDIA); +// The text goes away: the channel is on the retired whole-directory layout +// (`data` an absolute link to <MEDIA>/<slug>/data, `dataDir` recorded) with its +// drive not mounted — the text moved to AWAY, the link dangling. `rename` +// keeps every mtime, as `rsync -a` would. +const unmount = () => { + mkdirSync(path.dirname(AWAY_TEXT()), { recursive: true }); + renameSync(driveDataLink(), AWAY_TEXT()); + const target = path.join(MEDIA, DRIVE_CHANNEL, "data"); + symlinkSync(target, driveDataLink()); + writeChannel(DRIVE_CHANNEL, { dataDir: target }); +}; +// The text home again, in a real data/ (what `migrate-tier` leaves). +const remount = () => { + rmSync(driveDataLink()); + renameSync(AWAY_TEXT(), driveDataLink()); + writeChannel(DRIVE_CHANNEL, { mediaDir: DRIVE_MEDIA() }); +}; +// The MEDIA drive alone goes away and comes back. +const unmountMedia = () => renameSync(MEDIA, `${MEDIA}-unmounted`); +const remountMedia = () => renameSync(`${MEDIA}-unmounted`, MEDIA); async function runIndex(log: string[] = []) { const res = await buildIndex({ paths, onLog: (s) => log.push(s) }); @@ -327,7 +387,7 @@ test("(a) an unmounted drive: the channel's records, transcript pages and subs s assert.ok(line, log.join("\n")); assert.match( line, - /its media is not reachable \(drive not mounted\?\), on location "USB drive"; its 2 indexed video\(s\) are kept as they are, not rescanned/, + /its media layout is the retired whole-directory one \(run archilyzer storage migrate-tier\), on location "USB drive"; its 2 indexed video\(s\) are kept as they are, not rescanned/, ); assert.ok(!line.includes(ROOT), line); assert.ok( @@ -387,7 +447,7 @@ test("(c) a full rebuild with a channel held refuses without the override, and h err.message, new RegExp( `must be rebuilt in full \\(index schema ${current - 1} -> ${current}\\), but 1 channel\\(s\\) cannot be read: ` + - `drive-channel \\(its media is not reachable \\(drive not mounted\\?\\), on location "USB drive"\\)`, + `drive-channel \\(its media layout is the retired whole-directory one \\(run archilyzer storage migrate-tier\\), on location "USB drive"\\)`, ), ); // Why, the ways out (mounting first), and the override by name; no path. @@ -438,7 +498,7 @@ test("(c) a full rebuild with a channel held refuses without the override, and h assert.ok( log.some( (l) => - l.startsWith(`Channel ${DRIVE_CHANNEL}: its media is not reachable`) && + l.startsWith(`Channel ${DRIVE_CHANNEL}: its media layout is the retired`) && l.includes(`held under ${INDEX_ALLOW_HELD_ENV}: this full rebuild cleared its index records`), ), log.join("\n"), @@ -629,12 +689,118 @@ test("(i) the drive lost MID-WALK: the second look holds the channel instead of assert.deepEqual(sharedSubs(), subsBefore); assert.ok( log.some((l) => - l.startsWith(`Channel ${DRIVE_CHANNEL}: its media is not reachable (drive not mounted?), on location "USB drive"; its 4 indexed video(s) are kept`), + l.startsWith(`Channel ${DRIVE_CHANNEL}: its media layout is the retired whole-directory one (run archilyzer storage migrate-tier), on location "USB drive"; its 4 indexed video(s) are kept`), ), log.join("\n"), ); }); +test("(j) release 17: an unmounted MEDIA drive does not hold the index — the text is read", async () => { + resetCorpus(); + seedVideo("local"); + seedDriveChannel(); + // d2's audio is tiered: a relative link into the channel's media/. + mkdirSync(path.join(DRIVE_MEDIA(), "d2"), { recursive: true }); + writeFileSync(path.join(DRIVE_MEDIA(), "d2", "audio.mp3"), "AUDIO"); + symlinkSync("../../media/d2/audio.mp3", path.join(videoDir("d2", DRIVE_CHANNEL), "audio.mp3")); + await runIndex(); + assert.deepEqual(indexed(), [...DRIVE_VIDEOS, `${CHANNEL}/local`]); + + unmountMedia(); + try { + seedVideo("d3", DRIVE_CHANNEL); // arrived meanwhile, text on the corpus disk + const { res, log } = await runIndex(); + assert.deepEqual(res.heldChannels, [], log.join("\n")); + assert.equal(res.added, 1); + assert.equal(res.removed, 0); + assert.deepEqual(indexed(), [...DRIVE_VIDEOS, `${DRIVE_CHANNEL}/d3`, `${CHANNEL}/local`]); + } finally { + remountMedia(); + } +}); + +test("(k) release 17: a STALLED media drive with a tiered live chat holds neither the index nor the stats build, and is not asked", async () => { + resetCorpus(); + seedVideo("local"); + seedDriveChannel(); + // d1's raw live chat is tiered: the bytes on the drive, a relative link in + // data/d1 carrying the file's times, and no cues yet (so the build reads + // the raw — the fallback this case pins). + const chat = + JSON.stringify({ replayChatItemAction: { actions: [{ addChatItemAction: { item: { liveChatTextMessageRenderer: { message: { runs: [{ text: "hello" }] }, authorName: { simpleText: "a" }, timestampUsec: "1000000" } } } }], videoOffsetTimeMsec: "1000" } }) + "\n"; + mkdirSync(path.join(DRIVE_MEDIA(), "d1"), { recursive: true }); + const bytes = path.join(DRIVE_MEDIA(), "d1", "transcript.live_chat.json"); + writeFileSync(bytes, chat); + symlinkSync("../../media/d1/transcript.live_chat.json", path.join(videoDir("d1", DRIVE_CHANNEL), "transcript.live_chat.json")); + const first = await runIndex(); + assert.deepEqual(first.res.heldChannels, []); + const liveChatOf = () => + withIndex((db) => + [...db("subs").getRange()] + .filter(({ key }) => (key as unknown as string[])[2] === "d1") + .flatMap(({ value }) => (value as { track: string }[]).map((t) => t.track)), + ); + assert.ok(liveChatOf().includes("live_chat"), "read through the link while the drive answered"); + + // The drive stalls: marked (the health pass does this in the editor), and + // every call that would reach it hangs. A metadata rewrite makes d1 changed. + const { recordLocationHealth, resetStorageHealth } = await import("../lib/storageHealth"); + recordLocationHealth({ id: "usb", label: "USB drive", root: MEDIA }, "stalled"); + hangUnder = MEDIA; + driveCalls.length = 0; + try { + seedVideo("d1", DRIVE_CHANNEL, { subs: true, title: "Retitled" }); + const { res, log } = await runIndex(); + assert.deepEqual(res.heldChannels, [], log.join("\n")); + assert.equal(res.changed, 1); + assert.equal(res.removed, 0); + assert.ok(liveChatOf().includes("live_chat"), "the cues the last build held are kept"); + const { buildStats } = await import("./buildStats"); + const stats = await buildStats({ + paths, + onLog: () => {}, + wholePoolStatsDir: path.join(ROOT, "pool-stats"), + }); + assert.deepEqual(stats.heldChannels, []); + assert.deepEqual(driveCalls, [], "no call reached the stalled drive"); + } finally { + hangUnder = null; + resetStorageHealth(); + } +}); + +test("(l) release 17: a tiered raw live chat that cannot be read, with no cues to keep, is retried by the next build", async () => { + resetCorpus(); + seedVideo("local"); + seedDriveChannel(); + const chat = + JSON.stringify({ replayChatItemAction: { actions: [{ addChatItemAction: { item: { liveChatTextMessageRenderer: { message: { runs: [{ text: "hi" }] }, authorName: { simpleText: "a" }, timestampUsec: "1000000" } } } }], videoOffsetTimeMsec: "1000" } }) + "\n"; + mkdirSync(path.join(DRIVE_MEDIA(), "d2"), { recursive: true }); + writeFileSync(path.join(DRIVE_MEDIA(), "d2", "transcript.live_chat.json"), chat); + symlinkSync("../../media/d2/transcript.live_chat.json", path.join(videoDir("d2", DRIVE_CHANNEL), "transcript.live_chat.json")); + const tracksOfD2 = () => + withIndex((db) => + [...db("subs").getRange()] + .filter(({ key }) => (key as unknown as string[])[2] === "d2") + .flatMap(({ value }) => (value as { track: string }[]).map((t) => t.track)), + ); + // The media drive is away for the first build: nothing to read, nothing kept. + unmountMedia(); + try { + const away = await runIndex(); + assert.deepEqual(away.res.heldChannels, []); + assert.equal(tracksOfD2().includes("live_chat"), false); + } finally { + remountMedia(); + } + // Back: nothing on disk moved, yet the record is retried and the chat read. + const back = await runIndex(); + assert.equal(back.res.changed, 1, "d2 is retried"); + assert.ok(tracksOfD2().includes("live_chat")); + // And then it settles. + assert.equal((await runIndex()).res.changed, 0); +}); + test("(z) no write this file caused landed outside its temp root", () => { // LMDB writes natively, past the spy: its file must be under the root too. assert.ok(paths.lmdbPath.startsWith(ROOT + path.sep), paths.lmdbPath); diff --git a/common/controller/buildIndex.ts b/common/controller/buildIndex.ts @@ -28,6 +28,7 @@ import path from "node:path"; import { createHash } from "node:crypto"; import { + lstat, mkdir, readdir, readFile, @@ -95,11 +96,12 @@ import { } from "../lib/channelConfig"; import { readChannelConfigFile } from "./channels"; import { inspectChannelMedia } from "../lib/channelMedia"; +import { isTierable } from "../lib/mediaTier"; +import { onDrive } from "../lib/storageHealth"; import { HELD_WAYS_OUT, describeHeld, heldReason, - isMediaHeld, } from "../lib/channelMediaHold"; import { resolveChannelGroupId } from "../lib/channelGroups"; import type { Paths } from "../lib/paths"; @@ -314,17 +316,19 @@ async function scanSource( live: LiveEntry[]; channels: Map<string, ChannelConfig>; held: Map<string, string>; + mediaAccess: Map<string, MediaAccess>; }> { const locations = getSettings().storage.locations; const channels = new Map<string, ChannelConfig>(); const live: LiveEntry[] = []; const held = new Map<string, string>(); + const mediaAccess = new Map<string, MediaAccess>(); let channelEntries: Dirent[]; try { channelEntries = await readdir(channelsDir, { withFileTypes: true }); } catch { // Fresh transcripts dir with no channels yet. - return { live, channels, held }; + return { live, channels, held, mediaAccess }; } for (const ch of channelEntries) { if (!ch.isDirectory()) continue; @@ -343,25 +347,35 @@ async function scanSource( // and routing it through the video scan would only ever produce noise. Its // posts tree is built from the JSONL shards further down. if (isSocialChannel(cfg)) continue; - // An unmounted drive is not an empty channel (lib/channelMedia.ts): the - // readdir below would fail, and every record the channel has would be + // An unreadable text tier is not an empty channel (lib/channelMedia.ts): + // the readdir below would fail, and every record the channel has would be // removed as gone. + // + // THE TEXT GUARD, NOT THE MEDIA ONE (release 17): the index reads the text + // tier only — presence of a media file is a name in the one readdir, and + // every stat is a sidecar — so a moving, stalled or unmounted MEDIA drive + // never holds it. Held: a `legacy` channel (its text is on the far drive), + // a `data/` that is not a directory, a tier migration in flight. const media = await inspectChannelMedia({ channelsDir }, ch.name, cfg, { fresh: true, }); - if (isMediaHeld(media.status)) { - held.set(ch.name, heldReason(media, cfg.dataDir, locations)); + if (!media.text.readable) { + held.set(ch.name, heldReason(media, cfg.mediaDir ?? cfg.dataDir, locations)); continue; } + mediaAccess.set(ch.name, { + readable: media.status === "ok" || media.status === "in-place", + drive: cfg.mediaDir?.trim() || undefined, + }); const dataDir = path.join(channelDir, "data"); let videoEntries: Dirent[]; try { videoEntries = await readdir(dataDir, { withFileTypes: true }); } catch (err) { const code = errCode(err); - // No data/ at all on a channel whose media was never moved: it has - // downloaded nothing yet (or its media was deleted), and it IS empty. - if (code === "ENOENT" && media.status === "in-place") { + // No data/ at all: the text tier is on the corpus disk, so the channel + // has downloaded nothing yet (or its files were deleted), and it IS empty. + if (code === "ENOENT") { log(`Channel ${ch.name}: no data/ directory; indexed as a channel with no videos.`); continue; } @@ -412,7 +426,13 @@ async function scanSource( let subsMs: number | null = null; for (const t of subTracks) { try { - const ms = (await stat(path.join(fullVideoDir, t.filename))).mtimeMs; + // A TIERED FILE IS `lstat`ED (release 17): `transcript.live_chat.json` + // is a sub track AND media — on a tiered channel a link into + // channels/<slug>/media, possibly on another drive. The link carries + // the file's times (the tier hook's `lutimes`), so this answers from + // the corpus disk and never reaches the media drive. + const p = path.join(fullVideoDir, t.filename); + const ms = (await (isTierable(t.filename) ? lstat(p) : stat(p))).mtimeMs; if (subsMs === null || ms > subsMs) subsMs = ms; } catch { // ignore @@ -462,21 +482,25 @@ async function scanSource( held.set(ch.name, `a video in its data directory could not be read (${readFailure})`); continue; } - // Asked again after the walk: a drive that went away DURING it leaves the - // videos after that point missing from this scan, which would remove them. - // Three syscalls a channel. + // Asked again after the walk: a text tier that became unreadable DURING it + // (a tier migration that started) leaves the videos after that point + // missing from this scan, which would remove them. A few syscalls a channel. const after = await inspectChannelMedia({ channelsDir }, ch.name, cfg, { fresh: true, }); - if (isMediaHeld(after.status)) { - held.set(ch.name, heldReason(after, cfg.dataDir, locations)); + if (!after.text.readable) { + held.set(ch.name, heldReason(after, cfg.mediaDir ?? cfg.dataDir, locations)); continue; } for (const e of channelLive) live.push(e); } - return { live, channels, held }; + return { live, channels, held, mediaAccess }; } +// Whether a channel's media tier may be read in this build, and through which +// drive (scanSource, from the same inspect the text guard asked). +type MediaAccess = { readable: boolean; drive?: string }; + function pathKeyId(k: PathKey): string { return `${k[0]}\x00${k[1]}`; } @@ -646,6 +670,7 @@ export async function buildIndex({ live, channels: channelConfigs, held, + mediaAccess, } = await scanSource(channelsDir, log); // A full rebuild clears every channel's records, and a held channel cannot be @@ -856,6 +881,10 @@ export async function buildIndex({ const pk: PathKey = [s.channelSlug, s.videoDir]; const prev = mtimes.get(pk); + // What the last build held for this video's sub tracks, read before + // a re-keyed record is removed: a tiered live chat whose raw cannot + // be read now keeps the cues it had (below). + const prevSubs = prev ? subs.get(prev.indexKey) : undefined; if (prev && !indexKeysEqual(prev.indexKey, indexKey)) { sums.remove(prev.indexKey); cues.remove(prev.indexKey); @@ -868,6 +897,9 @@ export async function buildIndex({ else cues.remove(indexKey); const parsedSubs: StoredSubs = []; + // A tiered track that could be neither read nor kept: `subsMs` is + // stored null so the next build's scan sees a change and retries. + let retrySubs = false; for (const t of s.subTracks) { try { let trackCues: Cue[] | null = null; @@ -881,12 +913,43 @@ export async function buildIndex({ } } if (!trackCues) { - const raw = await readFile( - path.join(path.dirname(s.metaPath), t.filename), - "utf8", - ); - trackCues = - t.track === "live_chat" ? parseLiveChat(raw) : parseVtt(raw); + const rawPath = path.join(path.dirname(s.metaPath), t.filename); + // THE RAW REPLAY IS MEDIA (release 17). Tiered — a link into + // channels/<slug>/media — it is read only while the channel's + // media is reachable: the drive is PROBED through the watchdog + // (one stat, which is what its budget is sized for), then the + // file is read directly — a raw replay of hundreds of MB on a + // slow platter would outlast the budget, be refused, and, where + // the disk's counters are unknown, mark the location stalled. + // Not read (unreachable, not answering, unreadable): the cues + // the last build held are kept; with none to keep, the record + // is written with no sub-track time, so the next build retries. + const tiered = + t.track === "live_chat" && + (await lstat(rawPath).then((l) => l.isSymbolicLink(), () => false)); + if (tiered) { + const access = mediaAccess.get(s.channelSlug); + let raw: string | null = null; + if (access?.readable) { + try { + if (access.drive) await onDrive(access.drive, () => stat(rawPath)); + raw = await readFile(rawPath, "utf8"); + } catch { + raw = null; + } + } + if (raw === null) { + const kept = prevSubs?.find((x) => x.track === t.track); + if (kept) parsedSubs.push(kept); + else retrySubs = true; + continue; + } + trackCues = parseLiveChat(raw); + } else { + const raw = await readFile(rawPath, "utf8"); + trackCues = + t.track === "live_chat" ? parseLiveChat(raw) : parseVtt(raw); + } } if (trackCues.length > 0) { parsedSubs.push({ track: t.track, cues: trackCues }); @@ -988,7 +1051,7 @@ export async function buildIndex({ mtimes.put(pk, { metaMs: s.metaMs, transcriptMs: s.transcriptMs, - subsMs: s.subsMs, + subsMs: retrySubs ? null : s.subsMs, availabilityMs: s.availabilityMs, digestMs: s.digestMs, ...(availability ? { availability } : {}), diff --git a/common/controller/buildStats.test.ts b/common/controller/buildStats.test.ts @@ -455,19 +455,22 @@ test("(h) a video the index skipped is not announced as pending on every run", a assert.equal(res.notIndexable, 1, "the undated one stays, and is said as such"); }); -// A second channel whose data/ is a relocated symlink, the way the editor's -// Storage panel leaves it: channels/<slug>/data -> <root>/<slug>/data, with -// config.dataDir recording the target. -function seedDriveChannel(): { target: string } { - const target = path.join(ROOT, "media", DRIVE_CHANNEL, "data"); - mkdirSync(target, { recursive: true }); +// A second channel whose MEDIA is relocated (release 17): its text in a real +// data/ on the corpus disk, channels/<slug>/media -> <root>/<slug>/media, with +// config.mediaDir recording the target. +function driveConfig(extra: Record<string, unknown>): void { writeJson(path.join(paths.channelsDir, DRIVE_CHANNEL, "config.json"), { handling: "youtube", name: "Drive Channel", url: "https://www.youtube.com/@drive/videos", - dataDir: target, + ...extra, }); - symlinkSync(target, path.join(paths.channelsDir, DRIVE_CHANNEL, "data")); +} +function seedDriveChannel(): { target: string } { + const target = path.join(ROOT, "media", DRIVE_CHANNEL, "media"); + mkdirSync(target, { recursive: true }); + driveConfig({ mediaDir: target }); + symlinkSync(target, path.join(paths.channelsDir, DRIVE_CHANNEL, "media")); for (const id of ["d1", "d2"]) { seedVideo(id, "2026-07-11T11:00:00Z", {}, DRIVE_CHANNEL); addCaptions(id, "2026-07-11T12:00:00Z", DRIVE_CHANNEL); @@ -475,7 +478,45 @@ function seedDriveChannel(): { target: string } { return { target }; } -test("(i) an unmounted media drive keeps its channel's stats; a cache clear refuses", async () => { +// The text goes away: the channel put on the retired whole-directory layout +// (`data` an absolute link to <root>/<slug>/data, `dataDir` recorded) with its +// drive not mounted — the text moved aside, the link dangling. And back. +function retireAndUnmount(): void { + const data = path.join(paths.channelsDir, DRIVE_CHANNEL, "data"); + const away = path.join(ROOT, "media-away", DRIVE_CHANNEL, "data"); + mkdirSync(path.dirname(away), { recursive: true }); + renameSync(data, away); + const target = path.join(ROOT, "media", DRIVE_CHANNEL, "data"); + symlinkSync(target, data); + driveConfig({ dataDir: target }); +} +function migrateHome(): void { + const data = path.join(paths.channelsDir, DRIVE_CHANNEL, "data"); + rmSync(data); + renameSync(path.join(ROOT, "media-away", DRIVE_CHANNEL, "data"), data); + driveConfig({ mediaDir: path.join(ROOT, "media", DRIVE_CHANNEL, "media") }); +} + +test("(i2) release 17: an unmounted MEDIA drive does not hold the stats build", async () => { + resetCorpus(); + seedVideo("local"); + seedDriveChannel(); + await runIndex(); + await runStats(); + const media = path.join(ROOT, "media"); + renameSync(media, `${media}-unmounted`); + try { + const log: string[] = []; + const away = await runStats(log); + assert.deepEqual(away.res.heldChannels, [], log.join("\n")); + assert.equal(away.res.removed, 0); + assert.equal(statOf(away.byId, "d1").hasTranscript, true); + } finally { + renameSync(`${media}-unmounted`, media); + } +}); + +test("(i) a channel whose text cannot be read (the retired layout, drive away) keeps its stats; a cache clear refuses", async () => { resetCorpus(); seedVideo("local"); seedDriveChannel(); @@ -490,8 +531,7 @@ test("(i) an unmounted media drive keeps its channel's stats; a cache clear refu paths.settingsFile, JSON.stringify({ storage: { locations: [{ id: "usb", label: "USB drive", root: media, autoRepoint: false }] } }), ); - // Unmount: the link now dangles, exactly as an absent USB drive leaves it. - renameSync(media, `${media}-away`); + retireAndUnmount(); const log: string[] = []; const away = await runStats(log); assert.equal(away.res.removed, 0, "not read as a channel with no videos"); @@ -499,7 +539,7 @@ test("(i) an unmounted media drive keeps its channel's stats; a cache clear refu assert.deepEqual(away.res.heldChannels, [DRIVE_CHANNEL]); assert.equal(statOf(away.byId, "d1").hasTranscript, true); assert.ok( - log.some((l) => l.startsWith(`Channel ${DRIVE_CHANNEL}: its media is not reachable`) && l.includes("its 2 cached stat(s) are kept")), + log.some((l) => l.startsWith(`Channel ${DRIVE_CHANNEL}: its media layout is the retired`) && l.includes("its 2 cached stat(s) are kept")), log.join("\n"), ); @@ -509,7 +549,7 @@ test("(i) an unmounted media drive keeps its channel's stats; a cache clear refu await assert.rejects(runStats(), (err: Error) => { assert.match( err.message, - /must be rebuilt .* cannot be read: drive-channel \(its media is not reachable \(drive not mounted\?\), on location "USB drive"\)/, + /must be rebuilt .* cannot be read: drive-channel \(its media layout is the retired whole-directory one \(run archilyzer storage migrate-tier\), on location "USB drive"\)/, ); // The ways out, mounting first, and no path in the message. assert.match( @@ -522,7 +562,7 @@ test("(i) an unmounted media drive keeps its channel's stats; a cache clear refu assert.equal(readStoredSchema(), STATS_SCHEMA_VERSION - 1); assert.equal(countStats(), 3); - renameSync(`${media}-away`, media); + migrateHome(); const back = await runStats(); assert.deepEqual(back.res.heldChannels, []); assert.equal(readStoredSchema(), STATS_SCHEMA_VERSION); diff --git a/common/controller/buildStats.ts b/common/controller/buildStats.ts @@ -68,7 +68,6 @@ import { HELD_WAYS_OUT, describeHeld, heldReason, - isMediaHeld, } from "../lib/channelMediaHold"; import { getSettings } from "../lib/settings"; import type { Paths } from "../lib/paths"; @@ -286,13 +285,15 @@ async function scanSource( } if (cfg.excludeFromBuild) continue; channels.set(ch.name, cfg); - // An unmounted drive is not an empty channel (lib/channelMedia.ts): the - // readdir below would fail and every one of its stats would be removed. + // An unreadable text tier is not an empty channel (lib/channelMedia.ts): + // the readdir below would fail and every one of its stats would be removed. + // THE TEXT GUARD (release 17): the stats build reads the text tier only and + // is never held by the media tier — see buildIndex.ts's scanSource. const media = await inspectChannelMedia({ channelsDir }, ch.name, cfg, { fresh: true, }); - if (isMediaHeld(media.status)) { - held.set(ch.name, heldReason(media, cfg.dataDir, locations)); + if (!media.text.readable) { + held.set(ch.name, heldReason(media, cfg.mediaDir ?? cfg.dataDir, locations)); continue; } const dataDir = path.join(channelDir, "data"); diff --git a/common/controller/channelSnapshot.test.ts b/common/controller/channelSnapshot.test.ts @@ -12,7 +12,7 @@ import { foldBucketLaneEntry, generateChannelSnapshot, } from "./channelSnapshot"; -import { ChannelMediaUnreachableError } from "../lib/channelMedia"; +import { ChannelTextUnreadableError } from "../lib/channelMedia"; import type { Paths } from "../lib/paths"; import { emptyOperationCounts, @@ -332,11 +332,13 @@ test("only the reachable diarization states will ever clear a hold", () => { } }); -test("generateChannelSnapshot refuses an unreachable channel rather than writing an empty snapshot", async () => { +test("generateChannelSnapshot refuses a LEGACY channel rather than writing an empty snapshot", async () => { // The readdir inside it swallows ENOENT as "no videos", so without the guard a - // relocated channel whose drive is unmounted would publish a snapshot saying + // channel whose text is on an unmounted drive would publish a snapshot saying // every video is undownloaded — and all four lanes read that as work to do. // The throw is what makes the scheduler keep the last good snapshot.json. + // Since release 17 only the retired whole-directory layout puts the text on + // another drive, and the text guard refuses it whatever the drive is doing. const dir = await mkdtemp(path.join(tmpdir(), "ttb-snap-media-")); try { const paths = { channelsDir: path.join(dir, "channels") } as Paths; @@ -352,13 +354,52 @@ test("generateChannelSnapshot refuses an unreachable channel rather than writing await assert.rejects( () => generateChannelSnapshot(paths, "alpha"), - ChannelMediaUnreachableError, + (err: unknown) => + err instanceof ChannelTextUnreadableError && /migrate-tier alpha/.test(err.message), ); } finally { await rm(dir, { recursive: true, force: true }); } }); +test("release 17: an unmounted MEDIA drive does not stop the snapshot; its media bytes are unknown", async () => { + const dir = await mkdtemp(path.join(tmpdir(), "ttb-snap-tier-")); + try { + const paths = { channelsDir: path.join(dir, "channels") } as Paths; + const channelDir = path.join(paths.channelsDir, "alpha"); + const videoDir = path.join(channelDir, "data", "v1"); + await mkdir(videoDir, { recursive: true }); + const target = path.join(dir, "platter", "alpha", "media"); + await writeFile( + path.join(channelDir, "config.json"), + JSON.stringify({ handling: "youtube", url: "https://example.com/c", mediaDir: target }), + ); + await writeFile(path.join(channelDir, "playlist"), ""); + const meta = JSON.stringify({ id: "v1" }); + await writeFile(path.join(videoDir, "metadata.info.json"), meta); + await writeFile(path.join(videoDir, "transcript.json"), "{}"); + // The audio is tiered; the media link dangles (an unmounted drive). + await symlink(path.join("..", "..", "media", "v1", "audio.mp3"), path.join(videoDir, "audio.mp3")); + await symlink(target, path.join(channelDir, "media")); + + const snap = await generateChannelSnapshot(paths, "alpha"); + assert.equal(snap.totals.downloaded, 1, "the link is a name in the listing"); + assert.equal(snap.totalMediaBytes, undefined, "unknown, never 0"); + assert.equal(snap.totalAudioBytes, undefined); + assert.equal(snap.totalTextBytes, meta.length + 2); + + // The drive back: the link's bytes are counted as media. + await mkdir(path.join(target, "v1"), { recursive: true }); + await writeFile(path.join(target, "v1", "audio.mp3"), Buffer.alloc(300)); + const back = await generateChannelSnapshot(paths, "alpha"); + assert.equal(back.totalMediaBytes, 300); + assert.equal(back.totalAudioBytes, 300); + assert.equal(back.totalTextBytes, meta.length + 2); + } finally { + await rm(dir, { recursive: true, force: true }); + } +}); + // --------------------------------------------------------------------------- // The filtered channel, end to end through generateChannelSnapshot. // @@ -682,7 +723,7 @@ test("a chat already on disk stays a member after the mode is turned off", async }); -test("a video dir's clips/ counts into totalMediaBytes AND into totalClipsBytes", async () => { +test("a video dir's clips/ counts into totalClipsBytes, apart from the media and text tiers", async () => { // `clips/` is the one subdirectory a video dir has, and until this it was // counted by nothing: the loop did `if (!st.isFile()) continue` under a // comment saying a video dir is flat, which the clip-window feature made @@ -690,8 +731,9 @@ test("a video dir's clips/ counts into totalMediaBytes AND into totalClipsBytes" // /storage and every relocation estimate under-reported a channel that had // been walked by a report. // - // The two numbers OVERLAP on purpose: totalMediaBytes is every byte under - // data/<id>/, totalClipsBytes is the "of which". + // Release 17: the three are SIBLINGS. totalMediaBytes is the media tier's + // (the tierable names — audio, the raw live chat), totalTextBytes the rest of + // data/<id>/, totalClipsBytes the clip cache, which is never tiered. const dir = await mkdtemp(path.join(tmpdir(), "ttb-snap-clips-")); try { const paths = { channelsDir: path.join(dir, "channels") } as Paths; @@ -714,7 +756,8 @@ test("a video dir's clips/ counts into totalMediaBytes AND into totalClipsBytes" const snap = await generateChannelSnapshot(paths, "alpha"); assert.equal(snap.totalClipsBytes, 57); - assert.equal(snap.totalMediaBytes, 100 + meta.length + 57); + assert.equal(snap.totalMediaBytes, 100); + assert.equal(snap.totalTextBytes, meta.length); } finally { await rm(dir, { recursive: true, force: true }); } diff --git a/common/controller/channelSnapshot.ts b/common/controller/channelSnapshot.ts @@ -1,7 +1,7 @@ import { writeJsonAtomic } from "../lib/jsonFile-server"; import path from "node:path"; import type { Dirent } from "node:fs"; -import { readdir, readFile, stat } from "node:fs/promises"; +import { lstat, readdir, readFile, stat } from "node:fs/promises"; import pLimit from "p-limit"; import { readArchive } from "../lib/archive"; import { @@ -21,7 +21,8 @@ import { type Availability, } from "../lib/availability"; import { resolveCookiePolicy } from "../lib/cookiePolicy"; -import { assertChannelMediaReachable } from "../lib/channelMedia"; +import { assertChannelTextReadable } from "../lib/channelMedia"; +import { isTierable } from "../lib/mediaTier"; import { isDriveNotAnswering, onDrive } from "../lib/storageHealth"; import { CLIPS_DIR_NAME } from "../lib/clipWindow"; import { getSettings } from "../lib/settings"; @@ -349,21 +350,34 @@ export type ChannelSnapshot = { // must default to 0 — and MUST render that as "—", not "0", because a zero // here would claim a measurement nobody took. totalAudioBytes?: number; - // EVERY byte under `data/<id>/` for every video — audio, transcripts, cues, - // metadata, thumbnails, a persisted container. The figure `/storage` and the - // `/channels` Size column are priced in, and the one a relocation carries; + // THE MEDIA TIER's bytes (release 17): every file `isTierable` names — the + // audio and the raw live-chat replay — whether it is a link into + // `channels/<slug>/media/` already or still a real file a move's preflight + // would tier. The figure a media move carries and a media location holds; // `totalAudioBytes` is a fraction of it and is about what a CLEANUP could - // reclaim, which is a different question. + // reclaim, which is a different question. (Before release 17 this was every + // byte under `data/<id>/`; the text is now `totalTextBytes`.) // // Optional, and the distinction is load-bearing: a snapshot written before // this field existed lacks it, and a reader MUST render that as "size unknown // until Refresh report", never as 0 — a zero would rank a 400 GB channel - // bottom of a "free up N GB" list. + // bottom of a "free up N GB" list. ABSENT TOO when the channel's media tier + // could not be read (an unmounted, stalled or moving media drive): the + // snapshot is still written — its text is readable — and the bytes are + // unknown, never 0. totalMediaBytes?: number; - // OF WHICH: the bytes held by `data/<id>/clips/` — the clip windows umtool - // and the video page fetch a few seconds at a time. A SUBSET of - // `totalMediaBytes`, not a sibling of it: a window lives under the video dir, - // `rsync -a` carries it with everything else, and a volume is holding it. + // EVERYTHING ELSE UNDER `data/<id>/` BUT `clips/` — the bytes this channel + // holds on the corpus disk outside the tier and the clip cache: the text + // (transcripts, cues, metadata, sidecars, thumbnails) AND whatever is not + // tierable — a persisted source container not yet moved to the store, a + // partial download, scratch. Named for its bulk; on a channel with a few + // containers or partials it is more than the text alone. + // Absent on a snapshot written before release 17: "unknown", never 0. + totalTextBytes?: number; + // The bytes held by `data/<id>/clips/` — the clip windows umtool and the + // video page fetch a few seconds at a time. A SIBLING of the two tiers since + // release 17 (before it, a subset of `totalMediaBytes`): `clips/` stays on + // the corpus disk and is never tiered. // // Split out because it is the one part of a channel's bytes that is a CACHE: // nothing prunes a window, so a channel walked by many reports accumulates @@ -770,8 +784,13 @@ export async function generateChannelSnapshot( // the guard does not open config.json a second time for the one field it // needs. One extra await in front of a function that then walks the whole // channel. + // + // THE TEXT GUARD (release 17): the snapshot is a reading of the TEXT tier — + // names, sidecars, metadata — so a moving, stalled or unmounted MEDIA drive + // does not hold it (its media bytes are then unknown, below). Refused: a + // `legacy` channel, whose text is on the far drive. const config = await readChannelConfig(paths, slug); - await assertChannelMediaReachable(paths, slug, config); + const mediaLocation = await assertChannelTextReadable(paths, slug, config); // A CHANNEL ON ANOTHER DRIVE IS WALKED THROUGH THE WATCHDOG // (lib/storageHealth.ts `onDrive`): the data/ listing, the keep-latest keys and @@ -783,9 +802,18 @@ export async function generateChannelSnapshot( // scheduler keeps the last good snapshot.json, as on any failed refresh. // The reconcile pass just below is sequential (one read at a time) and is // not raced. - const drive = config?.dataDir?.trim() || undefined; - const through = <T>(read: () => Promise<T>): Promise<T> => - drive ? onDrive(drive, read) : read(); + // + // SINCE RELEASE 17 THE TEXT IS NEVER ON ANOTHER DRIVE: the text guard above + // refuses the one layout where it was (`legacy`), so `through` reads + // directly. What may be on another drive is the media tier, and its stats + // are the only calls that go through the watchdog (per video, below). + const through = <T>(read: () => Promise<T>): Promise<T> => read(); + const mediaDrive = config?.mediaDir?.trim() || undefined; + // Whether the media tier may be read at all this run: its links are statted + // only when the media is `ok` or `in-place`. Otherwise (unmounted, stalled, + // moving, inconsistent) no call is made to it and the bytes are unknown. + let mediaBytesKnown = + mediaLocation.status === "ok" || mediaLocation.status === "in-place"; // Heal any video dir that drifted from the canonical id layout before we read // data/* (best-effort; never fail snapshot generation on a reconcile error). @@ -882,54 +910,71 @@ export async function generateChannelSnapshot( limit(() => through(async () => { const dir = path.join(dataDir, id); const files = await readVideoFiles(dir, { checkUntranscribable: true }); - // EVERY FILE IN THE DIR, STATTED ONCE, feeding two numbers. + // EVERY FILE IN THE DIR, `lstat`ED ONCE, feeding the byte figures. // // `audioSizes` is what it always was: the real audio files, keyed by - // name, for the cleanup reclaim estimate. `mediaBytes` is new and is - // every byte this video dir holds — audio, transcripts, sidecars, - // thumbnails, a persisted container — because THAT is the number a - // relocation moves and a volume holds, and the audio total is only a - // fraction of it (a channel's transcripts, cues and metadata are not - // free). - // - // NOT A SECOND WALK: the readdir is `files.entries`, already in hand, - // and the audio stats this loop replaces were being paid anyway. What - // it adds is a stat per NON-audio entry — six to ten per video, warm - // inode cache, on a pass that already reads several sidecars per video. + // name, for the cleanup reclaim estimate. `mediaBytes` is the media + // tier's share (what `isTierable` names) and `textBytes` everything + // else; `clips/` (CLIPS_DIR_NAME, the fetched windows) is recursed ONE + // LEVEL into `clipsBytes`, apart from both — it is never tiered. // - // A VIDEO DIR IS NO LONGER FLAT, and exactly one subdirectory is the - // reason: `clips/` (CLIPS_DIR_NAME), the fetched clip windows. It is - // recursed ONE LEVEL — the windows and their `.json` sidecars are files, - // and nothing writes a directory under it — and its bytes are counted - // BOTH into `clipsBytes` and into `mediaBytes`. Every OTHER directory - // still counts nothing: there is no other one today, and a blanket - // recursion here would be the second walk this loop exists to avoid. - // - // Why clips count toward `mediaBytes` at all: `mediaBytes` is every - // byte under `data/<id>/`, which is what the volume is holding and what - // `rsync -a` carries in a relocation. A window was being left out of - // both numbers while sitting on the platter. + // BY FILE KIND (release 17): a real file's size is its `lstat`, on the + // corpus disk, no watchdog. A LINK is a tiered media file whose bytes + // are on the media tier — possibly another drive — so the links are + // statted together, as ONE call through the watchdog per video + // (`onDrive(mediaDir)`), and only while the media is readable. A drive + // that does not answer makes the channel's media bytes unknown; the + // snapshot is still written. const audioSet = new Set(files.audioFiles); const audioSizes: Record<string, number> = {}; let mediaBytes = 0; + let textBytes = 0; let clipsBytes = 0; + const linked: string[] = []; for (const name of files.entries) { try { - const st = await stat(path.join(dir, name)); + const st = await lstat(path.join(dir, name)); + if (st.isSymbolicLink()) { + linked.push(name); + continue; + } if (!st.isFile()) { if (st.isDirectory() && name === CLIPS_DIR_NAME) { - const bytes = await dirFileBytes(path.join(dir, name)); - clipsBytes += bytes; - mediaBytes += bytes; + clipsBytes += await dirFileBytes(path.join(dir, name)); } continue; } - mediaBytes += st.size; + if (isTierable(name)) mediaBytes += st.size; + else textBytes += st.size; if (audioSet.has(name)) audioSizes[name] = st.size; } catch { // ignore — file vanished or is unreadable } } + if (linked.length > 0 && mediaBytesKnown) { + const statLinks = async () => { + const sizes: Array<[string, number]> = []; + for (const name of linked) { + try { + const st = await stat(path.join(dir, name)); + if (st.isFile()) sizes.push([name, st.size]); + } catch { + // a dangling link: its bytes are not on the media tier + } + } + return sizes; + }; + try { + const sizes = await (mediaDrive ? onDrive(mediaDrive, statLinks) : statLinks()); + for (const [name, size] of sizes) { + mediaBytes += size; + if (audioSet.has(name)) audioSizes[name] = size; + } + } catch (err) { + if (!isDriveNotAnswering(err)) throw err; + mediaBytesKnown = false; + } + } let nativeId: string | null = null; if (wantNativeIndex && files.hasMeta) { try { @@ -999,6 +1044,7 @@ export async function generateChannelSnapshot( backfill, audioSizes, mediaBytes, + textBytes, clipsBytes, nativeId, availability, @@ -1121,12 +1167,10 @@ export async function generateChannelSnapshot( // hand on the perVideo entry, so this adds ZERO I/O to a pass that runs over // ~79,000 videos. let totalAudioBytes = 0; - // Every byte under `data/`, not only the audio. The figure the storage - // surfaces are priced in: how much a volume is holding for this channel, and - // how much a move would carry. + // The media tier's bytes, the text tier's and the clip cache's — three + // siblings since release 17 (see the field comments). let totalMediaBytes = 0; - // The `clips/` share of the above. Counted in BOTH, deliberately: see the - // field comment on `totalClipsBytes`. + let totalTextBytes = 0; let totalClipsBytes = 0; const heldAudioBytes = emptyHeldAudio(); const heldAudioCounts = emptyHeldAudio(); @@ -1144,6 +1188,7 @@ export async function generateChannelSnapshot( files, audioSizes, mediaBytes, + textBytes, clipsBytes, backfill, effectiveAvailability, @@ -1155,6 +1200,7 @@ export async function generateChannelSnapshot( if (isVideoTranscribed(files)) transcribed++; if (isVideoDownloaded(files)) downloaded++; totalMediaBytes += mediaBytes; + totalTextBytes += textBytes; totalClipsBytes += clipsBytes; // --- Hold attribution --------------------------------------------------- @@ -1662,8 +1708,9 @@ export async function generateChannelSnapshot( multipleAudioFormats: multipleAudioFormatsBytes, foreignAudio: foreignAudioBytes, }, - totalAudioBytes, - totalMediaBytes, + // Unknown, never a partial sum, when the media tier could not be read. + ...(mediaBytesKnown ? { totalAudioBytes, totalMediaBytes } : {}), + totalTextBytes, totalClipsBytes, heldAudioBytes, heldAudioCounts, diff --git a/common/controller/channelWriters.test.ts b/common/controller/channelWriters.test.ts @@ -220,6 +220,25 @@ test("through the live registry: cancel leaves it stopping, the job's end stamps assert.deepEqual(channelWriters(slug), []); }); +test("mediaOnly (release 17): a digest job and the digest lane are not media writers", () => { + const jobs = [ + job({ id: "J1", kind: "digest-channel-local", channelSlug: "alpha" }), + job({ id: "J2", kind: "normalize-transcripts", channelSlug: "alpha" }), + job({ id: "J3", kind: "whisper-all", channelSlug: "alpha" }), + ]; + const units = [ + unit("digest", { videoId: "d1" }), + unit("transcription", { videoId: "t1" }), + ]; + const all = channelWriters("alpha", { source: source(jobs, units) }); + assert.equal(all.length, 5); + const media = channelWriters("alpha", { source: source(jobs, units), mediaOnly: true }); + assert.deepEqual( + media.map((w) => (w.source === "job" ? w.jobId : `${w.lane}:${w.videoId}`)), + ["J3", "transcription:t1"], + ); +}); + // A GHOST NEVER HOLDS A MOVE (release 17 slice D0): a `running` meta a dead // process left on disk is no writer — before the boot pass closes it, and // after. diff --git a/common/controller/channelWriters.ts b/common/controller/channelWriters.ts @@ -4,7 +4,7 @@ import { type JobStatus, type JobTaskKind, } from "../jobs/registry"; -import { jobKindLabel } from "../jobs/jobKinds"; +import { jobKindLabel, kindNeedsMedia } from "../jobs/jobKinds"; import { LANES, type AutoQueueKind } from "../lib/autoQueueTypes"; import { getAutoRunnerStatus, type AutoRunnerInFlight } from "./autoRunner"; @@ -90,6 +90,11 @@ export type ChannelWritersOptions = { // `relocate-channel-media`: it is running on this channel's slug, it is the // asker, and the relocation queue runs one move at a time. ignoreKinds?: ReadonlyArray<string>; + // MEDIA WRITERS ONLY (release 17): a move of the media tier holds only the + // jobs that open or write a big file (`kindNeedsMedia`) and the lanes that + // do (every lane but digest). A digest, a normalize, an availability check + // reads and writes the text, which never moves, so it may run during a move. + mediaOnly?: boolean; source?: ChannelWritersSource; }; @@ -104,6 +109,7 @@ export function channelWriters( const queued: ChannelWriter[] = []; for (const j of source.jobs()) { if (j.channelSlug !== slug || ignore.has(j.kind)) continue; + if (opts.mediaOnly && !kindNeedsMedia(j.kind)) continue; // A CANCEL IS A REQUEST, NOT AN EXIT. The registry marks a running job // `cancelled` the moment it is asked to stop, and stamps `endedAt` when // its function has returned (registry.ts `finalize`). Until then it is @@ -132,6 +138,7 @@ export function channelWriters( const units: ChannelWriter[] = source .units() .filter(({ unit }) => unit.channelSlug === slug) + .filter(({ lane }) => !opts.mediaOnly || lane !== "digest") .map(({ lane, unit }) => ({ source: "lane" as const, lane, diff --git a/common/controller/channels.ts b/common/controller/channels.ts @@ -16,7 +16,7 @@ import { readVideoFiles, } from "../lib/videoStatus"; import { loadDigest } from "../lib/digest-server"; -import { channelMediaStall, readRelocationMarker } from "../lib/channelMedia"; +import { channelTextStall, readRelocationMarker } from "../lib/channelMedia"; import { isDriveNotAnswering, onDrive } from "../lib/storageHealth"; // TYPE-ONLY, and it must stay that way: ./channelSnapshot imports // readChannelConfig from this module, and it drags in the snapshot generator's @@ -88,8 +88,11 @@ async function hasDigestWithItems(videoDir: string): Promise<boolean> { ); } -// `drive` is the channel's configured target (`config.dataDir`) when its media -// is on another drive. Then every read goes through `onDrive`: none while that +// `drive` is the retired `config.dataDir` of a `legacy` channel, the one layout +// whose TEXT is on another drive (release 17: a relocated channel's text stays +// on the corpus disk, and this walk reads only names and text sidecars, so its +// `mediaDir` is never a reason to wrap it — onDrive is keyed by file kind). +// For a legacy channel every read goes through `onDrive`: none while that // location is stalled, at most `inFlightPerLocation` (4) in flight on it, and // one that has not answered within the budget (`storage.health.budgetMs`, 3 s // by default) marks it stalled — and the walk answers null ("the drive did @@ -212,8 +215,9 @@ export async function readChannelStat( // one-second job-list poll asks this for every channel with a job listed, and // on a stalled drive each readdir and stat in the walk would hold an I/O // thread until the drive came back. No counts is what a caller already - // handles (the row draws no progress bar). - if (channelMediaStall(config)) return null; + // handles (the row draws no progress bar). Only a legacy channel's text is + // on a drive that can stall; a stalled MEDIA drive does not stop the count. + if (channelTextStall(config)) return null; const channelDir = path.join(paths.channelsDir, slug); const counts = await countDataFiles( path.join(channelDir, "data"), diff --git a/common/controller/cleanAudioFromTranscribed.ts b/common/controller/cleanAudioFromTranscribed.ts @@ -1,3 +1,4 @@ +import { removeMediaFile } from "../lib/mediaTier-server"; import path from "node:path"; import fs from "fs-extra"; import type { Paths } from "../lib/paths"; @@ -10,7 +11,7 @@ import { computeKeptVideoIds } from "./keptVideos"; import { pruneSavedVideos } from "./pruneSavedVideos"; import { excludedIds, verifyBeforeClean } from "./verifyBeforeClean"; -const { pathExists, readdir, remove } = fs; +const { pathExists, readdir } = fs; export type CleanAudioOptions = { channelSlug: string; @@ -155,7 +156,8 @@ export async function cleanAudioFromTranscribed({ if (signal?.aborted) break; if (excluded.has(candidate.id)) continue; for (const f of candidate.audioFiles) { - await remove(path.join(candidate.videoDir, f)); + // Through its link when tiered (release 17): the bytes on the media tier go too. + await removeMediaFile(candidate.videoDir, f); log(`Removed ${candidate.id}/${f}`); removedFiles++; } diff --git a/common/controller/cleanExtraAudioFormats.ts b/common/controller/cleanExtraAudioFormats.ts @@ -1,3 +1,4 @@ +import { removeMediaFile } from "../lib/mediaTier-server"; import path from "node:path"; import fs from "fs-extra"; import type { Paths } from "../lib/paths"; @@ -72,7 +73,7 @@ export async function cleanExtraAudioFormats({ continue; } for (const f of extras) { - await remove(path.join(videoDir, f)); + await removeMediaFile(videoDir, f); // derefs a tiered link (release 17) log(`Removed ${id}/${f}`); removedFiles++; } diff --git a/common/controller/evictClipWindows.test.ts b/common/controller/evictClipWindows.test.ts @@ -117,21 +117,24 @@ test("a window whose sidecar was never written is still evictable", async () => }); }); -test("a channel whose media is in transition is SKIPPED, not evicted from", async () => { - // `.relocating.json` means there may be two copies and the mover's verify - // pass compares trees — deleting under it turns a copy that had verified into - // a failed move of a channel that was fine. +test("a tier migration in flight is SKIPPED; a media move is not (release 17)", async () => { + // A tier migration rebuilds `data/` itself (`scope: "tier-migration"`): its + // copy and verify compare the text trees, `clips/` included, so deleting + // under it turns a copy that had verified into a failed migration. A media + // move carries `media/` only — `clips/` is never in it — so eviction goes on. await withTmp(async (paths) => { const clipsDir = await seed(paths, "alpha", { "10.00-40.00.mp4": { ageDays: 90, size: 1000 }, }); + const marker = path.join(paths.channelsDir, "alpha", ".relocating.json"); await writeFile( - path.join(paths.channelsDir, "alpha", ".relocating.json"), + marker, JSON.stringify({ - target: "/mnt/platter/alpha/data", + target: "/mnt/platter/alpha/media", direction: "out", startedAt: new Date().toISOString(), phase: "copy", + scope: "tier-migration", }), ); const r = await evictClipWindows({ paths, olderThanDays: 30 }); @@ -143,6 +146,21 @@ test("a channel whose media is in transition is SKIPPED, not evicted from", asyn // A skip is the ANSWER, so it reaches the operator in the summary rather // than being swallowed by a clean-looking zero. assert.match(evictClipWindowsSummary(r), /Skipped: alpha:/); + + await writeFile( + marker, + JSON.stringify({ + target: "/mnt/platter/alpha/media", + direction: "out", + startedAt: new Date().toISOString(), + phase: "copy", + scope: "media", + }), + ); + const moving = await evictClipWindows({ paths, olderThanDays: 30 }); + assert.deepEqual(moving.skipped, []); + assert.equal(moving.windows, 1); + assert.equal(await exists(path.join(clipsDir, "10.00-40.00.mp4")), false); }); }); @@ -215,7 +233,28 @@ test("an unmounted drive is a SKIP, never a clean eviction of nothing", async () assert.equal(r.channels, 0); assert.equal(r.windows, 0); assert.equal(r.skipped.length, 1); - assert.match(r.skipped[0], /^alpha: media unreachable/); + // Release 17: a data link is the retired layout, `legacy`, with its way out. + assert.match(r.skipped[0], /^alpha: media legacy \(.*migrate-tier alpha\)/); + }); +}); + +test("release 17: an unmounted MEDIA drive does not stop eviction — clips/ is on the corpus disk", async () => { + await withTmp(async (paths) => { + const clipsDir = await seed(paths, "alpha", { + "10.00-40.00.mp4": { ageDays: 90, size: 1000 }, + }); + const channelDir = path.join(paths.channelsDir, "alpha"); + const target = path.join(paths.transcriptsDir, "platter", "alpha", "media"); + await writeFile( + path.join(channelDir, "config.json"), + JSON.stringify({ url: "https://example.com/c", mediaDir: target }), + ); + // The media link dangles: the media drive is not mounted. + await symlink(target, path.join(channelDir, "media")); + const r = await evictClipWindows({ paths, olderThanDays: 30 }); + assert.deepEqual(r.skipped, []); + assert.equal(r.windows, 1); + assert.equal(await exists(path.join(clipsDir, "10.00-40.00.mp4")), false); }); }); diff --git a/common/controller/evictClipWindows.ts b/common/controller/evictClipWindows.ts @@ -82,24 +82,24 @@ async function evictChannel( ): Promise<void> { const channelDir = path.join(opts.paths.channelsDir, slug); - // AN UNMOUNTED DRIVE IS NOT AN EMPTY CHANNEL, and `inspectChannelMedia` is - // the one module that can tell the two apart (AGENTS.md: a path that reads + // AN UNREADABLE TEXT TIER IS NOT AN EMPTY CHANNEL, and `inspectChannelMedia` + // is the one module that can tell the two apart (AGENTS.md: a path that reads // `data/` guards there). Every enumerator else swallows ENOENT as "no - // videos" — which here would report a clean eviction of zero bytes about a - // platter full of windows, and an operator would read that as "nothing to - // reclaim". + // videos" — which here would report a clean eviction of zero bytes, and an + // operator would read that as "nothing to reclaim". // - // `in-transition` is a refusal for a sharper reason than unreachability: - // there may be two copies, `data/` may be a link whose target is half - // written, and the mover's verify pass compares trees. Deleting under it - // turns a copy that had verified into a failed move of a channel that was - // fine. (The JOB also declares `needsMedia: true`, which covers a - // single-channel run before it starts; this covers the corpus-wide one, - // where there is no slug for that guard to check.) + // THE TEXT GUARD (release 17): `clips/` is on the corpus disk and is never + // tiered, so a channel whose MEDIA is moving, stalled or unmounted is evicted + // as usual — a media move carries `media/`, never `data/`. Skipped: a + // `legacy` channel (the retired whole-directory layout, its `data/` — clips + // included — on the far drive, where a move's verify compares trees) and a + // tier migration in flight, which rebuilds `data/` itself. (The JOB also + // declares `needsText`, which covers a single-channel run before it starts; + // this covers the corpus-wide one, where there is no slug for that guard.) const media = await inspectChannelMedia(opts.paths, slug, undefined, { fresh: true, }); - if (media.status !== "ok" && media.status !== "in-place") { + if (!media.text.readable) { out.skipped.push( `${slug}: media ${media.status}${media.detail ? ` (${media.detail})` : ""} — nothing was touched`, ); @@ -111,8 +111,8 @@ async function evictChannel( try { videoIds = await readdir(dataDir); } catch { - // Past the guard above this really is "no videos": a channel whose media - // is in place and whose `data/` has never been created. + // Past the guard above this really is "no videos": a channel whose text + // tier is on the corpus disk and whose `data/` has never been created. out.channels += 1; return; } diff --git a/common/controller/mediaTierHooks.test.ts b/common/controller/mediaTierHooks.test.ts @@ -0,0 +1,145 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { + chmod, + lstat, + mkdir, + mkdtemp, + readFile, + readlink, + rename, + rm, + symlink, + writeFile, +} from "node:fs/promises"; +import { existsSync } from "node:fs"; +import { tmpdir } from "node:os"; +import path from "node:path"; +import type { Paths } from "../lib/paths"; +import { tierVideoDir } from "../lib/mediaTier-server"; +import { transcodeAudio } from "./transcode"; +import { normalizeLiveChat, isLiveChatCuesFresh } from "./normalizeLiveChat"; +import { removeWrongFormatAudio } from "./removeWrongFormatAudio"; +import { normalizeAllLiveChat } from "./normalizeAllLiveChat"; + +// Run with: +// pnpm --filter yt-dlp-transcript-common exec tsx --test controller/mediaTierHooks.test.ts +// +// RELEASE 17, THE MEDIA TIER, AT ITS CALL SITES: the finalisers call the hook +// and the deleters go through the link. A channel tiered IN PLACE — a real +// `channels/<slug>/media/` on the same disk — is the fixture: the hook's +// cross-device path is lib/mediaTier-server.test.ts's. + +async function fixture() { + const root = await mkdtemp(path.join(tmpdir(), "media-hooks-")); + const paths = { + transcriptsDir: root, + channelsDir: path.join(root, "channels"), + ffmpegBin: path.join(root, "ffmpeg"), + } as Paths; + const channelDir = path.join(paths.channelsDir, "chan"); + const videoDir = path.join(channelDir, "data", "v1"); + const mediaDir = path.join(channelDir, "media"); + await mkdir(videoDir, { recursive: true }); + await mkdir(mediaDir); + await writeFile( + path.join(channelDir, "config.json"), + JSON.stringify({ handling: "transcribe", url: "https://example.com/c", audioFormat: "mp3" }), + ); + await writeFile( + path.join(videoDir, "metadata.info.json"), + JSON.stringify({ id: "v1", title: "A video", duration: 60 }), + ); + // A fake ffmpeg: copies its input (-i) to its last argument. + await writeFile( + paths.ffmpegBin, + '#!/bin/sh\nin=""\nprev=""\nfor a in "$@"; do if [ "$prev" = "-i" ]; then in="$a"; fi; prev="$a"; last="$a"; done\ncp "$in" "$last"\n', + ); + await chmod(paths.ffmpegBin, 0o755); + return { root, paths, channelDir, videoDir, mediaDir }; +} + +test("transcodeAudio tiers its output, and re-tiers over a link the rename replaced", async () => { + const f = await fixture(); + try { + await writeFile(path.join(f.videoDir, "audio.m4a"), "SOURCE-1"); + const run = () => + transcodeAudio({ + paths: f.paths, + videoDir: f.videoDir, + sourceFilename: "audio.m4a", + targetFormat: "mp3", + onLog: () => {}, + signal: new AbortController().signal, + }); + await run(); + const out = path.join(f.videoDir, "audio.mp3"); + assert.ok((await lstat(out)).isSymbolicLink()); + assert.equal(await readlink(out), "../../media/v1/audio.mp3"); + assert.equal(await readFile(out, "utf8"), "SOURCE-1"); + // Again, over the link: the stale bytes in media/ are replaced, the link + // made again — nothing orphaned, nothing left real. + await writeFile(path.join(f.videoDir, "audio.m4a"), "SOURCE-2"); + await run(); + assert.ok((await lstat(out)).isSymbolicLink()); + assert.equal(await readFile(path.join(f.mediaDir, "v1", "audio.mp3"), "utf8"), "SOURCE-2"); + } finally { + await rm(f.root, { recursive: true, force: true }); + } +}); + +test("normalizeLiveChat tiers the raw replay after writing the cues; freshness never reaches the media drive", async () => { + const f = await fixture(); + try { + await writeFile(path.join(f.videoDir, "transcript.live_chat.json"), ""); + const first = await normalizeLiveChat({ videoDir: f.videoDir, channelSlug: "chan" }); + assert.equal(first.status, "wrote"); + const raw = path.join(f.videoDir, "transcript.live_chat.json"); + assert.ok((await lstat(raw)).isSymbolicLink()); + assert.ok((await lstat(path.join(f.videoDir, "live_chat.cues.json"))).isFile(), "the cues are text"); + // The media drive goes away: the cues are still fresh, from the link's mtime. + await rename(f.mediaDir, `${f.mediaDir}-away`); + assert.equal((await isLiveChatCuesFresh(f.videoDir)).fresh, true); + assert.equal((await normalizeLiveChat({ videoDir: f.videoDir, channelSlug: "chan" })).status, "fresh"); + await rename(`${f.mediaDir}-away`, f.mediaDir); + } finally { + await rm(f.root, { recursive: true, force: true }); + } +}); + +test("a deleter derefs: removeWrongFormatAudio removes the bytes in media/, not just the link", async () => { + const f = await fixture(); + try { + await writeFile(path.join(f.videoDir, "audio.mp3"), "KEEP"); + await writeFile(path.join(f.videoDir, "audio.m4a"), "WRONG"); + await tierVideoDir(f.videoDir); + assert.ok((await lstat(path.join(f.videoDir, "audio.m4a"))).isSymbolicLink()); + const r = await removeWrongFormatAudio({ channelSlug: "chan", paths: f.paths, onLog: () => {} }); + assert.equal(r.removedFiles, 1); + assert.equal(existsSync(path.join(f.mediaDir, "v1", "audio.m4a")), false, "no orphan on the media tier"); + await assert.rejects(lstat(path.join(f.videoDir, "audio.m4a"))); + assert.equal(await readFile(path.join(f.videoDir, "audio.mp3"), "utf8"), "KEEP"); + } finally { + await rm(f.root, { recursive: true, force: true }); + } +}); + +test("normalize-live-chat skips a channel whose media is not reachable", async () => { + const f = await fixture(); + try { + // Relocated media on an unmounted drive. + await rm(f.mediaDir, { recursive: true }); + const target = path.join(f.root, "platter", "chan", "media"); + await writeFile( + path.join(f.channelDir, "config.json"), + JSON.stringify({ handling: "transcribe", url: "https://example.com/c", mediaDir: target }), + ); + await symlink(target, f.mediaDir); + const lines: string[] = []; + const r = await normalizeAllLiveChat({ paths: f.paths, onLog: (l) => lines.push(l) }); + assert.equal(r.failed, 1); + assert.ok(lines.some((l) => /Normalize live chat chan: SKIPPED — .*not reachable/.test(l)), lines.join("\n")); + } finally { + await rm(f.root, { recursive: true, force: true }); + } +}); diff --git a/common/controller/normalizeAll.ts b/common/controller/normalizeAll.ts @@ -9,7 +9,7 @@ import pLimit from "p-limit"; import { listChannelStatsFromDisk } from "./channels"; import { normalizeTranscript } from "./normalizeTranscript"; import type { Paths } from "../lib/paths"; -import { assertChannelMediaReachable } from "../lib/channelMedia"; +import { assertChannelTextReadable } from "../lib/channelMedia"; export type NormalizeAllOptions = { paths: Paths; @@ -59,8 +59,10 @@ export async function normalizeAllTranscripts( // "this channel's drive is not mounted", and here the second reads as a // clean run over zero videos that reports 0/0/0/0 and moves on. A sweep // that silently skips a channel is worse than one that stops on it. + // THE TEXT GUARD (release 17): normalize reads and writes text only, so a + // moving, stalled or unmounted MEDIA drive does not stop it. try { - await assertChannelMediaReachable(opts.paths, ch.slug, ch.config); + await assertChannelTextReadable(opts.paths, ch.slug, ch.config); } catch (err) { log(`Normalize ${ch.slug}: SKIPPED — ${(err as Error).message}`); result.failed++; diff --git a/common/controller/normalizeAllLiveChat.ts b/common/controller/normalizeAllLiveChat.ts @@ -9,6 +9,7 @@ import pLimit from "p-limit"; import { listChannelStatsFromDisk } from "./channels"; import { normalizeLiveChat } from "./normalizeLiveChat"; import type { Paths } from "../lib/paths"; +import { assertChannelMediaReachable } from "../lib/channelMedia"; export type NormalizeAllLiveChatOptions = { paths: Paths; @@ -38,6 +39,18 @@ export async function normalizeAllLiveChat( }; for (const ch of channels) { if (opts.signal?.aborted) break; + // THE MEDIA GUARD (release 17): the raw replay is media — on a tiered + // channel a link into channels/<slug>/media, possibly on another drive — + // and a stale cues file is re-derived by READING it. A channel whose media + // is not reachable (unmounted, stalled, moving, legacy) is skipped and + // counted, never read as a clean pass over zero videos. + try { + await assertChannelMediaReachable(opts.paths, ch.slug, ch.config); + } catch (err) { + log(`Normalize live chat ${ch.slug}: SKIPPED — ${(err as Error).message}`); + result.failed++; + continue; + } const dataDir = path.join(opts.paths.channelsDir, ch.slug, "data"); const videoIds = await readdir(dataDir).catch(() => [] as string[]); log(`Normalize live chat ${ch.slug}: ${videoIds.length} videos`); diff --git a/common/controller/normalizeLiveChat.ts b/common/controller/normalizeLiveChat.ts @@ -4,8 +4,9 @@ // matches NormalizedTranscript with source: "live_chat". import path from "node:path"; -import { readFile, stat } from "node:fs/promises"; +import { lstat, readFile, stat } from "node:fs/promises"; import { writeJsonAtomic } from "../lib/jsonFile-server"; +import { tierMediaFile } from "../lib/mediaTier-server"; import { parseLiveChat } from "../lib/liveChat"; import type { Cue } from "../lib/vtt"; import { summarize, type RawMetadata } from "../lib/transcripts-server"; @@ -26,6 +27,9 @@ export type NormalizeLiveChatOptions = { configName?: string; log?: (msg: string) => void; force?: boolean; + // Tier the raw replay after writing the cues (default true). The export + // build's archive pass passes false: a build reads, it does not move media. + tier?: boolean; }; export type NormalizeLiveChatOutcome = @@ -41,6 +45,20 @@ async function mtimeMs(p: string): Promise<number | null> { } } +// THE RAW REPLAY'S mtime, WITHOUT FOLLOWING A LINK (release 17). The raw file +// is media: on a tiered channel `transcript.live_chat.json` is a relative link +// into channels/<slug>/media, possibly on another drive. The tier hook copies +// the file's mtime onto the link (`lutimes`), so the link answers the same +// freshness question from the corpus disk — the index build asks it per video +// and must never reach the media drive to do so. +async function rawMtimeMs(p: string): Promise<number | null> { + try { + return (await lstat(p)).mtimeMs; + } catch { + return null; + } +} + // TODO: live_chat files can be hundreds of MB. parseLiveChat already splits // on "\n", so a streaming readline variant is straightforward if we hit a // memory wall. For now this matches the readFile pattern used by the @@ -54,7 +72,7 @@ export async function normalizeLiveChat( const [metaStatMs, rawStatMs, cuesMs] = await Promise.all([ mtimeMs(metaPath), - mtimeMs(rawPath), + rawMtimeMs(rawPath), mtimeMs(cuesPath), ]); @@ -101,6 +119,12 @@ export async function normalizeLiveChat( opts.log?.( `Normalized live chat ${opts.channelSlug}/${path.basename(opts.videoDir)} (${cues.length} cues)`, ); + // THE MEDIA TIER'S HOOK (release 17): the raw replay is read once, here, and + // every other reader uses the cues just written — so it moves into + // channels/<slug>/media now (a relative link stays). Never throws. + if (opts.tier !== false) await tierMediaFile(opts.videoDir, LIVE_CHAT_FILENAME, { + onLog: opts.log ? (line) => opts.log?.(line.trimEnd()) : undefined, + }); return { status: "wrote", cuesPath }; } @@ -130,7 +154,7 @@ export async function isLiveChatCuesFresh( const [cuesMs, metaMs, rawMs] = await Promise.all([ mtimeMs(cuesPath), mtimeMs(metaPath), - mtimeMs(rawPath), + rawMtimeMs(rawPath), ]); if (cuesMs === null || metaMs === null || rawMs === null) { return { fresh: false, cuesPath }; diff --git a/common/controller/operationBatch.ts b/common/controller/operationBatch.ts @@ -40,6 +40,7 @@ import { getSettings, type SiteSettings } from "../lib/settings"; import { isGateHeld } from "../lib/pauseGates"; import { assertChannelMediaReachable, + assertChannelTextReadable, readRelocationMarker, } from "../lib/channelMedia"; import { runPool } from "../jobs/concurrentRunner"; @@ -1590,7 +1591,16 @@ export async function runOperationBatch( // of this function's callers are already behind guard 1 today; this is here so // a fourth in-process caller that is not cannot quietly find "no candidates" // on a channel whose drive is unmounted. - await assertChannelMediaReachable(opts.paths, opts.channelSlug); + // + // BY LANE (release 17): the digest lane reads only the text tier, so it asks + // the text guard and runs while the channel's media is moving, stalled or + // unmounted; the backfill lane's operations open the audio and keep the + // media guard. + if (opts.lane === "digest") { + await assertChannelTextReadable(opts.paths, opts.channelSlug); + } else { + await assertChannelMediaReachable(opts.paths, opts.channelSlug); + } const dataDir = path.join(opts.paths.channelsDir, opts.channelSlug, "data"); const allDirs = await readdir(dataDir).catch(() => [] as string[]); const onDisk = new Set(allDirs); diff --git a/common/controller/purgeSupersededAutoSubs.ts b/common/controller/purgeSupersededAutoSubs.ts @@ -1,3 +1,4 @@ +import { removeMediaFile } from "../lib/mediaTier-server"; import path from "node:path"; import fs from "fs-extra"; import type { Paths } from "../lib/paths"; @@ -9,7 +10,7 @@ import { listTranscriptVtts, } from "../lib/videoStatus"; -const { pathExists, readdir, remove } = fs; +const { pathExists, readdir } = fs; // Delete the YouTube auto-caption VTTs that our own transcript has superseded — // the manual counterpart to the `supersededAutoSubs` snapshot bucket, and the @@ -88,7 +89,7 @@ export async function purgeSupersededAutoSubs({ skipped++; continue; } - await remove(path.join(videoDir, name)); + await removeMediaFile(videoDir, name); // the one rm of a video-dir entry (release 17) log(`Removed ${id}/${name}`); removedFiles++; removedHere++; diff --git a/common/controller/recencyIndex.ts b/common/controller/recencyIndex.ts @@ -196,8 +196,12 @@ async function readTailUploadDate(file: string): Promise<string | null> { // Date the ids in `wanted` from their on-disk metadata. Removes each id it // keys from `wanted`, like interpolateFromPlaylist. // -// `drives` maps a channel whose media is on another drive to its configured -// target. Its reads go through `onDrive`: none while that drive's location is +// `drives` maps a channel whose TEXT is on another drive — a `legacy` channel, +// the retired whole-directory layout, keyed by its `dataDir` — to that target. +// Since release 17 a relocated channel's text (`metadata.info.json` here) is on +// the corpus disk and only its media is on the far drive, so its reads pass no +// drive at all: onDrive is keyed by file kind. A legacy channel's reads go +// through `onDrive`: none while that drive's location is // stalled (lib/storageHealth.ts) — each would hold an I/O thread until the drive // came back, 32 at a time — at most `inFlightPerLocation` (4) in flight on it, // and one that has not answered within the budget (3 s by default) marks it @@ -349,7 +353,7 @@ export type BuildRecencyKeysArgs = { // are on another drive (their reads go through the stall watchdog). meta: ReadonlyArray<{ slug: string; - config?: Pick<ChannelConfig, "dataDir"> | null; + config?: Pick<ChannelConfig, "dataDir" | "mediaDir"> | null; }>; // Every video id the caller might sort. Ids outside this set are not keyed. candidateIds: ReadonlySet<string>; @@ -508,6 +512,7 @@ export async function buildRecencyKeys({ if (owner && missing.size > 0) { const drives = new Map<string, string>(); for (const m of meta) { + // The retired `dataDir` only: `mediaDir` holds no metadata. const dir = m.config?.dataDir?.trim(); if (dir) drives.set(m.slug, dir); } diff --git a/common/controller/relocateChannelMedia.test.ts b/common/controller/relocateChannelMedia.test.ts @@ -42,6 +42,11 @@ import { } from "../lib/storageVolumes"; import { readChannelConfig } from "./channels"; +// RELEASE 17 SLICE T1 made a channel whose `data/` is a link (or whose config +// carries `dataDir`) `legacy`; these cases still build that retired layout and +// expect it to read `ok`. Slice T2 rebases them on `media/` and un-skips them. +const T1_SKIP = "release 17 T2 rebases the mover on media/"; + // Run with: // pnpm --filter yt-dlp-transcript-common exec tsx --test controller/relocateChannelMedia.test.ts // @@ -135,7 +140,7 @@ async function seed( return channelDir; } -test("out: copies, links, records the target, keeps mtimes and reclaims the source", async () => { +test("out: copies, links, records the target, keeps mtimes and reclaims the source", { skip: T1_SKIP }, async () => { await withTmp(async (paths, root) => { const channelDir = await seed(paths, "alpha", { v1: { "audio.m4a": "one".repeat(500), "transcript.json": "{}" }, @@ -192,7 +197,7 @@ test("out: copies, links, records the target, keeps mtimes and reclaims the sour }); }); -test("abort from an onLog hook leaves the source intact, and the rerun completes", async () => { +test("abort from an onLog hook leaves the source intact, and the rerun completes", { skip: T1_SKIP }, async () => { await withTmp(async (paths, root) => { const channelDir = await seed(paths, "alpha", { v1: { "audio.m4a": "x".repeat(20000) }, @@ -329,7 +334,7 @@ test("a verify failure keeps the source and does not write the config", async () }); }); -test("back: restores a real directory, clears the config and reclaims the target", async () => { +test("back: restores a real directory, clears the config and reclaims the target", { skip: T1_SKIP }, async () => { await withTmp(async (paths, root) => { const channelDir = await seed(paths, "alpha", { v1: { "audio.m4a": "one", "transcript.json": "{}" }, @@ -557,7 +562,7 @@ async function copyTree(src: string, dest: string): Promise<void> { } } -test("out @ swap: crash before the rename — the rerun re-verifies and completes", async () => { +test("out @ swap: crash before the rename — the rerun re-verifies and completes", { skip: T1_SKIP }, async () => { await withTmp(async (paths, root) => { const channelDir = await seed(paths, "alpha", { v1: { "audio.m4a": "one".repeat(500) }, @@ -583,7 +588,7 @@ test("out @ swap: crash before the rename — the rerun re-verifies and complete }); }); -test("out @ swap: crash after the config write — the first run's parked copy is still reclaimed", async () => { +test("out @ swap: crash after the config write — the first run's parked copy is still reclaimed", { skip: T1_SKIP }, async () => { await withTmp(async (paths, root) => { const channelDir = await seed(paths, "alpha", { v1: { "audio.m4a": "one".repeat(500) }, @@ -621,7 +626,7 @@ test("out @ swap: crash after the config write — the first run's parked copy i }); }); -test("out @ reclaim: the rerun sweeps every parked copy and clears the marker", async () => { +test("out @ reclaim: the rerun sweeps every parked copy and clears the marker", { skip: T1_SKIP }, async () => { await withTmp(async (paths, root) => { const channelDir = await seed(paths, "alpha", { v1: { "audio.m4a": "one".repeat(500) }, @@ -832,7 +837,7 @@ async function pathIsThere(p: string): Promise<boolean> { // the run copied the target to data.incoming, verified it, found data/ already a // real dir, logged "the swap had completed", cleared the config, rm -r'd the // target and then swept data.incoming. Three copies in, zero out. -test("back: an inconsistent channel is refused, and the target keeps its bytes", async () => { +test("back: an inconsistent channel is refused, and the target keeps its bytes", { skip: T1_SKIP }, async () => { await withTmp(async (paths, root) => { const channelDir = await seed(paths, "alpha", { v1: { "audio.m4a": "one" }, @@ -886,7 +891,7 @@ test("back: an inconsistent channel is refused, and the target keeps its bytes", // An `unreachable` channel (the drive is not mounted) is refused for the same // reason: nothing can vouch for what the target holds, and the run ends by // deleting it. -test("back: an unreachable channel is refused", async () => { +test("back: an unreachable channel is refused", { skip: T1_SKIP }, async () => { await withTmp(async (paths, root) => { await seed(paths, "alpha", { v1: { "audio.m4a": "one" } }); const target = relocatedDataDir(root, "alpha"); @@ -1020,7 +1025,7 @@ test("a relative root is refused by the preview, not only by the job", async () // `rsync -a` ahead of it. That is the shape that can see the difference: in the // copy phase the transfer itself would have set the timestamps. -test("out @ swap: a directory timestamp is settled by one more pass, not refused", async () => { +test("out @ swap: a directory timestamp is settled by one more pass, not refused", { skip: T1_SKIP }, async () => { await withTmp(async (paths, root) => { const channelDir = await seed(paths, "alpha", { v1: { "audio.m4a": "one".repeat(500), "transcript.json": "{}" }, @@ -1067,7 +1072,7 @@ test("out @ swap: a directory timestamp is settled by one more pass, not refused // still the live directory, and one change gets one more mirror pass, exactly // as in the copy phase. A second change is still a refusal ("a verify failure // keeps the source …" above). -test("out @ swap: a file the target is missing is mirrored by one more pass", async () => { +test("out @ swap: a file the target is missing is mirrored by one more pass", { skip: T1_SKIP }, async () => { await withTmp(async (paths, root) => { const channelDir = await seed(paths, "alpha", { v1: { "audio.m4a": "one".repeat(500) }, @@ -1242,7 +1247,7 @@ test("a writer seen once the marker is written refuses, and a fresh move's marke // on the destination that has since gone from the source. The resume used to // copy everything else and refuse on the counts (1755 against 1750), and no // rerun could settle it; the mirror pass now does. -test("a resume with a stale extra dir on the destination completes", async () => { +test("a resume with a stale extra dir on the destination completes", { skip: T1_SKIP }, async () => { await withTmp(async (paths, root) => { const channelDir = await seed(paths, "alpha", { v50t5yt: { "audio.mp3": "a".repeat(64), "transcript.json": "{}" }, @@ -1313,7 +1318,7 @@ test("back: a stale extra on the copy coming home is removed on resume", async ( // RECONCILE AND RESUME — the remediation (the ruling's last bullet). An extra // file and a changed one on the destination: the job says what it found, by // kind, makes the copy match the source and finishes the move. -test("reconcile: an extra and a changed file on the destination are settled, and the move completes", async () => { +test("reconcile: an extra and a changed file on the destination are settled, and the move completes", { skip: T1_SKIP }, async () => { await withTmp(async (paths, root) => { const channelDir = await seed(paths, "alpha", { v1: { "audio.m4a": "one".repeat(100), "transcript.json": '{"v":2}' }, @@ -1355,7 +1360,7 @@ test("reconcile: an extra and a changed file on the destination are settled, and // A marker past the copy phase has nothing to reconcile: the run is a plain // resume, and says so (the review's L3). -test("reconcile: a marker past the copy phase resumes, and does not claim a reconcile", async () => { +test("reconcile: a marker past the copy phase resumes, and does not claim a reconcile", { skip: T1_SKIP }, async () => { await withTmp(async (paths, root) => { const channelDir = await seed(paths, "alpha", { v1: { "audio.m4a": "one".repeat(100) }, @@ -1412,6 +1417,7 @@ test("reconcile: with no marker there is nothing to reconcile", async () => { const FAKE_FINDMNT = `#!/usr/bin/env node import { readFileSync } from "node:fs"; import path from "node:path"; + const control = JSON.parse( readFileSync(path.join(import.meta.dirname, "control.json"), "utf8"), ); diff --git a/common/controller/removeWrongFormatAudio.ts b/common/controller/removeWrongFormatAudio.ts @@ -1,3 +1,4 @@ +import { removeMediaFile } from "../lib/mediaTier-server"; import path from "node:path"; import fs from "fs-extra"; import type { Paths } from "../lib/paths"; @@ -5,7 +6,7 @@ import { isDoNotClean } from "../lib/doNotClean-server"; import { audioFilesToRemove } from "../lib/videoStatus"; import { readChannelConfig } from "./channels"; -const { pathExists, readdir, remove } = fs; +const { pathExists, readdir } = fs; export type RemoveWrongFormatAudioOptions = { channelSlug: string; @@ -70,7 +71,7 @@ export async function removeWrongFormatAudio({ continue; } for (const f of wrongFormat) { - await remove(path.join(videoDir, f)); + await removeMediaFile(videoDir, f); // derefs a tiered link (release 17) log(`Removed ${id}/${f}`); removedFiles++; } diff --git a/common/controller/renameChannel.test.ts b/common/controller/renameChannel.test.ts @@ -37,6 +37,12 @@ import { RELOCATION_MARKER_FILENAME, } from "../lib/channelMedia"; +// RELEASE 17 SLICE T1 made a channel whose `data/` is a link (or whose config +// carries `dataDir`) `legacy`; these cases still build that retired layout and +// expect it to read `ok`. Slice T2 rebases them on `media/` and un-skips them. +const T1_SKIP = "release 17 T2 rebases the mover on media/"; + + // Run with: // pnpm --filter yt-dlp-transcript-common exec tsx --test controller/renameChannel.test.ts @@ -150,7 +156,7 @@ test("renameChannel rejects invalid, same, and existing targets", async () => { // --- relocated media ------------------------------------------------------ -test("rename re-points a convention-shaped relocated media dir", async () => { +test("rename re-points a convention-shaped relocated media dir", { skip: T1_SKIP }, async () => { await withPaths(async (paths) => { const dir = path.dirname(paths.channelsDir); const mediaRoot = path.join(dir, "platter"); diff --git a/common/controller/storageLocations.test.ts b/common/controller/storageLocations.test.ts @@ -28,6 +28,12 @@ import { } from "./storageLocations"; import { readChannelConfig } from "./channels"; +// RELEASE 17 SLICE T1 made a channel whose `data/` is a link (or whose config +// carries `dataDir`) `legacy`; these cases still build that retired layout and +// expect it to read `ok`. Slice T2 rebases them on `media/` and un-skips them. +const T1_SKIP = "release 17 T2 rebases the mover on media/"; + + // Run with: // pnpm --filter yt-dlp-transcript-common exec tsx --test controller/storageLocations.test.ts // @@ -145,7 +151,7 @@ function loc(id: string, root: string, extra: Partial<StorageLocation> = {}): St return { id, label: id, root, autoRepoint: false, ...extra }; } -test("channelsOnLocation buckets ok / unreachable / moving and ignores channels elsewhere", async () => { +test("channelsOnLocation buckets ok / unreachable / moving and ignores channels elsewhere", { skip: T1_SKIP }, async () => { await withTmp(async (h) => { await seedRelocated(h, "alpha", h.rootA); // Media dir never created: the link dangles, which is what an unmounted @@ -180,7 +186,7 @@ test("channelsOnLocation buckets ok / unreachable / moving and ignores channels }); }); -test("re-point rewrites both channels' links and configs, then the location", async () => { +test("re-point rewrites both channels' links and configs, then the location", { skip: T1_SKIP }, async () => { await withTmp(async (h) => { await seedRelocated(h, "alpha", h.rootA); await seedRelocated(h, "beta", h.rootA); @@ -227,7 +233,7 @@ test("re-point rewrites both channels' links and configs, then the location", as }); }); -test("re-point refuses a target that has no media for a channel, naming the slug", async () => { +test("re-point refuses a target that has no media for a channel, naming the slug", { skip: T1_SKIP }, async () => { await withTmp(async (h) => { await seedRelocated(h, "alpha", h.rootA); await seedRelocated(h, "beta", h.rootA); @@ -263,7 +269,7 @@ test("re-point refuses a target that has no media for a channel, naming the slug }); }); -test("a failure on the second channel rolls the first one back", async () => { +test("a failure on the second channel rolls the first one back", { skip: T1_SKIP }, async () => { await withTmp(async (h) => { await seedRelocated(h, "alpha", h.rootA); await seedRelocated(h, "beta", h.rootA); @@ -306,7 +312,7 @@ test("a failure on the second channel rolls the first one back", async () => { }); }); -test("re-point refuses a busy channel and names it", async () => { +test("re-point refuses a busy channel and names it", { skip: T1_SKIP }, async () => { await withTmp(async (h) => { await seedRelocated(h, "alpha", h.rootA); await mkdir(path.join(h.rootB, "alpha", "data"), { recursive: true }); @@ -324,7 +330,7 @@ test("re-point refuses a busy channel and names it", async () => { }); }); -test("a rerun after a crash finishes the channels that were left", async () => { +test("a rerun after a crash finishes the channels that were left", { skip: T1_SKIP }, async () => { await withTmp(async (h) => { await seedRelocated(h, "alpha", h.rootA); await seedRelocated(h, "beta", h.rootA); @@ -375,7 +381,7 @@ test("a rerun after a crash finishes the channels that were left", async () => { }); }); -test("a channel killed between its symlink and its config write is resumed, not refused", async () => { +test("a channel killed between its symlink and its config write is resumed, not refused", { skip: T1_SKIP }, async () => { await withTmp(async (h) => { await seedRelocated(h, "alpha", h.rootA); await seedRelocated(h, "beta", h.rootA); diff --git a/common/controller/storageStall.test.ts b/common/controller/storageStall.test.ts @@ -2,6 +2,14 @@ // storage location's drive in-process asks the health state first // (lib/storageHealth.ts) and, on a stalled location, answers WITHOUT the call. // +// RELEASE 17, THE MEDIA TIER: a relocated channel's TEXT is on the corpus disk +// and only its media is on the drive (`channels/<slug>/media` -> the drive, +// one relative link per big file in `data/<id>/`). So a stalled drive holds +// what opens a big file and nothing that reads text: the channel's counts, +// its recency, its snapshot are read as usual, with no call on the drive. The +// cases that need a channel whose TEXT is on the drive use the retired +// whole-directory layout (`seedLegacy`), the one layout where it still is. +// // No test stalls a real drive. The drive here is an ordinary temp directory // that answers every call at once; the health state is TOLD it is stalled. // Every node:fs and node:fs/promises call this file's code makes is recorded @@ -19,6 +27,8 @@ import { existsSync, mkdirSync, mkdtempSync, + readFileSync, + renameSync, rmSync, symlinkSync, writeFileSync, @@ -31,6 +41,7 @@ import type { StorageLocation } from "../lib/storageLocations"; import { assertChannelMediaReachable, ChannelMediaUnreachableError, + ChannelTextUnreadableError, CHANNEL_MEDIA_MEMO_MS, clearRelocationMarker, forgetChannelMedia, @@ -128,30 +139,45 @@ const LOC: StorageLocation = { root: DRIVE, autoRepoint: false, }; -const TARGET = path.join(DRIVE, SLUG, "data"); -const LINK = path.join(paths.channelsDir, SLUG, "data"); -const CONFIG = { dataDir: TARGET }; +// The media tier on the drive, and its one link: channels/<slug>/media. +const TARGET = path.join(DRIVE, SLUG, "media"); +const LINK = path.join(paths.channelsDir, SLUG, "media"); +// The text, on the corpus disk; vid1's audio is a tiered link into media/. +const DATA = path.join(paths.channelsDir, SLUG, "data"); +const TIERED = path.join(DATA, "vid1", "audio.mp3"); +const CONFIG = { mediaDir: TARGET }; +// The retired layout: data/ itself a link to the drive. +const LEGACY_TARGET = path.join(DRIVE, SLUG, "data"); +const LEGACY_CONFIG = { dataDir: LEGACY_TARGET }; +let legacy = false; const FINDMNT_MARK = path.join(ROOT, "findmnt-ran"); -function seed(): void { - rmSync(CORPUS, { recursive: true, force: true }); - rmSync(DRIVE, { recursive: true, force: true }); - mkdirSync(path.join(TARGET, "vid1"), { recursive: true }); - writeFileSync( - path.join(TARGET, "vid1", "metadata.info.json"), - JSON.stringify({ id: "vid1", upload_date: "20260601" }), - ); - writeFileSync(path.join(TARGET, "vid1", "transcript.en.vtt"), "WEBVTT\n"); - mkdirSync(path.join(paths.channelsDir, SLUG), { recursive: true }); +function writeConfig(extra: Record<string, unknown>): void { writeFileSync( path.join(paths.channelsDir, SLUG, "config.json"), JSON.stringify({ handling: "youtube", name: SLUG, url: `https://www.youtube.com/@${SLUG}/videos`, - dataDir: TARGET, + ...extra, }), ); +} + +function seed(): void { + legacy = false; + rmSync(CORPUS, { recursive: true, force: true }); + rmSync(DRIVE, { recursive: true, force: true }); + mkdirSync(path.join(DATA, "vid1"), { recursive: true }); + writeFileSync( + path.join(DATA, "vid1", "metadata.info.json"), + JSON.stringify({ id: "vid1", upload_date: "20260601" }), + ); + writeFileSync(path.join(DATA, "vid1", "transcript.en.vtt"), "WEBVTT\n"); + mkdirSync(path.join(TARGET, "vid1"), { recursive: true }); + writeFileSync(path.join(TARGET, "vid1", "audio.mp3"), "AUDIO"); + symlinkSync(path.join("..", "..", "media", "vid1", "audio.mp3"), TIERED); + writeConfig({ mediaDir: TARGET }); symlinkSync(TARGET, LINK); writeFileSync(paths.findmntBin, `#!/bin/sh\ntouch ${FINDMNT_MARK}\nexit 1\n`, { mode: 0o755, @@ -159,12 +185,29 @@ function seed(): void { rmSync(FINDMNT_MARK, { force: true }); } -// Anything that reaches the drive: a path under its root, or through the -// channel's link (`data/` itself, followed, or anything below it). +// The channel on the retired layout: its whole data/ on the drive, `data` an +// absolute link to it, `dataDir` recorded (and no media tier). +function seedLegacy(): void { + rmSync(LINK); + rmSync(TIERED); + writeFileSync(path.join(DATA, "vid1", "audio.mp3"), "AUDIO"); + renameSync(DATA, LEGACY_TARGET); + symlinkSync(LEGACY_TARGET, DATA); + writeConfig({ dataDir: LEGACY_TARGET }); + legacy = true; +} + +// Anything that reaches the drive: a path under its root; through the media +// link; a tiered file opened or statted through its link (an lstat or a +// readlink of the link itself is on the corpus disk); and, on the legacy +// layout, anything through `data/`. function onDrive(c: Call): boolean { const under = (p: string, base: string) => p === base || p.startsWith(base + path.sep); - return under(c.path, DRIVE) || under(c.path, LINK); + if (under(c.path, DRIVE) || under(c.path, LINK)) return true; + const ofTheLink = ["lstat", "lstatSync", "readlink", "readlinkSync"].includes(c.fn); + if (c.path === TIERED && !ofTheLink) return true; + return legacy && under(c.path, DATA) && !(c.path === DATA && ofTheLink); } function stall(): void { @@ -190,13 +233,18 @@ test("inspect with the config in hand: on a stalled drive only the marker is rea stall(); calls = []; const media = await inspectChannelMedia(paths, SLUG, CONFIG); - assert.deepEqual(calls, [{ fn: "readFile", path: MARKER }]); + assert.deepEqual(calls, [ + { fn: "lstat", path: DATA }, + { fn: "readFile", path: MARKER }, + ]); assert.deepEqual(calls.filter(onDrive), []); assert.equal(media.status, "stalled"); assert.equal(media.target, TARGET); assert.match(String(media.detail), /^drive not answering \(location "USB drive", since /); - // Held by both pool-wide builds, with a reason that names no path. + // Held for the media, with a reason that names no path — and its text is + // readable, so the builds do not hold it (release 17). assert.equal(isMediaHeld(media.status), true); + assert.equal(media.text.readable, true); assert.doesNotMatch(HELD_REASON.stalled, /\//); }); @@ -209,6 +257,7 @@ test("inspect without the config reads config.json and nothing on the drive", as calls.map((c) => [c.fn, path.relative(ROOT, c.path)]), [ ["readFile", path.join("corpus", "channels", SLUG, "config.json")], + ["lstat", path.join("corpus", "channels", SLUG, "data")], ["readFile", path.join("corpus", "channels", SLUG, ".relocating.json")], ], ); @@ -268,7 +317,7 @@ test("the memo is keyed by the configured target, and a mover's forget clears it calls.filter((c) => c.fn === "stat" && c.path === TARGET).length; await inspectChannelMedia(paths, SLUG, CONFIG, { now }); // Another configured target is another key. - const other = await inspectChannelMedia(paths, SLUG, { dataDir: path.join(DRIVE, "x", "data") }, { now }); + const other = await inspectChannelMedia(paths, SLUG, { mediaDir: path.join(DRIVE, "x", "media") }, { now }); assert.equal(other.status, "inconsistent"); forgetChannelMedia(SLUG); await inspectChannelMedia(paths, SLUG, CONFIG, { now }); @@ -316,7 +365,20 @@ test("volumeFreeBytes: a stalled location reads unknown, with no call on it", as assert.deepEqual(calls.filter(onDrive), []); }); -test("readChannelStat: no walk of data/ on a stalled drive", async () => { +test("readChannelStat: a stalled MEDIA drive does not stop the count, and is not asked", async () => { + const before = await readChannelStat(paths, SLUG); + assert.equal(before?.videoCount, 1); + assert.equal(before?.downloadCount, 1); + stall(); + calls = []; + const during = await readChannelStat(paths, SLUG); + assert.equal(during?.videoCount, 1); + assert.equal(during?.downloadCount, 1, "the tiered audio is a name in the listing"); + assert.deepEqual(calls.filter(onDrive), []); +}); + +test("readChannelStat: no walk of a LEGACY channel's data/ on a stalled drive", async () => { + seedLegacy(); const before = await readChannelStat(paths, SLUG); assert.equal(before?.videoCount, 1); stall(); @@ -325,9 +387,26 @@ test("readChannelStat: no walk of data/ on a stalled drive", async () => { assert.deepEqual(calls.filter(onDrive), []); }); -test("recency: no tail read on a stalled drive, and no miss remembered for it", async () => { +test("recency: a stalled MEDIA drive does not stop the tail read (metadata is text)", async () => { + const args = { + paths, + meta: [{ slug: SLUG, config: CONFIG }], + candidateIds: new Set(["vid1"]), + owner: new Map([["vid1", SLUG]]), + interpolate: false, + fresh: true, + }; + stall(); + calls = []; + const keys = await buildRecencyKeys(args); + assert.deepEqual(keys.get("vid1"), { key: "20260601", estimated: false }); + assert.deepEqual(calls.filter(onDrive), []); +}); + +test("recency: no tail read on a LEGACY channel's stalled drive, and no miss remembered for it", async () => { + seedLegacy(); const owner = new Map([["vid1", SLUG]]); - const meta = [{ slug: SLUG, config: CONFIG }]; + const meta = [{ slug: SLUG, config: LEGACY_CONFIG }]; const args = { paths, meta, @@ -363,13 +442,32 @@ test("a move onto a stalled location is refused without a stat of its root", asy assert.deepEqual(calls.filter(onDrive), []); }); -test("a snapshot refresh of a stalled channel throws before its walk", async () => { +test("a snapshot refresh of a stalled-MEDIA channel is written from its text; its media bytes are unknown", async () => { + const ok = await generateChannelSnapshot(paths, SLUG); + assert.equal(ok.totalMediaBytes, 5, "the tiered audio, statted through its link"); + assert.equal(typeof ok.totalTextBytes, "number"); + stall(); + calls = []; + const snap = await generateChannelSnapshot(paths, SLUG); + assert.deepEqual(calls.filter(onDrive), []); + assert.equal(snap.totals.downloaded, 1); + assert.equal("totalMediaBytes" in snap, false, "unknown, never 0"); + assert.equal("totalAudioBytes" in snap, false); + assert.equal(typeof snap.totalTextBytes, "number"); + const onDisk = JSON.parse( + readFileSync(path.join(paths.channelsDir, SLUG, "snapshot.json"), "utf8"), + ) as { totalMediaBytes?: number }; + assert.equal(onDisk.totalMediaBytes, undefined); +}); + +test("a snapshot refresh of a LEGACY channel throws before its walk", async () => { + seedLegacy(); stall(); calls = []; await assert.rejects( () => generateChannelSnapshot(paths, SLUG), (err: unknown) => - err instanceof ChannelMediaUnreachableError && err.status === "stalled", + err instanceof ChannelTextUnreadableError && err.status === "legacy", ); assert.deepEqual(calls.filter(onDrive), []); }); @@ -395,13 +493,15 @@ test("the saved-video store on a stalled drive reads unreachable, without a stat assert.deepEqual(followed, []); }); -test("the saved-video inventory, for a page: a stalled channel is named, not read", async () => { +test("the saved-video inventory, for a page: a stalled LEGACY channel is named, not read", async () => { // A saved video on the drive, so a channel that WAS read has an entry. // (savedVideoInventory reads through fs-extra, whose functions graceful-fs // captured before this file's spy was installed, so this case proves the skip - // by what comes back, not by the spy.) + // by what comes back, not by the spy.) The pointer is text: only on the + // retired layout is it on the drive at all. + seedLegacy(); writeFileSync( - path.join(TARGET, "vid1", "saved-video.json"), + path.join(LEGACY_TARGET, "vid1", "saved-video.json"), JSON.stringify({ storedAt: "", dir: "/store", file: "v.mp4", bytes: 7 }), ); assert.equal((await listSavedVideos({ paths, notAnswering: [] })).length, 1); @@ -444,11 +544,12 @@ test("watchdog: inspect's target stat never answers → stalled, marked, and not )) as { status: string; detail?: string }; assert.equal(media.status, "stalled"); assert.match(String(media.detail), /^drive not answering \(location "USB drive"/); - // The next inspect, the guard and the walk make no call on the drive. + // The next inspect, the guard and the walk make no call on the drive — and + // the walk, which reads text, still counts (release 17). calls = []; assert.equal((await inspectChannelMedia(paths, SLUG, CONFIG)).status, "stalled"); await assert.rejects(() => assertChannelMediaReachable(paths, SLUG, CONFIG)); - assert.equal(await readChannelStat(paths, SLUG), null); + assert.equal((await readChannelStat(paths, SLUG))?.videoCount, 1); assert.deepEqual(calls.filter(onDrive), []); // Cleared (two clean answers), the drive is asked again. recordLocationHealth(LOC, "ok"); @@ -460,20 +561,21 @@ test("watchdog: inspect's target stat never answers → stalled, marked, and not ); }); -test("watchdog: a video directory's read never answers mid-walk → readChannelStat answers null", async () => { +test("watchdog: a LEGACY video directory's read never answers mid-walk → readChannelStat answers null", async () => { + seedLegacy(); for (const id of ["vid2", "vid3", "vid4", "vid5", "vid6"]) { - mkdirSync(path.join(TARGET, id), { recursive: true }); + mkdirSync(path.join(LEGACY_TARGET, id), { recursive: true }); } const out = await watchdogCase(["readdir"], async () => { // The walk's first readdir (of data/ itself, through the link) answers; // the video directories' do not. hang = (c) => - c.fn === "readdir" && c.path.startsWith(path.join(LINK, "vid")); + c.fn === "readdir" && c.path.startsWith(path.join(DATA, "vid")); return readChannelStat(paths, SLUG); }); assert.equal(out, null); // At most four video directories were asked before the stall refused the rest. - const asked = calls.filter((c) => c.fn === "readdir" && c.path.startsWith(path.join(LINK, "vid"))); + const asked = calls.filter((c) => c.fn === "readdir" && c.path.startsWith(path.join(DATA, "vid"))); assert.ok(asked.length <= 4, `${asked.length} video dirs asked`); }); @@ -491,10 +593,11 @@ test("watchdog: volumeFreeBytes' stat never answers → unknown", async () => { assert.equal(out.usb, undefined); }); -test("watchdog: a recency tail read never answers → not dated, not remembered as a miss", async () => { +test("watchdog: a LEGACY recency tail read never answers → not dated, not remembered as a miss", async () => { + seedLegacy(); const args = { paths, - meta: [{ slug: SLUG, config: CONFIG }], + meta: [{ slug: SLUG, config: LEGACY_CONFIG }], candidateIds: new Set(["vid1"]), owner: new Map([["vid1", SLUG]]), interpolate: false, @@ -521,24 +624,26 @@ test("watchdog: a move onto a root whose stat never answers is refused", async ( assert.match(String(problem), /drive not answering/); }); -test("watchdog (M3): the snapshot walk's video unit never answers → the refresh throws and writes no snapshot", async () => { +test("watchdog (M3, release 17): a tiered file's stat never answers → the snapshot is written, its media bytes unknown", async () => { + // Six more videos, each with a tiered audio link, so the walk has units + // past the first to refuse. for (const id of ["vid2", "vid3", "vid4", "vid5", "vid6", "vid7"]) { + mkdirSync(path.join(DATA, id), { recursive: true }); + writeFileSync(path.join(DATA, id, "metadata.info.json"), JSON.stringify({ id })); mkdirSync(path.join(TARGET, id), { recursive: true }); + writeFileSync(path.join(TARGET, id, "audio.mp3"), "AUDIO"); + symlinkSync(path.join("..", "..", "media", id, "audio.mp3"), path.join(DATA, id, "audio.mp3")); } - const snapshotFile = path.join(paths.channelsDir, SLUG, "snapshot.json"); - await watchdogCase(["readdir"], async () => { - // data/ itself answers (the listing, the reconcile pass); the video - // directories do not. - hang = (c) => c.fn === "readdir" && c.path.startsWith(path.join(LINK, "vid")); - await assert.rejects( - () => generateChannelSnapshot(paths, SLUG), - (err: unknown) => err instanceof Error && err.name === "DriveNotAnsweringError", - ); - }); - assert.equal(existsSync(snapshotFile), false, "the last snapshot.json stands"); - // At most four video directories reached the drive. - const asked = calls.filter( - (c) => c.fn === "readdir" && c.path.startsWith(path.join(LINK, "vid")), - ); - assert.ok(asked.length <= 4, `${asked.length} video dirs asked`); + const tiered = (c: Call) => + c.fn === "stat" && c.path.startsWith(DATA + path.sep) && c.path.endsWith(`${path.sep}audio.mp3`); + const snap = (await watchdogCase(["stat"], async () => { + // The text answers; a stat through a tiered link does not. + hang = tiered; + return generateChannelSnapshot(paths, SLUG); + })) as Awaited<ReturnType<typeof generateChannelSnapshot>>; + assert.equal(snap.totals.downloaded, 7); + assert.equal("totalMediaBytes" in snap, false, "unknown, never a partial sum"); + // At most four video units reached the drive. + const asked = new Set(calls.filter(tiered).map((c) => path.dirname(c.path))); + assert.ok(asked.size <= 4, `${asked.size} video units asked`); }); diff --git a/common/controller/storageWatch.test.ts b/common/controller/storageWatch.test.ts @@ -30,6 +30,12 @@ import { type LocationHealthState, } from "../lib/storageHealth"; +// RELEASE 17 SLICE T1 made a channel whose `data/` is a link (or whose config +// carries `dataDir`) `legacy`; these cases still build that retired layout and +// expect it to read `ok`. Slice T2 rebases them on `media/` and un-skips them. +const T1_SKIP = "release 17 T2 rebases the storage watch on mediaDir"; + + // THE CONFIRMATION COUNT IS MODULE STATE (see storageWatch.ts rule 3), so each // case starts from a clean one — otherwise the second test inherits the first // test's suspicions and pauses on what should be its first pass. The health @@ -149,7 +155,7 @@ async function twoPasses(h: H) { return { first, second }; } -test("a channel whose target is gone is auto-paused, once, in one write", async () => { +test("a channel whose target is gone is auto-paused, once, in one write", { skip: T1_SKIP }, async () => { await withTmp(async (h) => { await seedRelocated(h, "gone-a", { targetExists: false }); await seedRelocated(h, "gone-b", { targetExists: false }); @@ -189,7 +195,7 @@ test("a channel whose target is gone is auto-paused, once, in one write", async }); }); -test("the drive coming back restores the tier it overwrote", async () => { +test("the drive coming back restores the tier it overwrote", { skip: T1_SKIP }, async () => { await withTmp(async (h) => { await seedRelocated(h, "away", { targetExists: false }); const settings = h.io.read(); @@ -295,7 +301,7 @@ test("a record on an in-place channel is restored", async () => { // `write: false` IS IDLE BOOT. It observes and reports; the write is the work, // and idle boot refuses work. -test("write: false reports the transition and changes nothing", async () => { +test("write: false reports the transition and changes nothing", { skip: T1_SKIP }, async () => { await withTmp(async (h) => { await seedRelocated(h, "gone", { targetExists: false }); // The first pass only suspects, whatever `write` says. @@ -334,7 +340,7 @@ test("no locations and nothing auto-paused is a free pass", async () => { // has spun down and needs a beat to answer is indistinguishable from "not // mounted" — and pausing on it rewrites the corpus's priority document for a // drive that is fine. -test("a drive that blips for one pass is never paused", async () => { +test("a drive that blips for one pass is never paused", { skip: T1_SKIP }, async () => { await withTmp(async (h) => { await seedRelocated(h, "blip", { targetExists: false }); const first = await runStorageWatchPass({ @@ -379,7 +385,7 @@ test("a drive that blips for one pass is never paused", async () => { // RESTORE STAYS SINGLE-PASS, and the asymmetry is the point: being slow to // pause costs a few refused units (the start-of-work guards catch those), while // being slow to restore leaves a lane off after the operator fixed the cable. -test("the restore needs only one good pass", async () => { +test("the restore needs only one good pass", { skip: T1_SKIP }, async () => { await withTmp(async (h) => { await seedRelocated(h, "back", { targetExists: false }); await twoPasses(h); @@ -409,7 +415,7 @@ function scripted(answers: LocationHealthState[]) { return async () => answers[Math.min(i++, answers.length - 1)]; } -test("one missed probe stalls the location; pages then answer 'stalled' without asking", async () => { +test("one missed probe stalls the location; pages then answer 'stalled' without asking", { skip: T1_SKIP }, async () => { await withTmp(async (h) => { await seedRelocated(h, "slow", { targetExists: true }); const lines: string[] = []; @@ -434,7 +440,7 @@ test("one missed probe stalls the location; pages then answer 'stalled' without }); }); -test("the stall clears only after two clean probes in a row", async () => { +test("the stall clears only after two clean probes in a row", { skip: T1_SKIP }, async () => { await withTmp(async (h) => { await seedRelocated(h, "slow", { targetExists: true }); const probe = scripted(["stalled", "ok", "stalled", "ok", "ok"]); diff --git a/common/controller/transcode.ts b/common/controller/transcode.ts @@ -3,6 +3,7 @@ import { rename, rm } from "node:fs/promises"; import { execa } from "execa"; import type { Paths } from "../lib/paths"; import type { AudioFormat } from "../lib/channelConfig"; +import { tierMediaFile } from "../lib/mediaTier-server"; const CODEC_ARGS: Record<AudioFormat, string[]> = { m4a: ["-c:a", "aac"], @@ -54,4 +55,11 @@ export async function transcodeAudio(opts: TranscodeAudioOptions): Promise<void> } await rename(tmp, out); opts.onLog(`Wrote ${out}\n`); + // THE MEDIA TIER'S HOOK (release 17). The rename above put a real file over + // the name — over a tiered LINK, when the channel had one — so the fresh + // file moves into channels/<slug>/media (over the stale copy there) and the + // link is made again. Every finalising transcode is this function: the app + // extraction, the audio-checked download, the video page's Transcode. Never + // throws; on a classic channel it leaves the file where it is. + await tierMediaFile(opts.videoDir, path.basename(out), { onLog: opts.onLog }); } diff --git a/common/jobs/jobKinds.test.ts b/common/jobs/jobKinds.test.ts @@ -1,6 +1,12 @@ import { test } from "node:test"; import assert from "node:assert/strict"; -import { isDrainableKind, jobKindLabel, getJobKind } from "./jobKinds"; +import { + isDrainableKind, + jobKindLabel, + getJobKind, + kindNeedsMedia, + kindNeedsText, +} from "./jobKinds"; // Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/jobKinds.test.ts // @@ -120,3 +126,72 @@ test("label-less and unknown kinds fall back to the raw kind", () => { assert.equal(jobKindLabel("totally-unknown"), "totally-unknown"); assert.equal(getJobKind("totally-unknown"), undefined); }); + + +// RELEASE 17: `needsMedia` means "opens or writes the BIG file". These kinds +// read only the text tier and flipped to `needsText`; the rest stay media. +const TEXT_KINDS = [ + "auto-digest", + "digest-channel-local", + "digest-channel-remote", + "digest-share-cluster", + "normalize-transcripts", + "purge-superseded-auto-subs", + "fetch-window", + "evict-clips", + "metadata-scan", + "download-missing-subs", + "check-availability", + "quick-availability-check", + "check-maybe-missing", + "check-kept-deleted", +]; +const STILL_MEDIA = [ + "auto-transcribe", + "auto-download", + "auto-download-unit", + "auto-backfill", + "whisper-all", + "whisper-bucket-downloaded-no-transcript", + "whisper-bucket-auto-subs", + "sync", + "import-one", + "download-from-playlist", + "download-missing", + "redownload-archive", + "redownload-incomplete-bucket", + "retry-bucket", + "persist-kept", + "whisper-video", + "transcribe-one", + "download-one-pipeline", + "transcode-audio", + "diarize-channel", + "backfill-channel", + "scan-media", + "scan-media-channel", + "clean-audio-transcribed", + "clean-extra-audio-formats", + "remove-wrong-format-audio", +]; + +test("release 17: the text kinds flipped from needsMedia to needsText", () => { + for (const k of TEXT_KINDS) { + assert.ok(getJobKind(k), `${k} is registered`); + assert.equal(kindNeedsMedia(k), false, `${k} does not open a big file`); + assert.equal(kindNeedsText(k), true, `${k} reads the text tier`); + } +}); + +test("release 17: every media kind keeps needsMedia, and none is also a text kind", () => { + for (const k of STILL_MEDIA) { + assert.ok(getJobKind(k), `${k} is registered`); + assert.equal(kindNeedsMedia(k), true, k); + assert.equal(kindNeedsText(k), false, k); + } + // The movers fix an unreachable channel and declare neither. + for (const k of ["relocate-channel-media", "relocate-saved-videos", "repoint-storage-location"]) { + assert.equal(kindNeedsMedia(k), false, k); + assert.equal(kindNeedsText(k), false, k); + } +}); diff --git a/common/jobs/jobKinds.ts b/common/jobs/jobKinds.ts @@ -37,12 +37,14 @@ export type JobKindMeta = { // Fallback scheduler tier when a record is not explicitly background. Left // undefined today so Phase 4 maps it to "foreground" (unchanged behavior). defaultTier?: SchedulerTier; - // Whether this kind reads or writes files under `channels/<slug>/data/`. - // Declarative, and consumed by exactly one thing: runManagedFunction refuses - // to enqueue a media kind for a channel whose media is not reachable (a - // relocated channel whose drive is unmounted, or one mid-relocation) rather - // than letting it read an empty dir as the truth. See - // common/lib/channelMedia.ts. + // Whether this kind OPENS OR WRITES A BIG FILE — the audio, a persisted + // source container, the raw live-chat replay (lib/mediaTier.ts says which): + // the files that may live on another drive (release 17, the media tier). + // Declarative, and consumed by runManagedFunction, which refuses a media kind + // for a channel whose media is not reachable (a relocated channel whose drive + // is unmounted or stalled, one mid-move, a `legacy` one) — at enqueue and + // again when its queue starts it — rather than letting it read an empty dir + // as the truth or write into a tree being copied. See common/lib/channelMedia.ts. // // ABSENT MEANS FALSE, and that is deliberate rather than lazy. The guard is // opt-in so a kind that is not listed keeps exactly its current behavior, and @@ -51,13 +53,24 @@ export type JobKindMeta = { // that genuinely never opens a video dir (store playlist, clear markers) must // not be refused for a drive it does not read. // - // THE TEST IS "DOES IT OPEN data/", NOT "IS IT BOOKKEEPING". Two kinds that - // read as bookkeeping declare it anyway, and both were misses: - // `normalize-transcripts` walks every video dir and writes a sidecar into - // each, and `sync` reads the data dir to decide what is already downloaded - // before downloading into it. Against an unmounted drive the first reports a - // clean run over zero videos and the second concludes nothing is downloaded. + // THE TEST IS "DOES IT OPEN OR WRITE THE BIG FILE", NOT "IS IT BOOKKEEPING" + // (the doctrine since release 17; before it, "does it open data/"). `sync` + // and the metadata scan's cousins that DOWNLOAD write media; a transcription + // reads it. A kind that reads only the text — a digest, normalize, the + // availability checks, the metadata scan, a clip-window fetch — declares + // `needsText` instead, and runs while the channel's media is moving, stalled + // or unmounted. needsMedia?: boolean; + // Whether this kind walks or writes the TEXT under `channels/<slug>/data/` + // (transcripts, cues, metadata, sidecars, `clips/`) and nothing else. The + // text guard (`assertChannelTextReadable`) refuses it only where the text + // itself cannot be read: a `legacy` channel (the retired whole-directory + // layout, its text on the far drive), a `data/` that is not a directory, a + // tier migration in flight. Against such a channel the walk would find + // nothing and report a clean run over zero videos — or, for the + // availability checks and the metadata scan, read every downloaded video as + // missing and re-request the whole channel. Absent means false. + needsText?: boolean; }; // One entry per kind known to the system. `label` is included only where the @@ -101,7 +114,8 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: true, replayable: false, queueKeyStrategy: "parallel", - needsMedia: true, + needsMedia: false, + needsText: true, }, "auto-backfill": { kind: "auto-backfill", @@ -149,7 +163,8 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: false, replayable: true, queueKeyStrategy: "custom", - needsMedia: true, + needsMedia: false, + needsText: true, }, // AI digest sweep, local (ollama) lane — the one that carries the corpus. Both // digest kinds are drainable (the batch honors the drain signal: it stops @@ -161,7 +176,8 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: true, replayable: true, queueKeyStrategy: "custom", - needsMedia: true, + needsMedia: false, + needsText: true, }, // Same batch, metered lane. Off unless settings.digest.remoteEnabled is true, // and it lands on its own queue key so it runs CONCURRENTLY with the local lane @@ -172,7 +188,8 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: true, replayable: true, queueKeyStrategy: "custom", - needsMedia: true, + needsMedia: false, + needsText: true, }, // Copy a duplicate cluster's canonical digest onto its aligned mirrors. A fast // file operation gated by the timestamp-alignment check, so it is not drainable @@ -183,7 +200,8 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: false, replayable: false, queueKeyStrategy: "parallel", - needsMedia: true, + needsMedia: false, + needsText: true, }, // Write the compact transcript.cues.json sidecar next to every raw transcript // that lacks a current one — corpus-wide from the Pool on /sites, or one @@ -196,16 +214,19 @@ const JOB_KINDS: Record<string, JobKindMeta> = { // write, so "let the in-flight one finish" is already how it behaves. Not // replayable either — that would need a JobSpec and a jobReplayRegistry // handler, and the button is one click from the card that reports the count. - // needsMedia: it walks `data/<id>/` for every video and writes a sidecar into - // each. Against an unmounted drive it finds nothing, reports a clean - // 0/0/0/0 run and moves on — a sweep that silently skips a channel. + // needsText (release 17; needsMedia before it): it walks `data/<id>/` for + // every video and writes a sidecar into each — text, on the corpus disk. + // Against an unreadable text tier (a `legacy` channel whose drive is + // unmounted) it finds nothing, reports a clean 0/0/0/0 run and moves on — a + // sweep that silently skips a channel. A stalled MEDIA drive does not stop it. "normalize-transcripts": { kind: "normalize-transcripts", label: "Normalize transcripts", drainable: false, replayable: false, queueKeyStrategy: "custom", - needsMedia: true, + needsMedia: false, + needsText: true, }, "redownload-incomplete-bucket": { kind: "redownload-incomplete-bucket", @@ -237,7 +258,8 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: true, replayable: true, queueKeyStrategy: "platform", - needsMedia: true, + needsMedia: false, + needsText: true, }, "import-one": { kind: "import-one", @@ -258,7 +280,8 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: false, replayable: true, queueKeyStrategy: "platform", - needsMedia: true, + needsMedia: false, + needsText: true, }, "redownload-archive": { kind: "redownload-archive", @@ -342,7 +365,8 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: false, replayable: true, queueKeyStrategy: "custom", - needsMedia: true, + needsMedia: false, + needsText: true, }, "persist-kept": { kind: "persist-kept", @@ -404,17 +428,19 @@ const JOB_KINDS: Record<string, JobKindMeta> = { // the cleanup lanes are about `audio.*`. This is the only thing that removes // one. // - // `needsMedia: true`, and it is the sharpest case for the flag in the table: - // it DELETES. Against an unmounted drive every `readdir` of `data/` throws - // and the walk would report a clean eviction of zero bytes — the operator - // would read "nothing to reclaim" about a platter full of windows. + // `needsText` (release 17): `clips/` is on the corpus disk and is never + // tiered, so the media drive does not concern it. The guard still matters: + // it DELETES, and against an unreadable text tier (a `legacy` channel whose + // drive is unmounted) every `readdir` of `data/` throws and the walk would + // report a clean eviction of zero bytes. "evict-clips": { kind: "evict-clips", label: "Evict fetched windows", drainable: false, replayable: false, queueKeyStrategy: "parallel", - needsMedia: true, + needsMedia: false, + needsText: true, }, // THE SAVED-VIDEO STORE, ONTO A LOCATION AND BACK. Same mechanism as the // channel move (relocateDir.ts is literally the same code) over one directory @@ -499,19 +525,20 @@ const JOB_KINDS: Record<string, JobKindMeta> = { // because it contends for the same thing a download does — the source's // patience. // - // `needsMedia` IS TRUE, despite the scan never creating a video directory, - // and the test in JobKindMeta is exactly why: does it open `data/`? It does — - // its whole target set is "listed, minus what is already on disk". Against an - // unmounted drive it would read every downloaded video as unfetched and - // re-request the entire channel. It writes no media; that is a different - // question from whether it READS the media dir. + // `needsText` (release 17; `needsMedia` before it): its whole target set is + // "listed, minus what is already on disk", and "on disk" is the video dirs in + // `data/` — text. Against an unreadable text tier (a `legacy` channel whose + // drive is unmounted) it would read every downloaded video as unfetched and + // re-request the entire channel. It writes no media, so a stalled media + // drive does not hold it. "metadata-scan": { kind: "metadata-scan", label: "Metadata scan", drainable: true, replayable: true, queueKeyStrategy: "platform", - needsMedia: true, + needsMedia: false, + needsText: true, }, // THE HUB'S AND THE HOMEPAGE'S BUILD AND DEPLOY (release 13 slice W1). They // ran from /sites — the hub since release 7, the homepage since release 11 — @@ -607,21 +634,24 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: false, replayable: false, queueKeyStrategy: "platform", - needsMedia: true, + needsMedia: false, + needsText: true, }, "quick-availability-check": { kind: "quick-availability-check", drainable: false, replayable: false, queueKeyStrategy: "platform", - needsMedia: true, + needsMedia: false, + needsText: true, }, "check-maybe-missing": { kind: "check-maybe-missing", drainable: false, replayable: false, queueKeyStrategy: "platform", - needsMedia: true, + needsMedia: false, + needsText: true, }, // Replayable kinds that never had a JOB_KIND_LABELS entry: label omitted so // jobKindLabel() keeps falling back to the raw kind (unchanged behavior). @@ -667,9 +697,15 @@ export function isDrainableKind(kind: string): boolean { return JOB_KINDS[kind]?.drainable ?? false; } -// Whether a kind's work reaches `channels/<slug>/data/`. Absent = false: the -// media guard is opt-in, so an unlisted (or unknown) kind behaves exactly as it -// did before the guard existed. +// Whether a kind's work opens or writes a big file. Absent = false: the media +// guard is opt-in, so an unlisted (or unknown) kind behaves exactly as it did +// before the guard existed. export function kindNeedsMedia(kind: string): boolean { return JOB_KINDS[kind]?.needsMedia ?? false; } + +// Whether a kind reads the text tier only (and so asks the text guard instead +// of the media one). Absent = false. +export function kindNeedsText(kind: string): boolean { + return JOB_KINDS[kind]?.needsText ?? false; +} diff --git a/common/jobs/streamCommand.test.ts b/common/jobs/streamCommand.test.ts @@ -1,6 +1,6 @@ import { test } from "node:test"; import assert from "node:assert/strict"; -import { mkdir, mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; +import { mkdir, mkdtemp, readFile, rm, symlink, writeFile } from "node:fs/promises"; import { tmpdir } from "node:os"; import path from "node:path"; import type { Paths } from "../lib/paths"; @@ -217,6 +217,68 @@ test("a media job queued before a move's marker and started after it refuses at } }); +// RELEASE 17: a TEXT kind (`needsText`) asks the text guard. It runs while the +// channel's media drive is unmounted, and is refused — before any record — +// only where the text itself cannot be read: the retired layout. +test("a text job runs over an unmounted media drive and is refused on a legacy channel", async () => { + const { paths: base, root } = await jobsDir(); + const paths = { ...base, channelsDir: path.join(root, "channels") } as Paths; + const channelDir = path.join(paths.channelsDir, "alpha"); + await mkdir(path.join(channelDir, "data"), { recursive: true }); + const target = path.join(root, "platter", "alpha", "media"); + await writeFile( + path.join(channelDir, "config.json"), + JSON.stringify({ handling: "youtube", url: "https://example.com/a", mediaDir: target }), + ); + await symlink(target, path.join(channelDir, "media")); // not mounted + try { + let ran = 0; + const job = ok( + await runManagedFunction({ + kind: "normalize-transcripts", + queueKey: `test:text-guard:${newJobId()}`, + paths, + channelSlug: "alpha", + fn: async () => { + ran += 1; + }, + }), + ); + assert.equal((await job.done).status, "done"); + assert.equal(ran, 1); + // A media kind on the same channel is refused, before any record. + const media = await runManagedFunction({ + kind: "whisper-all", + queueKey: `test:text-guard:${newJobId()}`, + paths, + channelSlug: "alpha", + fn: async () => {}, + }); + assert.equal(media.ok, false); + assert.match((media as { error: string }).error, /does not exist \(drive not mounted\?\)/); + + // The retired layout: the text kind is refused too, naming the way out. + await writeFile( + path.join(channelDir, "config.json"), + JSON.stringify({ handling: "youtube", url: "https://example.com/a", dataDir: "/mnt/platter/alpha/data" }), + ); + const refused = await runManagedFunction({ + kind: "normalize-transcripts", + queueKey: `test:text-guard:${newJobId()}`, + paths, + channelSlug: "alpha", + fn: async () => { + ran += 1; + }, + }); + assert.equal(refused.ok, false); + assert.match((refused as { error: string }).error, /archilyzer storage migrate-tier alpha/); + assert.equal(ran, 1); + } finally { + await rm(root, { recursive: true, force: true }); + } +}); + // THE ONE CANCEL THAT MUST NOT: the graceful-shutdown reaper cancels every // queued job only so the exit cannot promote one into a child. Nobody cancelled // it, and its `queued` sidecar is what the boot pass (bootQueuedJobs.ts) diff --git a/common/jobs/streamCommand.ts b/common/jobs/streamCommand.ts @@ -20,9 +20,10 @@ import { import { writeJobMeta } from "./jobMeta"; import { maybePruneJobLogs } from "./listJobs"; import type { JobSpec } from "./jobSpec"; -import { kindNeedsMedia } from "./jobKinds"; +import { kindNeedsMedia, kindNeedsText } from "./jobKinds"; import { assertChannelMediaReachable, + assertChannelTextReadable, ChannelMediaUnreachableError, } from "../lib/channelMedia"; import { mediaHoldText } from "../lib/channelMediaHold"; @@ -324,10 +325,25 @@ export async function runManagedCommand( // // A MOVE IS A HOLD, and the refusal says so in the hold's words ("held: its // media is moving …") rather than calling the media unreachable. +// +// A TEXT KIND ASKS THE TEXT GUARD (release 17): a kind that declares +// `needsText` reads only `data/`'s text, which stays on the corpus disk, so it +// runs while the channel's media is moving, stalled or unmounted, and is +// refused only where the text itself cannot be read (a `legacy` channel, a +// `data/` that is not a directory, a tier migration in flight). async function refuseForUnreachableMedia( opts: CommonOpts, ): Promise<string | null> { - if (!opts.channelSlug || !kindNeedsMedia(opts.kind)) return null; + if (!opts.channelSlug) return null; + if (!kindNeedsMedia(opts.kind)) { + if (!kindNeedsText(opts.kind)) return null; + try { + await assertChannelTextReadable(opts.paths, opts.channelSlug); + return null; + } catch (err) { + return (err as Error).message; + } + } try { await assertChannelMediaReachable(opts.paths, opts.channelSlug); return null; diff --git a/common/lib/channelConfig.ts b/common/lib/channelConfig.ts @@ -126,6 +126,7 @@ export type ChannelConfig = { extractionMode?: ExtractionMode; savedVideosDir?: string; dataDir?: string; + mediaDir?: string; ytdlpExtraArgs?: string[]; subLangs?: string; lastSyncedAt?: string; @@ -167,7 +168,9 @@ export const CHANNEL_CONFIG_FIELD_DOCS: FieldDocs<ChannelConfig> = { savedVideosDir: "Per-channel override for the saved-video store root: this channel's persisted source videos live under `<savedVideosDir>/<slug>/<videoId>/`. Trimmed; blank = the global store.", dataDir: - "Where this channel's media ACTUALLY lives when relocated to another drive: the absolute path `channels/<slug>/data` is a symlink to. Absent = in place. Written ONLY by the relocate / re-point jobs on success — a record of what is on disk, never free text, because a value that disagrees with the link is an \"inconsistent\" channel every guard refuses.", + "RETIRED (release 17). The whole-directory layout's record: the absolute path `channels/<slug>/data` was a symlink to. Still parsed for one release so a write never erases it: a channel that carries it — or whose `data/` is a link — is `legacy`, and every media job, lane and build holds it until `archilyzer storage migrate-tier <slug>` moves its text back and its media into `mediaDir`. Never written by anything but that migration, which removes it.", + mediaDir: + "Where this channel's big files live when relocated: `channels/<slug>/media` is a symlink to it, `<root>/<slug>/media`. Absent = in place. Written only by relocate / re-point / the tier migration — a record of what is on disk, never free text, because a value that disagrees with the link is an \"inconsistent\" channel every media guard refuses. The text (`data/`) never moves.", ytdlpExtraArgs: "Extra yt-dlp arguments, appended verbatim. Must be an array of strings or it is dropped.", subLangs: "yt-dlp `--sub-langs` value for caption downloads.", lastSyncedAt: @@ -368,6 +371,7 @@ export const CHANNEL_CONFIG_COERCIONS: { extractionMode: (v) => (v === "ytdlp" || v === "app" ? v : undefined), savedVideosDir: trimmedNonBlank, dataDir: trimmedNonBlank, + mediaDir: trimmedNonBlank, ytdlpExtraArgs: (v) => Array.isArray(v) && v.every((x) => typeof x === "string") ? (v as string[]) diff --git a/common/lib/channelConfigSchema.test.ts b/common/lib/channelConfigSchema.test.ts @@ -39,7 +39,7 @@ test("one key list: docs = coercions = schema shape, sync-state keys inside it", assert.deepEqual(Object.keys(CHANNEL_CONFIG_COERCIONS), [...CHANNEL_CONFIG_KEYS]); assert.deepEqual(Object.keys(channelConfigObjectSchema.shape), [...CHANNEL_CONFIG_KEYS]); assert.deepEqual(Object.keys(CHANNEL_CONFIG_FIELD_DOCS), [...CHANNEL_CONFIG_KEYS]); - assert.equal(CHANNEL_CONFIG_KEYS.length, 29); + assert.equal(CHANNEL_CONFIG_KEYS.length, 30); assert.equal(sameKeys, true); assert.equal(fits, true); for (const k of CHANNEL_SYNC_STATE_KEYS) assert.ok(CHANNEL_CONFIG_KEYS.includes(k), k); @@ -108,12 +108,14 @@ test("trims and normalises", () => { socialHandle: " @me ", savedVideosDir: " /x ", dataDir: " /d ", + mediaDir: " /m/chan/media ", cookiesFromBrowser: " firefox ", downloadFilter: { include: " a ", exclude: "", includeLivestreams: true, rejectedLivestreams: "skip" }, })!; assert.equal(cfg.socialHandle, "me"); assert.equal(cfg.savedVideosDir, "/x"); assert.equal(cfg.dataDir, "/d"); + assert.equal(cfg.mediaDir, "/m/chan/media"); assert.equal(cfg.cookiesFromBrowser, "firefox"); assert.deepEqual(cfg.downloadFilter, { include: "a", includeLivestreams: true }); }); diff --git a/common/lib/channelConfigSchema.ts b/common/lib/channelConfigSchema.ts @@ -57,6 +57,7 @@ export const channelConfigObjectSchema = z.object({ extractionMode: field("extractionMode"), savedVideosDir: field("savedVideosDir"), dataDir: field("dataDir"), + mediaDir: field("mediaDir"), ytdlpExtraArgs: field("ytdlpExtraArgs"), subLangs: field("subLangs"), lastSyncedAt: field("lastSyncedAt"), diff --git a/common/lib/channelMedia.test.ts b/common/lib/channelMedia.test.ts @@ -5,7 +5,11 @@ import { tmpdir } from "node:os"; import path from "node:path"; import { ChannelMediaUnreachableError, + ChannelTextUnreadableError, assertChannelMediaReachable, + assertChannelTextReadable, + channelMediaStall, + channelTextStall, clearRelocationMarker, inspectChannelMedia, readRelocationMarker, @@ -13,6 +17,9 @@ import { RELOCATION_MARKER_FILENAME, type ChannelMediaPaths, } from "./channelMedia"; +import { relocatedMediaDir } from "./mediaTier-server"; +import { recordLocationHealth, resetStorageHealth } from "./storageHealth"; +import { isMediaHeld, isTextHeld, HELD_REASON } from "./channelMediaHold"; // Run with: // pnpm --filter yt-dlp-transcript-common exec tsx --test lib/channelMedia.test.ts @@ -47,12 +54,12 @@ async function seedChannel( return channelDir; } -test("relocatedDataDir fixes the <root>/<slug>/data suffix", () => { +test("relocatedDataDir (retired) and relocatedMediaDir fix their suffixes", () => { assert.equal(relocatedDataDir("/mnt/p", "alpha"), "/mnt/p/alpha/data"); - assert.equal(relocatedDataDir(" /mnt/p ", "alpha"), "/mnt/p/alpha/data"); + assert.equal(relocatedMediaDir(" /mnt/p ", "alpha"), "/mnt/p/alpha/media"); }); -test("a real data dir with no config.dataDir is in-place", async () => { +test("a real data dir and no media link, no mediaDir: in-place (classic)", async () => { await withTmp(async (paths) => { const channelDir = await seedChannel(paths, "alpha"); await mkdir(path.join(channelDir, "data", "v1"), { recursive: true }); @@ -61,7 +68,21 @@ test("a real data dir with no config.dataDir is in-place", async () => { assert.equal(loc.relocated, false); assert.equal(loc.target, undefined); assert.equal(loc.dataDir, path.join(channelDir, "data")); + assert.equal(loc.mediaLink, path.join(channelDir, "media")); + assert.deepEqual(loc.text, { dir: path.join(channelDir, "data"), readable: true }); await assertChannelMediaReachable(paths, "alpha"); + await assertChannelTextReadable(paths, "alpha"); + }); +}); + +test("a real media/ directory on the corpus disk is in-place (tiered in place)", async () => { + await withTmp(async (paths) => { + const channelDir = await seedChannel(paths, "alpha"); + await mkdir(path.join(channelDir, "data", "v1"), { recursive: true }); + await mkdir(path.join(channelDir, "media", "v1"), { recursive: true }); + const loc = await inspectChannelMedia(paths, "alpha"); + assert.equal(loc.status, "in-place"); + assert.equal(loc.relocated, false); }); }); @@ -70,33 +91,42 @@ test("a channel that has downloaded nothing is in-place, not an error", async () await seedChannel(paths, "alpha"); const loc = await inspectChannelMedia(paths, "alpha"); assert.equal(loc.status, "in-place"); + assert.equal(loc.text.readable, true); await assertChannelMediaReachable(paths, "alpha"); + await assertChannelTextReadable(paths, "alpha"); }); }); -test("a link agreeing with config and pointing at a live dir is ok", async () => { +async function seedRelocated(paths: ChannelMediaPaths, root: string, opts: { mount?: boolean } = {}) { + const target = relocatedMediaDir(root, "alpha"); + if (opts.mount !== false) await mkdir(path.join(target, "v1"), { recursive: true }); + const channelDir = await seedChannel(paths, "alpha", { mediaDir: target }); + await mkdir(path.join(channelDir, "data", "v1"), { recursive: true }); + await symlink(target, path.join(channelDir, "media")); + return { target, channelDir }; +} + +test("a media link agreeing with mediaDir and pointing at a live dir is ok", async () => { await withTmp(async (paths, root) => { - const target = relocatedDataDir(root, "alpha"); - await mkdir(path.join(target, "v1"), { recursive: true }); - const channelDir = await seedChannel(paths, "alpha", { dataDir: target }); - await symlink(target, path.join(channelDir, "data")); + const { target } = await seedRelocated(paths, root); const loc = await inspectChannelMedia(paths, "alpha"); assert.equal(loc.status, "ok"); assert.equal(loc.relocated, true); assert.equal(loc.target, target); + assert.equal(loc.text.readable, true); await assertChannelMediaReachable(paths, "alpha"); + await assertChannelTextReadable(paths, "alpha"); }); }); -test("a dangling link (drive not mounted) is unreachable and throws", async () => { +test("a dangling media link (drive not mounted): media unreachable, text readable", async () => { await withTmp(async (paths, root) => { - const target = relocatedDataDir(root, "alpha"); - const channelDir = await seedChannel(paths, "alpha", { dataDir: target }); // The link is made WITHOUT creating the target: exactly an unmounted drive. - await symlink(target, path.join(channelDir, "data")); + await seedRelocated(paths, root, { mount: false }); const loc = await inspectChannelMedia(paths, "alpha"); assert.equal(loc.status, "unreachable"); assert.equal(loc.relocated, true); + assert.equal(loc.text.readable, true); assert.match(loc.detail ?? "", /does not exist/); await assert.rejects( () => assertChannelMediaReachable(paths, "alpha"), @@ -108,24 +138,44 @@ test("a dangling link (drive not mounted) is unreachable and throws", async () = return true; }, ); + // THE NEW RULE: the text is on the corpus disk, so it is not held. + await assertChannelTextReadable(paths, "alpha"); + assert.equal(isMediaHeld(loc.status), true); + assert.equal(isTextHeld(loc.status), false); }); }); test("an EMPTY mountpoint is still unreachable — the link points deep", async () => { await withTmp(async (paths, root) => { - const target = relocatedDataDir(root, "alpha"); - // The root exists (mountpoint present) but holds nothing. await mkdir(root, { recursive: true }); - const channelDir = await seedChannel(paths, "alpha", { dataDir: target }); - await symlink(target, path.join(channelDir, "data")); + await seedRelocated(paths, root, { mount: false }); const loc = await inspectChannelMedia(paths, "alpha"); assert.equal(loc.status, "unreachable"); }); }); -test("a marker makes the channel in-transition whatever the disk says", async () => { +test("a stalled media location: answered from memory, text still readable", async () => { await withTmp(async (paths, root) => { - const target = relocatedDataDir(root, "alpha"); + resetStorageHealth(); + try { + const { target } = await seedRelocated(paths, root); + recordLocationHealth({ id: "platter", label: "Platter", root }, "stalled"); + const loc = await inspectChannelMedia(paths, "alpha", undefined, { fresh: true }); + assert.equal(loc.status, "stalled"); + assert.equal(loc.text.readable, true); + assert.ok(channelMediaStall({ mediaDir: target })); + assert.equal(channelTextStall({ mediaDir: target }), null); + await assert.rejects(() => assertChannelMediaReachable(paths, "alpha")); + await assertChannelTextReadable(paths, "alpha"); + } finally { + resetStorageHealth(); + } + }); +}); + +test("a marker makes the channel in-transition; a media move leaves the text readable", async () => { + await withTmp(async (paths, root) => { + const target = relocatedMediaDir(root, "alpha"); const channelDir = await seedChannel(paths, "alpha"); await mkdir(path.join(channelDir, "data"), { recursive: true }); await writeFile( @@ -135,81 +185,186 @@ test("a marker makes the channel in-transition whatever the disk says", async () direction: "out", startedAt: new Date().toISOString(), phase: "copy", + scope: "media", }), ); const loc = await inspectChannelMedia(paths, "alpha"); assert.equal(loc.status, "in-transition"); assert.equal(loc.marker?.phase, "copy"); assert.equal(loc.marker?.direction, "out"); + assert.equal(loc.marker?.scope, "media"); assert.equal(loc.target, target); + assert.equal(loc.text.readable, true); await assert.rejects( () => assertChannelMediaReachable(paths, "alpha"), ChannelMediaUnreachableError, ); + await assertChannelTextReadable(paths, "alpha"); const marker = await readRelocationMarker(paths, "alpha"); assert.equal(marker?.target, target); }); }); -test("a link that disagrees with config is inconsistent, never guessed past", async () => { +test("a tier-migration marker holds the text too", async () => { + await withTmp(async (paths, root) => { + const channelDir = await seedChannel(paths, "alpha"); + await mkdir(path.join(channelDir, "data"), { recursive: true }); + await writeFile( + path.join(channelDir, RELOCATION_MARKER_FILENAME), + JSON.stringify({ + target: relocatedMediaDir(root, "alpha"), + direction: "out", + phase: "copy", + scope: "tier-migration", + }), + ); + const loc = await inspectChannelMedia(paths, "alpha"); + assert.equal(loc.status, "in-transition"); + assert.equal(loc.text.readable, false); + await assert.rejects( + () => assertChannelTextReadable(paths, "alpha"), + (err: unknown) => { + assert.ok(err instanceof ChannelTextUnreadableError); + assert.match(err.message, /being migrated/); + return true; + }, + ); + }); +}); + +test("LEGACY: a data link with dataDir — held by both guards, no call to the drive", async () => { await withTmp(async (paths, root) => { - const real = relocatedDataDir(root, "alpha"); - const recorded = relocatedDataDir(path.join(root, "other"), "alpha"); + const target = relocatedDataDir(root, "alpha"); + await mkdir(path.join(target, "v1"), { recursive: true }); + const channelDir = await seedChannel(paths, "alpha", { dataDir: target }); + await symlink(target, path.join(channelDir, "data")); + const loc = await inspectChannelMedia(paths, "alpha"); + assert.equal(loc.status, "legacy"); + assert.equal(loc.relocated, true); + assert.equal(loc.target, target); + assert.equal(loc.text.readable, false); + assert.match(loc.detail ?? "", /archilyzer storage migrate-tier alpha/); + assert.equal(isMediaHeld(loc.status), true); + assert.equal(isTextHeld(loc.status), true); + assert.match(HELD_REASON.legacy, /migrate-tier/); + await assert.rejects( + () => assertChannelMediaReachable(paths, "alpha"), + (err: unknown) => { + assert.ok(err instanceof ChannelMediaUnreachableError); + assert.match(err.message, /migrate-tier/); + return true; + }, + ); + await assert.rejects( + () => assertChannelTextReadable(paths, "alpha"), + (err: unknown) => { + assert.ok(err instanceof ChannelTextUnreadableError); + assert.equal(err.status, "legacy"); + assert.match(err.message, /migrate-tier alpha/); + return true; + }, + ); + }); +}); + +test("LEGACY: a dangling data link (unmounted) is legacy too, never 'no videos'", async () => { + await withTmp(async (paths, root) => { + const target = relocatedDataDir(root, "alpha"); + const channelDir = await seedChannel(paths, "alpha", { dataDir: target }); + await symlink(target, path.join(channelDir, "data")); + assert.equal((await inspectChannelMedia(paths, "alpha")).status, "legacy"); + }); +}); + +test("LEGACY: a data link with no dataDir, and dataDir with a real data dir", async () => { + await withTmp(async (paths, root) => { + const target = relocatedDataDir(root, "alpha"); + await mkdir(target, { recursive: true }); + const a = await seedChannel(paths, "alpha"); + await symlink(target, path.join(a, "data")); + assert.equal((await inspectChannelMedia(paths, "alpha")).status, "legacy"); + + const b = await seedChannel(paths, "beta", { dataDir: relocatedDataDir(root, "beta") }); + await mkdir(path.join(b, "data"), { recursive: true }); + const loc = await inspectChannelMedia(paths, "beta"); + assert.equal(loc.status, "legacy"); + assert.equal(loc.text.readable, false); + assert.ok(channelTextStall({ dataDir: relocatedDataDir(root, "beta") }) === null); + }); +}); + +test("a media link that disagrees with mediaDir is inconsistent, never guessed past", async () => { + await withTmp(async (paths, root) => { + const real = relocatedMediaDir(root, "alpha"); + const recorded = relocatedMediaDir(path.join(root, "other"), "alpha"); await mkdir(real, { recursive: true }); - const channelDir = await seedChannel(paths, "alpha", { dataDir: recorded }); - await symlink(real, path.join(channelDir, "data")); + const channelDir = await seedChannel(paths, "alpha", { mediaDir: recorded }); + await symlink(real, path.join(channelDir, "media")); const loc = await inspectChannelMedia(paths, "alpha"); assert.equal(loc.status, "inconsistent"); assert.match(loc.detail ?? "", /points at/); + assert.equal(loc.text.readable, true); await assert.rejects( () => assertChannelMediaReachable(paths, "alpha"), ChannelMediaUnreachableError, ); + await assertChannelTextReadable(paths, "alpha"); }); }); -test("a link with no config.dataDir is inconsistent", async () => { +test("a media link with no mediaDir is inconsistent", async () => { await withTmp(async (paths, root) => { - const target = relocatedDataDir(root, "alpha"); + const target = relocatedMediaDir(root, "alpha"); await mkdir(target, { recursive: true }); const channelDir = await seedChannel(paths, "alpha"); - await symlink(target, path.join(channelDir, "data")); + await symlink(target, path.join(channelDir, "media")); const loc = await inspectChannelMedia(paths, "alpha"); assert.equal(loc.status, "inconsistent"); assert.equal(loc.relocated, false); - assert.match(loc.detail ?? "", /records no dataDir/); + assert.match(loc.detail ?? "", /records no mediaDir/); }); }); -test("config.dataDir with a real directory on disk is inconsistent", async () => { +test("mediaDir with a real media/ directory on disk is inconsistent", async () => { await withTmp(async (paths, root) => { - const target = relocatedDataDir(root, "alpha"); - const channelDir = await seedChannel(paths, "alpha", { dataDir: target }); - await mkdir(path.join(channelDir, "data"), { recursive: true }); + const target = relocatedMediaDir(root, "alpha"); + const channelDir = await seedChannel(paths, "alpha", { mediaDir: target }); + await mkdir(path.join(channelDir, "media"), { recursive: true }); const loc = await inspectChannelMedia(paths, "alpha"); assert.equal(loc.status, "inconsistent"); assert.match(loc.detail ?? "", /never moved/); }); }); -test("config.dataDir with no data/ at all is inconsistent (link gone)", async () => { +test("mediaDir with no media link at all is inconsistent (link gone)", async () => { await withTmp(async (paths, root) => { - const target = relocatedDataDir(root, "alpha"); - await seedChannel(paths, "alpha", { dataDir: target }); + const target = relocatedMediaDir(root, "alpha"); + await seedChannel(paths, "alpha", { mediaDir: target }); const loc = await inspectChannelMedia(paths, "alpha"); assert.equal(loc.status, "inconsistent"); assert.match(loc.detail ?? "", /symlink is missing/); }); }); +test("a data/ that is a file is inconsistent and its text unreadable", async () => { + await withTmp(async (paths) => { + const channelDir = await seedChannel(paths, "alpha"); + await writeFile(path.join(channelDir, "data"), "not a dir"); + const loc = await inspectChannelMedia(paths, "alpha"); + assert.equal(loc.status, "inconsistent"); + assert.equal(loc.text.readable, false); + await assert.rejects(() => assertChannelTextReadable(paths, "alpha"), ChannelTextUnreadableError); + }); +}); + test("a passed config is used verbatim; config.json is only read when it is absent", async () => { await withTmp(async (paths, root) => { - const target = relocatedDataDir(root, "alpha"); + const target = relocatedMediaDir(root, "alpha"); await mkdir(target, { recursive: true }); // config.json on disk says NOTHING about a relocation... const channelDir = await seedChannel(paths, "alpha"); - await symlink(target, path.join(channelDir, "data")); + await symlink(target, path.join(channelDir, "media")); // ...so reading it itself gives "inconsistent"... assert.equal( @@ -219,7 +374,7 @@ test("a passed config is used verbatim; config.json is only read when it is abse // ...while a caller that hands over the config it already holds gets the // answer for THAT config, with no second read. const passed = await inspectChannelMedia(paths, "alpha", { - dataDir: target, + mediaDir: target, }); assert.equal(passed.status, "ok"); assert.equal(passed.target, target); @@ -231,13 +386,13 @@ test("a passed config is used verbatim; config.json is only read when it is abse }); }); -test("a blank config.dataDir means in place", async () => { +test("a blank mediaDir or dataDir means in place", async () => { await withTmp(async (paths) => { - const channelDir = await seedChannel(paths, "alpha", { dataDir: " " }); + const channelDir = await seedChannel(paths, "alpha", { dataDir: " ", mediaDir: " " }); await mkdir(path.join(channelDir, "data"), { recursive: true }); assert.equal((await inspectChannelMedia(paths, "alpha")).status, "in-place"); assert.equal( - (await inspectChannelMedia(paths, "alpha", { dataDir: " " })).status, + (await inspectChannelMedia(paths, "alpha", { dataDir: " ", mediaDir: "" })).status, "in-place", ); }); @@ -253,11 +408,8 @@ test("an unreadable or missing config.json is not a relocation", async () => { }); test("clearRelocationMarker removes the marker and touches nothing else", async () => { - await withTmp(async (paths) => { - const target = path.join(paths.channelsDir, "..", "platter", "alpha", "data"); - await mkdir(target, { recursive: true }); - const channelDir = await seedChannel(paths, "alpha", { dataDir: target }); - await symlink(target, path.join(channelDir, "data")); + await withTmp(async (paths, root) => { + const { target, channelDir } = await seedRelocated(paths, root); await writeFile( path.join(channelDir, RELOCATION_MARKER_FILENAME), JSON.stringify({ target, direction: "out", phase: "swap", startedAt: "" }), @@ -279,3 +431,20 @@ test("clearRelocationMarker removes the marker and touches nothing else", async await clearRelocationMarker(paths, "alpha"); }); }); + +test("a scope-less marker aimed at the retired <root>/<slug>/data holds the text (the old mover)", async () => { + await withTmp(async (paths, root) => { + const channelDir = await seedChannel(paths, "alpha"); + await mkdir(path.join(channelDir, "data"), { recursive: true }); + await writeFile( + path.join(channelDir, RELOCATION_MARKER_FILENAME), + JSON.stringify({ target: relocatedDataDir(root, "alpha"), direction: "out", phase: "copy" }), + ); + assert.equal((await inspectChannelMedia(paths, "alpha", undefined, { fresh: true })).text.readable, false); + await writeFile( + path.join(channelDir, RELOCATION_MARKER_FILENAME), + JSON.stringify({ target: relocatedMediaDir(root, "alpha"), direction: "out", phase: "copy" }), + ); + assert.equal((await inspectChannelMedia(paths, "alpha", undefined, { fresh: true })).text.readable, true); + }); +}); diff --git a/common/lib/channelMedia.ts b/common/lib/channelMedia.ts @@ -12,25 +12,33 @@ import { type LocationHealth, } from "./storageHealth"; import { secondsText } from "./storageHealthTimings"; +import { MEDIA_LINK_NAME } from "./mediaTier-server"; // WHERE A CHANNEL'S MEDIA ACTUALLY IS, and whether it can be reached. // -// A channel's downloaded media lives at `channels/<slug>/data/`. That path is -// joined inline at ~74 call sites and is the on-disk contract every reader, -// yt-dlp's cwd-relative output template and the LMDB index depend on, so -// relocating a channel to another drive does NOT change it: `data/` becomes an -// absolute SYMLINK to `<root>/<slug>/data` and `config.json` records the target -// in `dataDir`. Every existing reader follows the link transparently — there is -// no symlink-aware code anywhere in common/, editor/ or export/, and there does -// not need to be. +// A channel's files live at `channels/<slug>/data/<id>/`. That path is joined +// inline at ~74 call sites and is the on-disk contract every reader, yt-dlp's +// cwd-relative output template and the LMDB index depend on, so it never moves. +// Since release 17 only the BIG files leave it (lib/mediaTier.ts says which): +// each becomes a RELATIVE link `data/<id>/<name> -> ../../media/<id>/<name>`, +// and `channels/<slug>/media` is either a real directory on the corpus disk +// (tiered in place) or ONE absolute symlink to `<root>/<slug>/media` on another +// drive, recorded in `config.json` as `mediaDir` (relocated). The text — +// transcripts, cues, metadata, every sidecar — stays on the corpus disk. // -// What that buys in call-site churn it owes in one new failure mode: an -// unmounted drive. A dangling link reads as ENOENT, and the three places that -// enumerate `data/` swallow ENOENT as "this channel has no videos" — which to an -// unattended runner means *everything is undownloaded* and is an instruction to -// re-download hundreds of gigabytes onto the volume that was too full to hold -// them. This module is the one place that can tell those two apart, and the -// guards that call it are what make the symlink safe. +// What the link buys in call-site churn it owes in one failure mode: an +// unmounted drive. A dangling link reads as ENOENT. For the media alone that is +// survivable — a reader of the text never touches it — and this module is the +// one place that can tell the media's states apart, so the guards that call it +// are what make the link safe: `assertChannelMediaReachable` for a job that +// opens a big file, `assertChannelTextReadable` for one that reads only text. +// +// THE RETIRED LAYOUT. Before release 17 a relocation moved the whole `data/` +// (an absolute symlink `data -> <root>/<slug>/data`, `config.dataDir`). Such a +// channel is `legacy`: its text is on the far drive too, so it is held by BOTH +// guards — every lane, every media job, the index and stats builds — until +// `archilyzer storage migrate-tier <slug>` brings its text home. `dataDir` is +// still parsed one release for exactly that (channelConfig.ts). // // IT LIVES IN lib/ AND MAY NOT IMPORT controller/ (architecture.test.ts), which // is where readChannelConfig is. Hence the optional `config` argument: a caller @@ -52,18 +60,28 @@ export const RELOCATION_MARKER_FILENAME = ".relocating.json"; export type RelocationDirection = "out" | "back"; export type RelocationPhase = "copy" | "swap" | "reclaim"; +// What a marker is moving. `media` (or absent): the channel's media tier — its +// text stays readable, so only its media writers are held. `tier-migration`: +// the one-off migration off the retired layout (common/bin/migrate-media-tier.ts), +// which rebuilds `data/` itself, so the text is held too. +export type RelocationScope = "media" | "tier-migration"; + export type RelocationMarker = { - // Absolute path of the relocated data dir: <root>/<slug>/data. + // Absolute path of the relocated media dir: <root>/<slug>/media (the retired + // mover wrote <root>/<slug>/data). target: string; direction: RelocationDirection; startedAt: string; phase: RelocationPhase; + scope?: RelocationScope; }; export type ChannelMediaStatus = - // No relocation: `data/` is a real directory (or does not exist yet). + // No relocation: no `media` link (a classic channel, or one tiered into a + // real `media/` directory on the corpus disk) and no `mediaDir`. | "in-place" - // Relocated, link and config agree, and the target is a reachable directory. + // Relocated: the `media` link and `mediaDir` agree, and the target is a + // reachable directory. | "ok" // Relocated, but the target is not there — almost always an unmounted drive. | "unreachable" @@ -75,15 +93,30 @@ export type ChannelMediaStatus = // (`lib/storageHealth.ts`). Answered from memory, WITHOUT a filesystem call: // a call there would block one of the process's few I/O threads for as long // as the drive takes to come back. Held and refused like `unreachable`. - | "stalled"; + | "stalled" + // The RETIRED whole-directory layout: `data/` is a link, or config.json still + // records `dataDir`. Its text is not on the corpus disk, so it is held by the + // text guard as well as the media one until `archilyzer storage migrate-tier` + // runs. Answered without a call to the far drive. + | "legacy"; export type ChannelMediaLocation = { - // Always channelDir/data — the path every reader uses, relocated or not. + // Always channelDir/data — the real text dir every reader uses (a link only + // on a `legacy` channel). dataDir: string; - // Whether config.json records a relocation target. + // Always channelDir/media — the media tier's one name. + mediaLink: string; + // Whether config.json records a relocation target (`mediaDir`, or the + // retired `dataDir` on a legacy channel). relocated: boolean; - // config.dataDir (or, mid-transition with no config yet, the marker's target). + // config.mediaDir (or, mid-transition with no config yet, the marker's + // target; on a legacy channel the retired `dataDir`). target?: string; + // The text tier: `data/`, and whether a reader may walk it. Not readable on a + // `legacy` channel (its text is on the far drive), when `data/` is something + // other than a directory, or mid tier-migration. A channel with no `data/` + // yet is readable (it has downloaded nothing). + text: { dir: string; readable: boolean }; status: ChannelMediaStatus; // Operator-readable reason, set for every status except "in-place" and "ok". detail?: string; @@ -111,13 +144,50 @@ export class ChannelMediaUnreachableError extends Error { } } -// The relocated layout, fixed so one root can hold many channels and the shape -// mirrors the saved-video store (<root>/<slug>/<...>). The `<slug>/data` suffix -// is not configurable: deleteChannel and renameChannel recognise a target by it. +// Thrown by assertChannelTextReadable: the channel's TEXT cannot be walked — a +// `legacy` channel, a `data/` that is not a directory, a tier migration in +// flight. Same shape as the media error so a lane can skip on either. +export class ChannelTextUnreadableError extends Error { + readonly slug: string; + readonly status: ChannelMediaStatus; + readonly location: ChannelMediaLocation; + constructor(slug: string, location: ChannelMediaLocation, detail: string) { + super(`Channel "${slug}": its text is not readable — ${detail}`); + this.name = "ChannelTextUnreadableError"; + this.slug = slug; + this.status = location.status; + this.location = location; + } +} + +// THE RETIRED layout's target shape, `<root>/<slug>/data`. Still exported for +// the mover, the re-point and the rename until release 17 slice T2 rebases them +// on `relocatedMediaDir` (lib/mediaTier-server.ts). export function relocatedDataDir(root: string, slug: string): string { return path.join(root.trim(), slug, "data"); } +// WHETHER A MARKER HOLDS THE TEXT TOO. A media move (`scope` "media", or none +// with a `<root>/<slug>/media` target) carries `media/` only, so its channel's +// text stays readable. A tier migration rebuilds `data/` itself; and a marker +// with no scope whose target is the retired `<root>/<slug>/data` shape is the +// old whole-directory mover's, copying `data/` — both hold the text. +export function markerHoldsText( + marker: Pick<RelocationMarker, "scope" | "target">, +): boolean { + if (marker.scope === "tier-migration") return true; + if (marker.scope === "media") return false; + return path.basename(marker.target.replace(/\/+$/, "")) === "data"; +} + +// The sentence a legacy channel is refused with, naming the way out. +export function legacyDetail(slug: string): string { + return ( + `its media layout is the retired whole-directory one — run ` + + `archilyzer storage migrate-tier ${slug}` + ); +} + export function channelMediaDir( paths: ChannelMediaPaths, slug: string, @@ -139,12 +209,14 @@ function parseMarker(raw: unknown): RelocationMarker | null { const direction = r.direction === "back" ? "back" : "out"; const phase = r.phase === "swap" || r.phase === "reclaim" ? r.phase : "copy"; - return { + const marker: RelocationMarker = { target: r.target, direction, startedAt: typeof r.startedAt === "string" ? r.startedAt : "", phase, }; + if (r.scope === "media" || r.scope === "tier-migration") marker.scope = r.scope; + return marker; } export async function readRelocationMarker( @@ -180,24 +252,37 @@ export async function clearRelocationMarker( forgetChannelMedia(slug); } -// The `dataDir` field alone, read straight off config.json. Deliberately NOT -// parseChannelConfig: this runs in guards on hot paths and must not depend on -// the controller that owns the rest of the schema. -async function readConfiguredDataDir( +// The two fields this module needs, the way a ChannelConfig carries them. +export type ChannelMediaConfig = Pick<ChannelConfig, "mediaDir" | "dataDir">; + +type Configured = { mediaDir?: string; dataDir?: string }; + +function trimmed(v: unknown): string | undefined { + if (typeof v !== "string") return undefined; + const t = v.trim(); + return t === "" ? undefined : t; +} + +function configuredOf(config: ChannelMediaConfig | null | undefined): Configured { + return { mediaDir: trimmed(config?.mediaDir), dataDir: trimmed(config?.dataDir) }; +} + +// `mediaDir` and the retired `dataDir`, read straight off config.json. +// Deliberately NOT parseChannelConfig: this runs in guards on hot paths and +// must not depend on the controller that owns the rest of the schema. +async function readConfigured( paths: ChannelMediaPaths, slug: string, -): Promise<string | undefined> { +): Promise<Configured> { try { const raw = await readFile( path.join(paths.channelsDir, slug, "config.json"), "utf8", ); - const parsed = JSON.parse(raw) as { dataDir?: unknown }; - if (typeof parsed.dataDir !== "string") return undefined; - const trimmed = parsed.dataDir.trim(); - return trimmed === "" ? undefined : trimmed; + const parsed = JSON.parse(raw) as { mediaDir?: unknown; dataDir?: unknown }; + return { mediaDir: trimmed(parsed.mediaDir), dataDir: trimmed(parsed.dataDir) }; } catch { - return undefined; + return {}; } } @@ -216,8 +301,11 @@ export function stalledMediaLocation( ): ChannelMediaLocation { return { dataDir, + mediaLink: path.join(path.dirname(dataDir), MEDIA_LINK_NAME), relocated: true, target: configured, + // The text is on the corpus disk: a stalled MEDIA drive does not hold it. + text: { dir: dataDir, readable: true }, status: "stalled", detail: health ? `${NOT_ANSWERING} (location "${health.label}", ${sinceText(health.since)})` @@ -227,12 +315,25 @@ export function stalledMediaLocation( } // THE STALL, for a caller holding a parsed config: the stalled location the -// channel's media is on, or null. No I/O — the question every page and poll -// that reads a channel's `data/` asks before it does. +// channel's MEDIA is on, or null. No I/O — the question a page or poll that +// opens a big file asks before it does. Keys on `mediaDir`; on a legacy +// channel (no `mediaDir` yet) on its retired `dataDir`, where everything is. export function channelMediaStall( - config: Pick<ChannelConfig, "dataDir"> | null | undefined, + config: ChannelMediaConfig | null | undefined, +): LocationHealth | null { + const c = configuredOf(config); + const dir = c.mediaDir ?? c.dataDir; + return dir ? stalledLocationForPath(dir) : null; +} + +// THE STALL OF THE TEXT: non-null only for a legacy channel whose retired +// `dataDir` is on a stalled location — the one layout whose text is on another +// drive. A page that reads only text (a listing, a transcript) asks this, not +// `channelMediaStall`: a stalled media drive never holds the text. +export function channelTextStall( + config: ChannelMediaConfig | null | undefined, ): LocationHealth | null { - const dir = config?.dataDir?.trim(); + const dir = configuredOf(config).dataDir; return dir ? stalledLocationForPath(dir) : null; } @@ -240,12 +341,12 @@ export function channelMediaStall( // The memo // --------------------------------------------------------------------------- // -// FIVE SECONDS, PER CHANNEL, KEYED BY SLUG AND THE CONFIGURED TARGET. The home +// FIVE SECONDS, PER CHANNEL, KEYED BY SLUG AND THE CONFIGURED TARGETS. The home // page, /channels and the auto-queue status poll (every three seconds, four -// lanes) each inspect every channel; without this each of them costs three +// lanes) each inspect every channel; without this each of them costs a few // syscalls a relocated channel, every time, on a drive that may be the slow -// one. The key carries the configured `dataDir`, so a move that rewrites it is -// a new key at once, and the movers clear the memo outright +// one. The key carries the configured `mediaDir` and `dataDir`, so a move that +// rewrites either is a new key at once, and the movers clear the memo outright // (`forgetChannelMedia`) whenever a marker is written or removed. // // THE STALL GATE IS ASKED BEFORE A REMEMBERED ANSWER IS GIVEN, so a drive that @@ -300,33 +401,28 @@ export function forgetChannelMedia(slug?: string): void { } } -function memoKey(paths: ChannelMediaPaths, slug: string, configured?: string): string { - return `${paths.channelsDir}\u0000${slug}\u0000${configured ?? ""}`; +function memoKey(paths: ChannelMediaPaths, slug: string, c: Configured): string { + return `${paths.channelsDir}\u0000${slug}\u0000${c.mediaDir ?? ""}\u0000${c.dataDir ?? ""}`; } // Past this many entries, expired ones are swept on insert. A corpus has tens // of channels; this only matters to a process that inspects many corpora. const MEMO_SWEEP_AT = 512; -// Two stats and (at most) one small JSON read. Render-safe: nothing here walks a -// directory, so calling it per channel on a listing page costs three syscalls a -// row — or none, for five seconds after the last answer (see the memo above). -// On a stalled location the one call that reaches the drive (the target's -// `stat`) is not made; the marker, the link and config.json are on the corpus -// disk and are read as usual. +// A few lstats and (at most) one small JSON read. Render-safe: nothing here +// walks a directory, so calling it per channel on a listing page costs a few +// syscalls a row — or none, for five seconds after the last answer (see the +// memo above). On a stalled location the one call that reaches the drive (the +// media target's `stat`) is not made; the marker, the links and config.json are +// on the corpus disk and are read as usual. export async function inspectChannelMedia( paths: ChannelMediaPaths, slug: string, - config?: Pick<ChannelConfig, "dataDir"> | null, + config?: ChannelMediaConfig | null, opts: InspectOptions = {}, ): Promise<ChannelMediaLocation> { - const dataDir = channelMediaDir(paths, slug); const configured = - config === undefined - ? await readConfiguredDataDir(paths, slug) - : config?.dataDir && config.dataDir.trim() !== "" - ? config.dataDir.trim() - : undefined; + config === undefined ? await readConfigured(paths, slug) : configuredOf(config); const now = opts.now ?? Date.now(); const key = memoKey(paths, slug, configured); @@ -334,14 +430,17 @@ export async function inspectChannelMedia( if (!opts.fresh) { const hit = memo.get(key); if (hit && now - hit.at < CHANNEL_MEDIA_MEMO_MS) { - if (hit.location.status !== "in-transition" && configured) { - const stall = stalledLocationForPath(configured); - if (stall) return stalledMediaLocation(dataDir, configured, stall); + const st = hit.location.status; + if (st !== "in-transition" && st !== "legacy" && configured.mediaDir) { + const stall = stalledLocationForPath(configured.mediaDir); + if (stall) { + return stalledMediaLocation(hit.location.dataDir, configured.mediaDir, stall); + } } - return { ...hit.location }; + return cloneLocation(hit.location); } } - const location = await inspectOnDisk(paths, slug, dataDir, configured); + const location = await inspectOnDisk(paths, slug, configured); // A stall is not remembered: the health state is already its memory. if (!opts.fresh && location.status !== "stalled") { if (memo.size >= MEMO_SWEEP_AT) { @@ -351,25 +450,53 @@ export async function inspectChannelMedia( } memo.set(key, { at: now, location }); } - return { ...location }; + return cloneLocation(location); +} + +function cloneLocation(l: ChannelMediaLocation): ChannelMediaLocation { + return { ...l, text: { ...l.text } }; +} + +// `data/` on the corpus disk: absent (nothing downloaded yet — readable), a +// directory (readable), a link (the retired layout), or something else. +async function textState( + dataDir: string, +): Promise<"absent" | "dir" | "link" | "other"> { + try { + const l = await lstat(dataDir); + if (l.isSymbolicLink()) return "link"; + return l.isDirectory() ? "dir" : "other"; + } catch { + return "absent"; + } } async function inspectOnDisk( paths: ChannelMediaPaths, slug: string, - dataDir: string, - configured: string | undefined, + configured: Configured, ): Promise<ChannelMediaLocation> { + const dataDir = channelMediaDir(paths, slug); + const mediaLink = path.join(paths.channelsDir, slug, MEDIA_LINK_NAME); + const { mediaDir } = configured; + const text = await textState(dataDir); + const textReadable = text === "absent" || text === "dir"; + const base = { dataDir, mediaLink }; + // THE MARKER FIRST. It is in the channel dir, on the corpus disk, so reading // it costs the drive nothing — and a channel mid-move reads `in-transition` // whatever its drive is doing, which is what a resumed move and every guard - // key off. + // key off. A media move leaves the text readable; a tier migration does not. const marker = await readRelocationMarker(paths, slug); if (marker) { return { - dataDir, - relocated: Boolean(configured), - target: configured ?? marker.target, + ...base, + relocated: Boolean(mediaDir ?? configured.dataDir), + target: mediaDir ?? marker.target, + text: { + dir: dataDir, + readable: textReadable && !markerHoldsText(marker), + }, status: "in-transition", detail: `a media relocation (${marker.direction}) is in progress or was ` + @@ -378,30 +505,67 @@ async function inspectOnDisk( }; } - // THE GATE: a channel whose configured target is on a location whose drive - // is not answering is answered from memory, before the link is looked at and - // before the target's stat, which would hold an I/O thread for as long as the - // drive takes. - if (configured) { - const stall = stalledLocationForPath(configured); - if (stall) return stalledMediaLocation(dataDir, configured, stall); + // THE RETIRED LAYOUT, from the corpus disk alone: a `data` link or a + // recorded `dataDir`. Never a call to the far drive — the migration is what + // reads it, with the editor stopped. + if (text === "link" || configured.dataDir) { + let linkTarget: string | undefined; + if (text === "link") { + try { + linkTarget = await readlink(dataDir); + } catch { + /* unreadable: the config's value, if any, stands */ + } + } + return { + ...base, + relocated: true, + target: configured.dataDir ?? linkTarget, + text: { dir: dataDir, readable: false }, + status: "legacy", + detail: legacyDetail(slug), + }; + } + + const textRecord = { dir: dataDir, readable: textReadable }; + if (!textReadable) { + return { + ...base, + relocated: Boolean(mediaDir), + target: mediaDir, + text: textRecord, + status: "inconsistent", + detail: `${dataDir} is neither a directory nor a symlink`, + }; + } + + // THE GATE: a channel whose media is on a location whose drive is not + // answering is answered from memory, before the link is looked at and before + // the target's stat, which would hold an I/O thread for as long as the drive + // takes. Its text stays readable. + if (mediaDir) { + const stall = stalledLocationForPath(mediaDir); + if (stall) return stalledMediaLocation(dataDir, mediaDir, stall); } let link: Awaited<ReturnType<typeof lstat>> | null = null; try { - link = await lstat(dataDir); + link = await lstat(mediaLink); } catch { - // No data/ at all. With no configured target that is just a channel that - // has downloaded nothing yet — the overwhelmingly common case, and not an - // error. With one, the link this channel is supposed to have is gone. - if (!configured) return { dataDir, relocated: false, status: "in-place" }; + // No media link at all. With no configured target that is a classic + // channel — its media are real files in `data/<id>/` — the overwhelmingly + // common case, and not an error. With one, the link is gone. + if (!mediaDir) { + return { ...base, relocated: false, text: textRecord, status: "in-place" }; + } return { - dataDir, + ...base, relocated: true, - target: configured, + target: mediaDir, + text: textRecord, status: "inconsistent", detail: - `config.json records dataDir ${configured} but ${dataDir} does not ` + + `config.json records mediaDir ${mediaDir} but ${mediaLink} does not ` + `exist — the symlink is missing`, }; } @@ -409,34 +573,36 @@ async function inspectOnDisk( if (link.isSymbolicLink()) { let linkTarget = ""; try { - linkTarget = await readlink(dataDir); + linkTarget = await readlink(mediaLink); } catch { /* readlink of a link we just lstat'd: treat as unreadable below */ } - if (!configured) { + if (!mediaDir) { return { - dataDir, + ...base, relocated: false, target: linkTarget || undefined, + text: textRecord, status: "inconsistent", detail: - `${dataDir} is a symlink to ${linkTarget || "(unreadable)"} but ` + - `config.json records no dataDir`, + `${mediaLink} is a symlink to ${linkTarget || "(unreadable)"} but ` + + `config.json records no mediaDir`, }; } - if (path.resolve(linkTarget) !== path.resolve(configured)) { + if (path.resolve(linkTarget) !== path.resolve(mediaDir)) { return { - dataDir, + ...base, relocated: true, - target: configured, + target: mediaDir, + text: textRecord, status: "inconsistent", detail: - `${dataDir} points at ${linkTarget || "(unreadable)"} but ` + - `config.json records ${configured}`, + `${mediaLink} points at ${linkTarget || "(unreadable)"} but ` + + `config.json records ${mediaDir}`, }; } - // The link points at a DEEP path (<root>/<slug>/data), so an unmounted root - // gives ENOENT here. An empty mountpoint can never be mistaken for the + // The link points at a DEEP path (<root>/<slug>/media), so an unmounted + // root gives ENOENT here. An empty mountpoint can never be mistaken for the // media, which is the whole reason the suffix is fixed. // // THE ONE CALL HERE THAT REACHES THE DRIVE, so it goes through the @@ -444,55 +610,61 @@ async function inspectOnDisk( // not answered within the budget (`storage.health.budgetMs`, 3 s by // default) marks it stalled and answers `stalled` now. try { - const st = await onDrive(configured, () => stat(configured)); + const st = await onDrive(mediaDir, () => stat(mediaDir)); if (!st.isDirectory()) { return { - dataDir, + ...base, relocated: true, - target: configured, + target: mediaDir, + text: textRecord, status: "unreachable", - detail: `${configured} exists but is not a directory`, + detail: `${mediaDir} exists but is not a directory`, }; } } catch (err) { if (isDriveNotAnswering(err)) { - return stalledMediaLocation(dataDir, configured, err.health, err.message); + return stalledMediaLocation(dataDir, mediaDir, err.health, err.message); } return { - dataDir, + ...base, relocated: true, - target: configured, + target: mediaDir, + text: textRecord, status: "unreachable", - detail: `${configured} does not exist (drive not mounted?)`, + detail: `${mediaDir} does not exist (drive not mounted?)`, }; } - return { dataDir, relocated: true, target: configured, status: "ok" }; + return { ...base, relocated: true, target: mediaDir, text: textRecord, status: "ok" }; } if (!link.isDirectory()) { return { - dataDir, - relocated: Boolean(configured), - target: configured, + ...base, + relocated: Boolean(mediaDir), + target: mediaDir, + text: textRecord, status: "inconsistent", - detail: `${dataDir} is neither a directory nor a symlink`, + detail: `${mediaLink} is neither a directory nor a symlink`, }; } - if (configured) { + if (mediaDir) { return { - dataDir, + ...base, relocated: true, - target: configured, + target: mediaDir, + text: textRecord, status: "inconsistent", detail: - `config.json records dataDir ${configured} but ${dataDir} is a real ` + + `config.json records mediaDir ${mediaDir} but ${mediaLink} is a real ` + `directory — the media was never moved, or was moved back by hand`, }; } - return { dataDir, relocated: false, status: "in-place" }; + // Tiered in place: `media/` is a real directory on the corpus disk. + return { ...base, relocated: false, text: textRecord, status: "in-place" }; } +// THE MEDIA GUARD, for a job that opens or writes a BIG file. // "ok" and "in-place" pass; everything else throws. An in-transition or // inconsistent channel is refused for the same reason an unreachable one is: // the caller would otherwise read a half-populated or empty dir as the truth. @@ -500,18 +672,10 @@ async function inspectOnDisk( // // ALWAYS FRESH: this is the start-of-work guard, and a remembered "ok" from a // few seconds ago is not what a job about to read `data/` should be told. -// -// THIS CHECK HAS A TWIN. `checkChannelReachable` in -// `umtool/report-to-video/cues.mjs` repeats the same statuses in plain `.mjs`, -// because umtool's bins run under bare node with no `tsx` and cannot import -// this module. It guards the same failure: a relocated channel whose drive is -// not mounted reads as ENOENT, and the cue resolver would otherwise answer from -// the published archive — cutting clips from a snapshot's cues instead of the -// corpus's. Change the checks here and change them there. export async function assertChannelMediaReachable( paths: ChannelMediaPaths, slug: string, - config?: Pick<ChannelConfig, "dataDir"> | null, + config?: ChannelMediaConfig | null, ): Promise<ChannelMediaLocation> { const location = await inspectChannelMedia(paths, slug, config, { fresh: true, @@ -521,3 +685,39 @@ export async function assertChannelMediaReachable( } throw new ChannelMediaUnreachableError(slug, location); } + +// THE TEXT GUARD (release 17): passes for every channel whose `data/` is a +// readable directory on the corpus disk — whatever its MEDIA is doing. A +// stalled, unreachable, inconsistent or moving media tier does not hold a +// reader of the text: the index and stats builds, the snapshot, the digests, +// normalize. Refused: a `legacy` channel (its text is on the far drive), a +// `data/` that is not a directory, a tier migration in flight. +// +// ALWAYS FRESH, for the reason the media guard is. +// +// THIS CHECK HAS A TWIN. `checkChannelReachable` in +// `umtool/report-to-video/cues.mjs` repeats it in plain `.mjs` (refuse a +// legacy channel — a `data` link or a recorded `dataDir`; refuse a marker only +// when its `scope` is `tier-migration`), because umtool's bins run under bare +// node with no `tsx` and cannot import this module. The cue resolver reads +// TEXT; on a legacy channel whose drive is not mounted the cues read as +// ENOENT, and it would otherwise answer from the published archive — cutting +// clips from a snapshot's cues instead of the corpus's. Change the checks here +// and change them there. +export async function assertChannelTextReadable( + paths: ChannelMediaPaths, + slug: string, + config?: ChannelMediaConfig | null, +): Promise<ChannelMediaLocation> { + const location = await inspectChannelMedia(paths, slug, config, { + fresh: true, + }); + if (location.text.readable) return location; + const detail = + location.status === "legacy" + ? (location.detail ?? legacyDetail(slug)) + : location.status === "in-transition" + ? `its media layout is being migrated (${location.detail ?? "a marker is present"})` + : `${location.dataDir} is not a readable directory`; + throw new ChannelTextUnreadableError(slug, location, detail); +} diff --git a/common/lib/channelMediaHold.ts b/common/lib/channelMediaHold.ts @@ -7,26 +7,39 @@ import { type StorageLocation, } from "./storageLocations"; -// THE HOLD, in the words both pool-wide builds use. +// THE HOLD, in the words the builds, the lanes and the surfaces use. // // The index build (controller/buildIndex.ts) and the stats build -// (controller/buildStats.ts) each walk every channel's `data/`. A channel whose -// media cannot be read — a relocated `data/` on an unmounted drive, a move in -// progress, a link and a config that disagree — is HELD by both: not rescanned, -// and what the last build knew of it kept, rather than read as a channel with -// no videos and removed (lib/channelMedia.ts says why that reading is the -// dangerous one). This module is only the shared vocabulary: which statuses +// (controller/buildStats.ts) each walk every channel's `data/`. Since release 17 +// they read the TEXT tier only and are never held by the media tier: a channel +// whose text cannot be read — the retired whole-directory layout (`legacy`), +// a `data/` that is not a directory — is HELD by both: not rescanned, and what +// the last build knew of it kept, rather than read as a channel with no videos +// and removed (lib/channelMedia.ts says why that reading is the dangerous one). +// The media hold (an unmounted media drive, a move in progress, a link and a +// config that disagree) holds the lanes and jobs that open a big file. This module is only the shared vocabulary: which statuses // hold, why, in words with no path in them (/storage shows the paths), and the // ways out a refusal names. Each build decides for itself what "kept" means. // // Pure: no I/O. The caller asks inspectChannelMedia and passes the answer in. // "ok" and "in-place" are read; every other status holds, including any a later -// inspectChannelMedia adds. +// inspectChannelMedia adds (`legacy` among them). The MEDIA hold: what the +// transcription, download and backfill lanes and every media job key on. export function isMediaHeld(status: ChannelMediaStatus): boolean { return status !== "ok" && status !== "in-place"; } +// THE TEXT HOLD (release 17): only the retired whole-directory layout holds the +// text, because only there is it on another drive. The digest lane keys on +// this, so a digest runs on a channel whose media is moving, stalled or +// unmounted. (The text guard, `assertChannelTextReadable`, also refuses a +// `data/` that is not a directory and a tier migration in flight — conditions +// a status alone does not carry; a lane that holds a location asks it.) +export function isTextHeld(status: ChannelMediaStatus): boolean { + return status === "legacy"; +} + // Why a channel is held, without the paths inspectChannelMedia's `detail` // carries. // @@ -40,6 +53,8 @@ export const HELD_REASON: Record<ChannelMediaStatus, string> = { "in-transition": "its media is moving (a move is in progress or was interrupted)", inconsistent: "its data link and its config disagree", stalled: "its drive is not answering (a stalled disk)", + legacy: + "its media layout is the retired whole-directory one (run archilyzer storage migrate-tier)", ok: "reachable", "in-place": "reachable", }; @@ -54,14 +69,15 @@ export function mediaHoldText(status: ChannelMediaStatus): string | null { } // The reason, and the storage location's label when the channel's media is on -// one. `dataDir` is the channel config's; the inspector's target stands in when -// the config names none (a move in flight). +// one. `mediaDir` is the channel config's (`mediaDir`, or the retired `dataDir` +// on a legacy channel); the inspector's target stands in when the config names +// none (a move in flight). export function heldReason( media: Pick<ChannelMediaLocation, "status" | "target">, - dataDir: string | undefined, + mediaDir: string | undefined, locations: StorageLocation[], ): string { - const label = locationLabelOfDataDir(dataDir ?? media.target, locations); + const label = locationLabelOfDataDir(mediaDir ?? media.target, locations); return `${HELD_REASON[media.status]}${label ? `, on location "${label}"` : ""}`; } diff --git a/common/lib/fileSchemaDocs.ts b/common/lib/fileSchemaDocs.ts @@ -162,7 +162,7 @@ export function renderChannelMarkdown(): string { "Every other key is optional and has NO default of its own: an absent " + "key means whatever its description says — for the per-channel " + "overrides, inherit the global setting of the same name; for `name`, " + - "`url`, `dataDir`, `subLangs` and the sync-state stamps, simply unset. " + + "`url`, `mediaDir`, `subLangs` and the sync-state stamps, simply unset. " + "So an ill-typed or out-of-range value is not coerced — it is DROPPED, " + "as if the file did not spell it. Unknown keys (including the retired " + "`excludeFromSync`, now a paused `sync` tier in the channel-priority " + @@ -177,7 +177,7 @@ export function renderChannelMarkdown(): string { "and changes only its own keys, so a stamp and a form save made at once " + "in the editor both land. The two exceptions write a whole config, and " + "only when there is no readable file to patch: a media move and a " + - "channel rename record `dataDir` from their own copy of the config.", + "channel rename record `mediaDir` from their own copy of the config.", ); out.push(""); out.push(REGENERATE); diff --git a/common/lib/mediaTier-server.test.ts b/common/lib/mediaTier-server.test.ts @@ -0,0 +1,340 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { + lstat, + mkdir, + stat, + mkdtemp, + readFile, + readlink, + readdir, + rm, + symlink, + utimes, + writeFile, +} from "node:fs/promises"; +import { existsSync } from "node:fs"; +import { tmpdir } from "node:os"; +import path from "node:path"; +import { RELOCATION_MARKER_FILENAME } from "./channelMedia"; +import { + TIER_RELOCATION_MARKER, + channelMediaLink, + relocatedMediaDir, + removeMediaFile, + removeVideoDirMedia, + tierChannelMedia, + tierLinkTarget, + tierMediaFile, + tierVideoDir, +} from "./mediaTier-server"; + +// Run with: +// pnpm --filter yt-dlp-transcript-common exec tsx --test lib/mediaTier-server.test.ts +// +// Everything happens inside one mkdtemp; nothing reads a real corpus. + +async function fixture() { + const root = await mkdtemp(path.join(tmpdir(), "media-tier-")); + const channelsDir = path.join(root, "channels"); + const slug = "chan"; + const videoDir = path.join(channelsDir, slug, "data", "vid1"); + await mkdir(videoDir, { recursive: true }); + await writeFile(path.join(videoDir, "audio.mp3"), "AUDIO"); + await writeFile(path.join(videoDir, "transcript.json"), "{}"); + await writeFile(path.join(videoDir, "transcript.live_chat.json"), "CHAT"); + await writeFile(path.join(videoDir, "source-media.mp4"), "VIDEO"); + return { root, channelsDir, slug, videoDir, paths: { channelsDir } }; +} + +test("names: the link, the relocated root, the relative link target", () => { + assert.equal(channelMediaLink({ channelsDir: "/c" }, "s"), "/c/s/media"); + assert.equal(relocatedMediaDir(" /mnt/p ", "s"), "/mnt/p/s/media"); + assert.equal(tierLinkTarget("v", "audio.mp3"), "../../media/v/audio.mp3"); +}); + +test("a classic channel (no media/) leaves every file real", async () => { + const f = await fixture(); + assert.equal(await tierMediaFile(f.videoDir, "audio.mp3"), "left"); + assert.ok((await lstat(path.join(f.videoDir, "audio.mp3"))).isFile()); + assert.equal(existsSync(channelMediaLink(f.paths, f.slug)), false); + await rm(f.root, { recursive: true, force: true }); +}); + +test("tiered in place: the bytes move to media/<id>, a relative link stays", async () => { + const f = await fixture(); + await mkdir(channelMediaLink(f.paths, f.slug)); + assert.equal(await tierMediaFile(f.videoDir, "audio.mp3"), "tiered"); + const link = path.join(f.videoDir, "audio.mp3"); + assert.ok((await lstat(link)).isSymbolicLink()); + assert.equal(await readlink(link), "../../media/vid1/audio.mp3"); + assert.equal(await readFile(link, "utf8"), "AUDIO"); + assert.equal( + await readFile(path.join(f.channelsDir, f.slug, "media", "vid1", "audio.mp3"), "utf8"), + "AUDIO", + ); + // Called again: already a link. + assert.equal(await tierMediaFile(f.videoDir, "audio.mp3"), "already"); + // No temp link left behind. + assert.deepEqual( + (await readdir(f.videoDir)).filter((n) => n.includes("tierlink")), + [], + ); + await rm(f.root, { recursive: true, force: true }); +}); + +test("a missing name is already; text and source-media are left", async () => { + const f = await fixture(); + await mkdir(channelMediaLink(f.paths, f.slug)); + assert.equal(await tierMediaFile(f.videoDir, "audio.m4a"), "already"); + assert.equal(await tierMediaFile(f.videoDir, "transcript.json"), "left"); + assert.equal(await tierMediaFile(f.videoDir, "source-media.mp4"), "left"); + assert.ok((await lstat(path.join(f.videoDir, "source-media.mp4"))).isFile()); + await rm(f.root, { recursive: true, force: true }); +}); + +test("tierVideoDir tiers the audio and the raw live chat, nothing else", async () => { + const f = await fixture(); + await mkdir(channelMediaLink(f.paths, f.slug)); + const c = await tierVideoDir(f.videoDir); + assert.deepEqual(c, { tiered: 2, left: 0, already: 0 }); + assert.ok((await lstat(path.join(f.videoDir, "transcript.live_chat.json"))).isSymbolicLink()); + assert.ok((await lstat(path.join(f.videoDir, "transcript.json"))).isFile()); + assert.ok((await lstat(path.join(f.videoDir, "source-media.mp4"))).isFile()); + assert.deepEqual(await tierVideoDir(f.videoDir), { tiered: 0, left: 0, already: 2 }); + await rm(f.root, { recursive: true, force: true }); +}); + +test("relocated: the media link points at another root; EXDEV copies (with the file's times), then links", async () => { + const f = await fixture(); + const platter = path.join(f.root, "platter"); + const target = relocatedMediaDir(platter, f.slug); + await mkdir(target, { recursive: true }); + await symlink(target, channelMediaLink(f.paths, f.slug)); + let copies = 0; + const result = await tierMediaFile(f.videoDir, "audio.mp3", { + fs: { + link: async () => { + const err = new Error("cross-device link not permitted") as NodeJS.ErrnoException; + err.code = "EXDEV"; + throw err; + }, + copyFileAtomic: async (src, dest) => { + copies += 1; + await writeFile(dest, await readFile(src)); + }, + }, + }); + assert.equal(result, "tiered"); + assert.equal(copies, 1); + const link = path.join(f.videoDir, "audio.mp3"); + assert.ok((await lstat(link)).isSymbolicLink()); + assert.equal(await readFile(link, "utf8"), "AUDIO"); + assert.equal(await readFile(path.join(target, "vid1", "audio.mp3"), "utf8"), "AUDIO"); + // stat (the copy) and lstat (the link) agree on the mtime. + assert.equal((await stat(link)).mtimeMs, (await lstat(link)).mtimeMs); + await rm(f.root, { recursive: true, force: true }); +}); + +test("a failed move is undone: the name is a real file again, no link", async () => { + const f = await fixture(); + await mkdir(channelMediaLink(f.paths, f.slug)); + const logs: string[] = []; + const result = await tierMediaFile(f.videoDir, "audio.mp3", { + onLog: (s) => logs.push(s), + fs: { + link: async () => { + const err = new Error("no space left on device") as NodeJS.ErrnoException; + err.code = "ENOSPC"; + throw err; + }, + }, + }); + assert.equal(result, "left"); + assert.ok((await lstat(path.join(f.videoDir, "audio.mp3"))).isFile()); + assert.match(logs.join(""), /left in place: no space left/); + assert.deepEqual( + (await readdir(f.videoDir)).filter((n) => n.includes("tierlink")), + [], + ); + await rm(f.root, { recursive: true, force: true }); +}); + +test("a dangling media link (unmounted drive) creates nothing and leaves the file", async () => { + const f = await fixture(); + const platter = path.join(f.root, "platter-unmounted"); + const target = relocatedMediaDir(platter, f.slug); + await symlink(target, channelMediaLink(f.paths, f.slug)); + assert.equal(await tierMediaFile(f.videoDir, "audio.mp3"), "left"); + assert.deepEqual(await tierVideoDir(f.videoDir), { tiered: 0, left: 2, already: 0 }); + assert.deepEqual(await tierChannelMedia(f.paths, f.slug), { tiered: 0, left: 0, already: 0 }); + // Nothing materialised under the missing mountpoint. + assert.equal(existsSync(platter), false); + assert.ok((await lstat(path.join(f.videoDir, "audio.mp3"))).isFile()); + await rm(f.root, { recursive: true, force: true }); +}); + +test("the per-video mkdir is not recursive: a media root that vanished is not rebuilt", async () => { + const f = await fixture(); + const platter = path.join(f.root, "platter"); + const target = relocatedMediaDir(platter, f.slug); + await mkdir(target, { recursive: true }); + await symlink(target, channelMediaLink(f.paths, f.slug)); + const calls: unknown[][] = []; + // The drive goes away between the readiness check and the mkdir. + const result = await tierMediaFile(f.videoDir, "audio.mp3", { + fs: { + mkdir: async (...args: unknown[]) => { + calls.push(args); + await rm(platter, { recursive: true, force: true }); + return mkdir(args[0] as string); + }, + }, + }); + assert.equal(result, "left"); + assert.equal(calls.length, 1); + assert.equal(calls[0].length, 1, "mkdir is called without options"); + assert.equal(existsSync(platter), false); + assert.ok((await lstat(path.join(f.videoDir, "audio.mp3"))).isFile()); + await rm(f.root, { recursive: true, force: true }); +}); + +test("tierChannelMedia: createMediaDir makes media/ real; since skips old dirs", async () => { + const f = await fixture(); + const old = path.join(f.channelsDir, f.slug, "data", "vid0"); + await mkdir(old); + await writeFile(path.join(old, "audio.mp3"), "OLD"); + const past = new Date(Date.now() - 3_600_000); + await utimes(old, past, past); + const since = Date.now() - 60_000; + // Without createMediaDir a classic channel tiers nothing. + assert.deepEqual(await tierChannelMedia(f.paths, f.slug, { since }), { + tiered: 0, + left: 0, + already: 0, + }); + const c = await tierChannelMedia(f.paths, f.slug, { since, createMediaDir: true }); + assert.deepEqual(c, { tiered: 2, left: 0, already: 0 }); + assert.ok((await lstat(channelMediaLink(f.paths, f.slug))).isDirectory()); + assert.ok((await lstat(path.join(old, "audio.mp3"))).isFile()); + // Without since, the old dir too. + assert.deepEqual(await tierChannelMedia(f.paths, f.slug), { tiered: 1, left: 0, already: 2 }); + await rm(f.root, { recursive: true, force: true }); +}); + +test("removeMediaFile derefs a link: the bytes in media/<id> go, then the link", async () => { + const f = await fixture(); + await mkdir(channelMediaLink(f.paths, f.slug)); + await tierVideoDir(f.videoDir); + const bytes = path.join(f.channelsDir, f.slug, "media", "vid1", "audio.mp3"); + assert.ok(existsSync(bytes)); + await removeMediaFile(f.videoDir, "audio.mp3"); + assert.equal(existsSync(bytes), false); + await assert.rejects(lstat(path.join(f.videoDir, "audio.mp3"))); + // A real file, and a missing one. + await removeMediaFile(f.videoDir, "transcript.json"); + await assert.rejects(lstat(path.join(f.videoDir, "transcript.json"))); + await removeMediaFile(f.videoDir, "nope"); + await rm(f.root, { recursive: true, force: true }); +}); + +test("removeMediaFile never follows a link out of the channel's media/<id>", async () => { + const f = await fixture(); + const outside = path.join(f.root, "precious.mp3"); + await writeFile(outside, "KEEP"); + await rm(path.join(f.videoDir, "audio.mp3")); + await symlink(outside, path.join(f.videoDir, "audio.mp3")); + await removeMediaFile(f.videoDir, "audio.mp3"); + assert.equal(await readFile(outside, "utf8"), "KEEP"); + await assert.rejects(lstat(path.join(f.videoDir, "audio.mp3"))); + await rm(f.root, { recursive: true, force: true }); +}); + +test("removeVideoDirMedia clears every link's bytes and the video's media dir", async () => { + const f = await fixture(); + await mkdir(channelMediaLink(f.paths, f.slug)); + await tierVideoDir(f.videoDir); + const tierDir = path.join(f.channelsDir, f.slug, "media", "vid1"); + await writeFile(path.join(tierDir, "metadata.info.json"), "{}"); // a leftover copy + await removeVideoDirMedia(f.videoDir); + assert.equal(existsSync(tierDir), false); + // The text is the caller's to remove with the dir. + assert.ok(existsSync(path.join(f.videoDir, "transcript.json"))); + await rm(f.root, { recursive: true, force: true }); +}); + +test("the link carries the file's mtime, so an lstat answers freshness without the media drive", async () => { + const f = await fixture(); + await mkdir(channelMediaLink(f.paths, f.slug)); + const file = path.join(f.videoDir, "transcript.live_chat.json"); + const when = new Date("2026-01-02T03:04:05Z"); + await utimes(file, when, when); + assert.equal(await tierMediaFile(f.videoDir, "transcript.live_chat.json"), "tiered"); + const l = await lstat(file); + assert.ok(l.isSymbolicLink()); + assert.equal(l.mtimeMs, when.getTime()); + await rm(f.root, { recursive: true, force: true }); +}); + +test("a move marker on the channel: the hook writes nothing into media/", async () => { + assert.equal(TIER_RELOCATION_MARKER, RELOCATION_MARKER_FILENAME); + const f = await fixture(); + await mkdir(channelMediaLink(f.paths, f.slug)); + await writeFile(path.join(f.channelsDir, f.slug, RELOCATION_MARKER_FILENAME), "{}"); + assert.equal(await tierMediaFile(f.videoDir, "audio.mp3"), "left"); + assert.deepEqual(await readdir(channelMediaLink(f.paths, f.slug)), []); + assert.ok((await lstat(path.join(f.videoDir, "audio.mp3"))).isFile()); + await rm(f.root, { recursive: true, force: true }); +}); + +test("the name always resolves: a crash after the bytes land leaves the real file, and the next sweep removes the stray link", async () => { + const f = await fixture(); + await mkdir(channelMediaLink(f.paths, f.slug)); + // A process killed between placing the bytes and the rename: its temp link + // is left beside the still-real file (pid 999999 is not alive). + await symlink("../../media/vid1/audio.mp3", path.join(f.videoDir, ".audio.mp3.tierlink-999999")); + assert.ok((await lstat(path.join(f.videoDir, "audio.mp3"))).isFile()); + const c = await tierVideoDir(f.videoDir); + assert.equal(c.tiered, 2); + assert.deepEqual( + (await readdir(f.videoDir)).filter((n) => n.includes("tierlink")), + [], + ); + await rm(f.root, { recursive: true, force: true }); +}); + +test("same filesystem: the bytes are hard-linked into the tier (one inode), then the name is replaced", async () => { + const f = await fixture(); + await mkdir(channelMediaLink(f.paths, f.slug)); + const ino = (await stat(path.join(f.videoDir, "audio.mp3"))).ino; + assert.equal(await tierMediaFile(f.videoDir, "audio.mp3"), "tiered"); + const bytes = path.join(f.channelsDir, f.slug, "media", "vid1", "audio.mp3"); + assert.equal((await stat(bytes)).ino, ino); + assert.equal((await stat(bytes)).nlink, 1, "the original name no longer holds the inode"); + await rm(f.root, { recursive: true, force: true }); +}); + +test("removeMediaFile derefs a link into ANOTHER id's media dir (a renamed video dir), and drops it when empty", async () => { + const f = await fixture(); + const media = channelMediaLink(f.paths, f.slug); + await mkdir(path.join(media, "oldid"), { recursive: true }); + await writeFile(path.join(media, "oldid", "audio.m4a"), "OLD"); + await symlink("../../media/oldid/audio.m4a", path.join(f.videoDir, "audio.m4a")); + await removeMediaFile(f.videoDir, "audio.m4a"); + assert.equal(existsSync(path.join(media, "oldid")), false); + await assert.rejects(lstat(path.join(f.videoDir, "audio.m4a"))); + await rm(f.root, { recursive: true, force: true }); +}); + +test("a dead process's half-placed bytes in media/<id>/ are swept with the video dir", async () => { + const f = await fixture(); + const tierDir = path.join(channelMediaLink(f.paths, f.slug), "vid1"); + await mkdir(tierDir, { recursive: true }); + await writeFile(path.join(tierDir, ".audio.mp3.tiering-999999"), "HALF"); + await tierVideoDir(f.videoDir); + assert.deepEqual( + (await readdir(tierDir)).filter((n) => n.includes("tiering")), + [], + ); + await rm(f.root, { recursive: true, force: true }); +}); diff --git a/common/lib/mediaTier-server.ts b/common/lib/mediaTier-server.ts @@ -0,0 +1,423 @@ +// THE MEDIA TIER ON DISK — the hook that moves a finished media file out of +// `data/<id>/` and leaves a link behind, and the one way to remove one. +// +// The layout (plans/release-17.md, "The model (A′)"): +// +// channels/<slug>/data/<id>/audio.mp3 -> ../../media/<id>/audio.mp3 (a RELATIVE link) +// channels/<slug>/media a real directory on the corpus disk (tiered in place), +// or ONE absolute symlink to <root>/<slug>/media (relocated), +// or absent (classic: the media is real files in data/<id>/) +// +// Readers never change: they open `data/<id>/<name>` and the kernel follows the +// link. Writers change in exactly two ways, and both are here: +// +// - EVERY SITE THAT FINALISES A MEDIA FILE calls the hook (`tierMediaFile`, +// `tierVideoDir`, `tierChannelMedia`) once it has renamed its temp into +// place. yt-dlp's postprocessor and this app's transcode both `rename()` a +// temp OVER the final name, which replaces a link with a real file — the +// hook then moves that file into the tier, over the stale copy there. +// - EVERY SITE THAT DELETES ONE calls `removeMediaFile` / `removeVideoDirMedia`. +// A plain `rm` of `data/<id>/audio.mp3` removes the LINK and orphans the +// bytes on the media drive, which no sweep would ever find again. +// +// THE HOOK NEVER THROWS INTO A DOWNLOAD. A classic channel (no `media/`), a +// relocated one whose drive is unmounted (the link dangles), a stalled drive, +// a full disk: the file stays a real file in `data/<id>/`, readers are +// unaffected, and the next sweep that calls the hook tiers it. +// +// lib/, so no controller import (architecture.test.ts). + +import path from "node:path"; +import { + link, + lstat, + lutimes, + mkdir, + readdir, + readlink, + rename, + rm, + rmdir, + stat, + symlink, + utimes, +} from "node:fs/promises"; +import type { Paths } from "./paths"; +import { copyFileAtomic } from "./jsonFile-server"; +import { isTierable } from "./mediaTier"; +import { onDrive, stalledLocationForPath } from "./storageHealth"; + +// The one name a channel's media tier is reached by: `channels/<slug>/media`. +export const MEDIA_LINK_NAME = "media"; + +// channelMedia.ts's RELOCATION_MARKER_FILENAME, as a literal because that +// module imports this one (mediaTier-server.test.ts pins that they agree). +export const TIER_RELOCATION_MARKER = ".relocating.json"; + +export function channelMediaLink( + paths: Pick<Paths, "channelsDir">, + slug: string, +): string { + return path.join(paths.channelsDir, slug, MEDIA_LINK_NAME); +} + +// A relocated channel's media root: `<root>/<slug>/media`. The suffix is fixed, +// not configurable, so an empty mountpoint can never be mistaken for the media +// and the movers can recognise a target by its shape (the same reason +// `relocatedDataDir` fixed `<slug>/data`). +export function relocatedMediaDir(root: string, slug: string): string { + return path.join(root.trim(), slug, MEDIA_LINK_NAME); +} + +// The link a tiered file leaves in `data/<id>/`: RELATIVE, so it survives a +// channel rename, `reconcileVideoDirs`' renames of a video dir (the link and +// its target move together only when the target's `<id>` is renamed too — open +// question 3 of the plan) and a move back, which makes `media/` a real +// directory without touching a single link. +export function tierLinkTarget(id: string, name: string): string { + return path.join("..", "..", MEDIA_LINK_NAME, id, name); +} + +// `data/<id>` → `channels/<slug>/media/<id>`, by shape (the video dir is +// always `channels/<slug>/data/<id>`). +function mediaDirOfVideoDir(videoDir: string): string { + const id = path.basename(videoDir); + return path.join(path.dirname(path.dirname(videoDir)), MEDIA_LINK_NAME, id); +} + +function errCode(err: unknown): string | undefined { + return (err as NodeJS.ErrnoException | null)?.code; +} + +// Test seam: the filesystem calls a test needs to fail on cue (an injected +// `link` throwing EXDEV, the copy the fallback makes). +export type TierFs = { + link: (existing: string, linkPath: string) => Promise<void>; + copyFileAtomic: (src: string, dest: string) => Promise<void>; + // Called with ONE argument — never `{ recursive: true }` (below). + mkdir: (dir: string) => Promise<unknown>; +}; + +const DEFAULT_FS: TierFs = { + link: (a, b) => link(a, b), + copyFileAtomic: (s, d) => copyFileAtomic(s, d), + mkdir: (d) => mkdir(d), +}; + +// A hard link cannot be made here: another filesystem (a relocated channel), +// or one that has none. +const NO_HARD_LINK = new Set(["EXDEV", "EPERM", "ENOTSUP", "EOPNOTSUPP", "EMLINK"]); + +async function markerStands(mediaRoot: string): Promise<boolean> { + try { + await lstat(path.join(path.dirname(mediaRoot), TIER_RELOCATION_MARKER)); + return true; + } catch { + return false; + } +} + +export type TierOptions = { + onLog?: (line: string) => void; + fs?: Partial<TierFs>; +}; + +// Whether the channel's media tier can take a file now: `channels/<slug>/media` +// resolves to a directory, on a drive that is not known to be stalled and that +// answers within the watchdog's budget. False for a classic channel (no +// `media`), a dangling link (an unmounted drive), a stalled drive. +async function mediaTierReady(mediaRoot: string): Promise<boolean> { + // A move of this channel's media in flight (or interrupted): its marker + // stands in the channel dir, and nothing is written into `media/` while it + // does — the file stays real and the next sweep tiers it. (The writers that + // call the hook are held by the media guard during a move; this is the + // backstop for one that started before the marker.) + if (await markerStands(mediaRoot)) return false; + let target = mediaRoot; + try { + const l = await lstat(mediaRoot); + if (l.isSymbolicLink()) { + target = path.resolve(path.dirname(mediaRoot), await readlink(mediaRoot)); + if (stalledLocationForPath(target)) return false; + } else { + // A real `media/` is on the corpus disk: no drive, no watchdog slot. + return l.isDirectory(); + } + } catch { + return false; + } + try { + const st = await onDrive(target, () => stat(mediaRoot)); + return st.isDirectory(); + } catch { + return false; + } +} + +export type TierResult = "tiered" | "left" | "already"; + +// MOVE ONE FINISHED MEDIA FILE INTO THE TIER and leave a relative link behind. +// +// "already" — the name is a link already, or there is no such file; +// "left" — it stays a real file (not tierable, no media tier, the drive is +// not there or not answering, or a step failed and was undone); +// "tiered" — the bytes are in `media/<id>/<name>` and `data/<id>/<name>` is +// the link. +// +// `mkdir(media/<id>)` is NOT recursive: `media/` was just seen to be a +// directory, and a recursive mkdir aimed at a mountpoint that went away in +// between would build the path on the root filesystem and fill it. +// +// THE NAME ALWAYS RESOLVES, crash or not. The bytes are put into the tier +// WITHOUT touching the name — same filesystem: a hard link (instant, the same +// inode), renamed into place; across filesystems (a relocated channel) or with +// no hard links: an atomic copy, given the file's times so `stat` and `lstat` +// agree — and only then does the temp link replace the name, in one rename. +// A process killed anywhere leaves the real file under its name (and at worst +// a stray temp link, which the next sweep removes). +// +// A MOVE THAT BEGAN MEANWHILE: the marker is asked again after the bytes land +// and before the name changes; if one stands, the tier's copy is removed and +// the file stays real. +export async function tierMediaFile( + videoDir: string, + name: string, + opts: TierOptions = {}, +): Promise<TierResult> { + const fsx: TierFs = { ...DEFAULT_FS, ...opts.fs }; + const file = path.join(videoDir, name); + let st; + try { + st = await lstat(file); + } catch { + return "already"; + } + if (st.isSymbolicLink()) return "already"; + if (!st.isFile() || !isTierable(name)) return "left"; + + const mediaRoot = path.dirname(mediaDirOfVideoDir(videoDir)); + if (!(await mediaTierReady(mediaRoot))) return "left"; + + const id = path.basename(videoDir); + const destDir = mediaDirOfVideoDir(videoDir); + const dest = path.join(destDir, name); + try { + await fsx.mkdir(destDir); + } catch (err) { + if (errCode(err) !== "EEXIST") { + opts.onLog?.(`[media tier] ${id}/${name} left in place: ${(err as Error).message}\n`); + return "left"; + } + } + + const tmpLink = path.join(videoDir, `.${name}.tierlink-${process.pid}`); + const destTmp = path.join(destDir, `.${name}.tiering-${process.pid}`); + let placed = false; + try { + await rm(tmpLink, { force: true }); + await symlink(tierLinkTarget(id, name), tmpLink); + // The FILE's times on the LINK, so an `lstat` of the name answers what a + // `stat` of the file did before it was tiered — the live-chat freshness + // check and the index's sub-track key read it there, on the corpus disk, + // instead of reaching the media drive. + await lutimes(tmpLink, st.atime, st.mtime).catch(() => {}); + await rm(destTmp, { force: true }); + try { + await fsx.link(file, destTmp); + await rename(destTmp, dest); + } catch (err) { + await rm(destTmp, { force: true }).catch(() => {}); + if (!NO_HARD_LINK.has(errCode(err) ?? "")) throw err; + await fsx.copyFileAtomic(file, dest); + await utimes(dest, st.atime, st.mtime).catch(() => {}); + } + placed = true; + if (await markerStands(mediaRoot)) { + throw new Error("a move of this channel's media began"); + } + await rename(tmpLink, file); + return "tiered"; + } catch (err) { + await rm(tmpLink, { force: true }).catch(() => {}); + // The name still holds the real file; the tier's copy goes. + if (placed) await rm(dest, { force: true }).catch(() => {}); + opts.onLog?.(`[media tier] ${id}/${name} left in place: ${(err as Error).message}\n`); + return "left"; + } +} + +function processAlive(pid: number): boolean { + if (pid === process.pid) return true; + try { + process.kill(pid, 0); + return true; + } catch (err) { + return errCode(err) === "EPERM"; + } +} + +export type TierCounts = { tiered: number; left: number; already: number }; + +function emptyCounts(): TierCounts { + return { tiered: 0, left: 0, already: 0 }; +} + +// Every tierable file in one video dir. A missing dir counts nothing. +export async function tierVideoDir( + videoDir: string, + opts: TierOptions = {}, +): Promise<TierCounts> { + const counts = emptyCounts(); + let names: string[]; + try { + names = await readdir(videoDir); + } catch { + return counts; + } + for (const name of names) { + // A stray temp link from a process that is gone (killed between placing + // the bytes and the rename): the name still holds the file, so it is only + // removed. + const stray = /^\..+\.tierlink-(\d+)$/.exec(name); + if (stray) { + if (!processAlive(Number(stray[1]))) { + await rm(path.join(videoDir, name), { force: true }).catch(() => {}); + } + continue; + } + if (!isTierable(name)) continue; + counts[await tierMediaFile(videoDir, name, opts)] += 1; + } + // The media side's strays: bytes a dead process was placing + // (`.<name>.tiering-<pid>` in `media/<id>/`). Only while the tier answers. + const destDir = mediaDirOfVideoDir(videoDir); + if (await mediaTierReady(path.dirname(destDir))) { + for (const name of await readdir(destDir).catch(() => [] as string[])) { + const stray = /^\..+\.tiering-(\d+)$/.exec(name); + if (stray && !processAlive(Number(stray[1]))) { + await rm(path.join(destDir, name), { force: true }).catch(() => {}); + } + } + } + return counts; +} + +export type TierChannelOptions = TierOptions & { + // Only video dirs whose mtime is at or after this instant (ms since epoch) — + // a batch download's run start: a dir yt-dlp wrote into had an entry added + // or renamed, which moves its mtime. + since?: number; + // Make `channels/<slug>/media` a real directory when neither a link nor a + // directory is there (the mover's preflight tiers a classic channel in + // place, same filesystem, before it copies `media/`). + createMediaDir?: boolean; +}; + +// Every video dir of a channel. A classic channel without `createMediaDir` +// tiers nothing (every file is "left" — no work is attempted, and none is +// counted). Never throws. +export async function tierChannelMedia( + paths: Pick<Paths, "channelsDir">, + slug: string, + opts: TierChannelOptions = {}, +): Promise<TierCounts> { + const counts = emptyCounts(); + const mediaRoot = channelMediaLink(paths, slug); + if (opts.createMediaDir) { + try { + await lstat(mediaRoot); + } catch (err) { + if (errCode(err) === "ENOENT") { + await mkdir(mediaRoot).catch(() => {}); + } + } + } + if (!(await mediaTierReady(mediaRoot))) return counts; + const dataDir = path.join(paths.channelsDir, slug, "data"); + let ids: string[]; + try { + ids = await readdir(dataDir); + } catch { + return counts; + } + for (const id of ids) { + const videoDir = path.join(dataDir, id); + if (opts.since !== undefined) { + try { + const st = await stat(videoDir); + if (!st.isDirectory() || st.mtimeMs < opts.since) continue; + } catch { + continue; + } + } + const c = await tierVideoDir(videoDir, opts); + counts.tiered += c.tiered; + counts.left += c.left; + counts.already += c.already; + } + return counts; +} + +// REMOVE ONE FILE FROM A VIDEO DIR, through its link when it is one: the link's +// target is removed first (only when it resolves inside this channel's +// `media/<id>/` — a link pointing anywhere else is removed alone, never +// followed out), then the name. Missing is not an error. +// +// THE ONE `rm` OF A VIDEO-DIR ENTRY. Every deleter in common/ and the editor +// goes through here (the grep gate in plans/release-17.md), text sidecars +// included, so no call site has to know which of its names may be a link. +export async function removeMediaFile(videoDir: string, name: string): Promise<void> { + const file = path.join(videoDir, name); + let st; + try { + st = await lstat(file); + } catch { + return; + } + if (st.isSymbolicLink()) { + // Any id's dir under the channel's own `media/`: a video dir renamed by + // reconcileVideoDirs keeps links into `media/<itsOldId>/`. Never followed + // out of `media/`. + const mediaRoot = path.dirname(mediaDirOfVideoDir(videoDir)); + let target = ""; + try { + target = path.resolve(videoDir, await readlink(file)); + } catch { + /* unreadable link: removed alone below */ + } + if (target && path.dirname(path.dirname(target)) === mediaRoot) { + await rm(target, { force: true }); + // Its dir goes when it empties (a non-recursive rmdir refuses otherwise). + await rmdir(path.dirname(target)).catch(() => {}); + } + } + await rm(file, { force: true }); +} + +// Before a whole video dir is deleted: every link's target in it, then the +// video's own `media/<id>/` (whatever else is left there — the migration's +// platter copies before a `--reclaim`). The caller removes the dir itself. +export async function removeVideoDirMedia(videoDir: string): Promise<void> { + let names: string[] = []; + try { + names = await readdir(videoDir); + } catch { + /* no dir: only the tier's side to clear */ + } + for (const name of names) { + let st; + try { + st = await lstat(path.join(videoDir, name)); + } catch { + continue; + } + if (st.isSymbolicLink()) await removeMediaFile(videoDir, name); + } + const tierDir = mediaDirOfVideoDir(videoDir); + try { + const l = await lstat(tierDir); + if (l.isDirectory()) await rm(tierDir, { recursive: true, force: true }); + } catch { + /* no tier dir for this video */ + } +} diff --git a/common/lib/mediaTier.test.ts b/common/lib/mediaTier.test.ts @@ -0,0 +1,88 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { + CLIPS_DIR, + LIVE_CHAT_MEDIA_FILENAME, + classifyEntry, + classifyVideoDir, + isTierable, +} from "./mediaTier"; +import { LIVE_CHAT_FILENAME } from "./videoStatus"; +import { CLIPS_DIR_NAME } from "./clipWindow"; + +// Run with: +// pnpm --filter yt-dlp-transcript-common exec tsx --test lib/mediaTier.test.ts + +test("the classifier's two copied names agree with their owners", () => { + assert.equal(LIVE_CHAT_MEDIA_FILENAME, LIVE_CHAT_FILENAME); + assert.equal(CLIPS_DIR, CLIPS_DIR_NAME); +}); + +// THE TABLE. By name, never by size; every row is a name this corpus holds. +const TABLE: ReadonlyArray<[string, "media" | "text" | "scratch", boolean]> = [ + // [name, tier, tierable] + ["audio.mp3", "media", true], + ["audio.m4a", "media", true], + ["audio.opus", "media", true], + ["audio.mp4", "media", true], + ["audio.webm", "media", true], + ["source-media.mp4", "media", false], + ["source-media.webm", "media", false], + ["transcript.live_chat.json", "media", true], + // A subtitle named like audio is TEXT (the isRealAudioFile anchoring). + ["audio.en-orig.vtt", "text", false], + ["audio.en.vtt", "text", false], + // Somebody's scratch. + ["audio.tmp-2760235.mp3", "scratch", false], + ["source-media.temp.mp4", "scratch", false], + ["audio.temp.mp3", "scratch", false], + ["audio.m4a.part", "scratch", false], + ["audio.m4a.part.good", "scratch", false], + ["audio.m4a.part.testing", "scratch", false], + ["audio.live_chat.json.part-Frag114", "scratch", false], + ["transcript.live_chat.json.part", "scratch", false], + ["audio.m4a.ytdl", "scratch", false], + [".audio.mp3.parakeet", "scratch", false], + [".audio.mp3.tierlink-4242", "scratch", false], + [".audio.mp3.tiering-4242", "scratch", false], + // The hot text. + ["transcript.json", "text", false], + ["transcript.en.vtt", "text", false], + ["transcript.en-orig.vtt", "text", false], + ["transcript.cues.json", "text", false], + ["live_chat.cues.json", "text", false], + ["metadata.info.json", "text", false], + ["metadata.history.json", "text", false], + ["diarization.json", "text", false], + ["saved-video.json", "text", false], + ["clips", "text", false], + ["thumbnail.jpg", "text", false], +]; + +for (const [name, tier, tierable] of TABLE) { + test(`classifyEntry(${name}) = ${tier}, tierable ${tierable}`, () => { + assert.equal(classifyEntry(name), tier); + assert.equal(isTierable(name), tierable); + }); +} + +test("classifyVideoDir splits a listing by tier, in order", () => { + const out = classifyVideoDir([ + "metadata.info.json", + "audio.mp3", + "audio.tmp-1.mp3", + "transcript.json", + "transcript.live_chat.json", + ]); + assert.deepEqual(out, { + media: ["audio.mp3", "transcript.live_chat.json"], + text: ["metadata.info.json", "transcript.json"], + scratch: ["audio.tmp-1.mp3"], + }); +}); + +test("every tierable name is media (tierable is narrower)", () => { + for (const [name] of TABLE) { + if (isTierable(name)) assert.equal(classifyEntry(name), "media", name); + } +}); diff --git a/common/lib/mediaTier.ts b/common/lib/mediaTier.ts @@ -0,0 +1,114 @@ +// THE MEDIA TIER'S CLASSIFIER — which files in a video dir are big and cold, +// which are the hot text, and which are somebody's scratch. +// +// Release 17 splits a channel's files across two tiers: the TEXT (transcripts, +// cues, `metadata.info.json`, every sidecar) stays in `channels/<slug>/data/<id>/` +// on the corpus disk, and the MEDIA (the audio, a persisted source container, +// the raw live-chat replay) may live under `channels/<slug>/media/<id>/` — a +// real directory, or one absolute symlink to `<root>/<slug>/media` on another +// drive — with a RELATIVE per-file link left in `data/<id>/` so every reader +// keeps opening the same path (`lib/mediaTier-server.ts` holds the links). +// +// BY NAME, NEVER BY SIZE. A file's tier is a function of its name alone, over +// `mediaFiles.ts`'s anchored predicates, so the answer is the same for a file +// half-written, a file on an unmounted drive (whose size nobody can read) and a +// file in a listing a bundle tool reads off another machine. Pure: no fs, no +// import that brings one in — the channel export/import bundle (the slice after +// release 17) reuses this module, and a client may too. +// +// `transcript.live_chat.json` IS MEDIA. It is read once, by +// `normalizeLiveChat`, which derives the small `live_chat.cues.json` every +// other reader uses; the raw replay is tens of GB on the big channels. The +// `clips/` cache is NEVER tiered: it stays on the corpus disk and is evicted by +// age (`evictClipWindows`). + +import { + isPartAudioFile, + isRealAudioFile, + isSourceMediaFile, +} from "./mediaFiles"; + +// The raw live-chat replay's name. `videoStatus.ts` exports the same string as +// `LIVE_CHAT_FILENAME`, but that module imports `node:fs`; mediaTier.test.ts +// pins that the two agree. +export const LIVE_CHAT_MEDIA_FILENAME = "transcript.live_chat.json"; + +// The clip-window cache dir (`clipWindow.ts`'s `CLIPS_DIR_NAME`, pinned by the +// test for the same reason). Never tiered, never classified as media. +export const CLIPS_DIR = "clips"; + +export type MediaTierKind = "media" | "text" | "scratch"; + +// Somebody's in-flight bytes — a downloader's partial, a transcoder's temp, a +// transcriber's window dir. Never tiered (the writer is about to rename it, or +// to resume it), never counted as text, and never carried by a copy that +// rebuilds a `data/` (the migration lists text and scratch separately so a +// verify can say which is which). +const SCRATCH_PATTERNS: ReadonlyArray<RegExp> = [ + // This app's transcode temp: `audio.tmp-<pid>.<ext>` (controller/transcode.ts). + /^audio\.tmp-\d+\./, + // parakeet's per-file window scratch dir: `.audio.<ext>.parakeet`. + /^\.audio\..*\.parakeet$/, + // yt-dlp's fragment downloads: `<name>.part-Frag<n>`. + /\.part-Frag\d+$/, + // yt-dlp's postprocessor temp: `audio.temp.mp3`, `source-media.temp.mp4`. + /\.temp\./, + // The audio check's snapshots of a partial: `audio.m4a.part.good`, `.part.testing`. + /\.part\.(good|testing)$/, + // Any other downloader partial (`transcript.live_chat.json.part`) and + // yt-dlp's resume-state file beside one (`audio.m4a.ytdl`). + /\.part$/, + /\.ytdl$/, + // The media-tier hook's own temps (`lib/mediaTier-server.ts`): the link + // beside the name, and the bytes being placed in `media/<id>/`. + /\.tierlink-\d+$/, + /\.tiering-\d+$/, +]; + +export function isScratchEntry(name: string): boolean { + if (isPartAudioFile(name)) return true; + return SCRATCH_PATTERNS.some((re) => re.test(name)); +} + +// Which tier a video-dir entry belongs to. Scratch first: `source-media.temp.mp4` +// and `audio.tmp-2760235.mp3` are already rejected by the anchored media +// predicates, and a scratch rule that matched a finalized name would be a bug +// the test table catches. +export function classifyEntry(name: string): MediaTierKind { + if (name === CLIPS_DIR) return "text"; + if (isScratchEntry(name)) return "scratch"; + if ( + isRealAudioFile(name) || + isSourceMediaFile(name) || + name === LIVE_CHAT_MEDIA_FILENAME + ) { + return "media"; + } + return "text"; +} + +// What the hook actually moves into the media tier: NARROWER than "media". +// +// - `source-media.*` stays a real file in `data/<id>/`: `persistSourceVideo` +// `rename`s it into the saved-video store, and a rename of a LINK would move +// the link, leaving the bytes behind in `media/` and a dangling pointer in +// the store. The store is already its own per-object tier. +// - `audio.*.part` stays real: it is yt-dlp's resumable partial, which yt-dlp +// appends to and renames. +export function isTierable(name: string): boolean { + return isRealAudioFile(name) || name === LIVE_CHAT_MEDIA_FILENAME; +} + +export type ClassifiedVideoDir = { + media: string[]; + text: string[]; + scratch: string[]; +}; + +// One video dir's entries (a `readdir`'s names), by tier, each list in the +// order given. +export function classifyVideoDir(entries: Iterable<string>): ClassifiedVideoDir { + const out: ClassifiedVideoDir = { media: [], text: [], scratch: [] }; + for (const name of entries) out[classifyEntry(name)].push(name); + return out; +} diff --git a/common/lib/storageHealth.test.ts b/common/lib/storageHealth.test.ts @@ -681,3 +681,11 @@ test("DT review L4: a timeout on no known location names the budget the call ran setDriveCallBudget(5_000); assert.equal(await refused, "drive not answering (a read did not answer within 0.08 s)"); }); + +test("rootOfUnknownPath strips <root>/<slug>/media (release 17) and the retired <root>/<slug>/data", async () => { + const { rootOfUnknownPath } = await import("./storageHealth"); + assert.equal(rootOfUnknownPath("/mnt/p/chan/media"), "/mnt/p"); + assert.equal(rootOfUnknownPath("/mnt/p/chan/media/"), "/mnt/p"); + assert.equal(rootOfUnknownPath("/mnt/p/chan/data"), "/mnt/p"); + assert.equal(rootOfUnknownPath("/mnt/p/chan/other"), "/mnt/p/chan/other"); +}); diff --git a/common/lib/storageHealth.ts b/common/lib/storageHealth.ts @@ -524,12 +524,14 @@ function slotKeyOfLocation(id: string): string { return `loc:${id}`; } -// The root a path on no configured location is under: `<root>/<slug>/data` -// (relocatedDataDir's shape) gives `<root>`; anything else is its own key. -function rootOfUnknownPath(p: string): string { +// The root a path on no configured location is under: `<root>/<slug>/media` +// (relocatedMediaDir's shape, release 17) or the retired `<root>/<slug>/data` +// gives `<root>`; anything else is its own key. +export function rootOfUnknownPath(p: string): string { const clean = p.replace(/\/+$/, ""); const parts = clean.split("/"); - return parts.length > 2 && parts[parts.length - 1] === "data" + const last = parts[parts.length - 1]; + return parts.length > 2 && (last === "media" || last === "data") ? parts.slice(0, -2).join("/") || "/" : clean; } diff --git a/common/ytdlp/audioCheckedDownload.ts b/common/ytdlp/audioCheckedDownload.ts @@ -17,6 +17,7 @@ // the probe) and promotes it to `.good` on success. SIGCONT is always sent in a // finally to avoid orphaning a suspended child. +import { removeMediaFile } from "../lib/mediaTier-server"; import { constants as fsConstants } from "node:fs"; import { isPartAudioFile, isRealAudioFile } from "../lib/mediaFiles"; import { @@ -803,7 +804,7 @@ export async function runAudioCheckedYtdlp( // container if !keepSourceVideo. await rm(goodPath(partPathFor(finalFile)), { force: true }); if (!opts.channelConfig.keepSourceVideo) { - await rm(finalFile, { force: true }); + await removeMediaFile(path.dirname(finalFile), path.basename(finalFile)); } return buildOutcome("ok"); } diff --git a/common/ytdlp/downloadOneManaged.ts b/common/ytdlp/downloadOneManaged.ts @@ -1,3 +1,5 @@ +import { removeMediaFile } from "../lib/mediaTier-server"; +import { tierVideoDir } from "../lib/mediaTier-server"; import path from "node:path"; import { appendFile, mkdir, readdir, readFile, rm, stat } from "node:fs/promises"; import { createWriteStream, type Dirent, type WriteStream } from "node:fs"; @@ -290,7 +292,7 @@ async function finalizeAppExtraction(opts: { return; } if (!opts.persist) { - await rm(path.join(opts.videoDir, source), { force: true }); + await removeMediaFile(opts.videoDir, source); opts.onLog(`Discarded source container ${source} (audio-only).\n`); return; } @@ -998,6 +1000,9 @@ async function runManagedDownload( } else { try { await mkdir(videoDir, { recursive: true }); + // THE MEDIA TIER'S HOOK (release 17): what this run finalised moves + // into channels/<slug>/media when the channel has one. Never throws. + await tierVideoDir(videoDir, { onLog: opts.onLog }); await writeDownloadOutcome(videoDir, record); } catch (err) { opts.onLog( @@ -1626,6 +1631,12 @@ async function writeOutcome( // exist yet; in that case `mkdir -p` it so the sidecar lands somewhere. try { await mkdir(videoDir, { recursive: true }); + // THE MEDIA TIER'S HOOK (release 17), after reconcileVideoDirs: what this + // run finalised — yt-dlp's own `-x` output, the app's extraction, the live + // chat — moves into channels/<slug>/media when the channel has one, and a + // relative link stays. Never throws: on a classic channel, or a media + // drive that is not there, the files stay real. + await tierVideoDir(videoDir, { onLog: opts.onLog }); await writeDownloadOutcome(videoDir, record); } catch (err) { opts.onLog( diff --git a/common/ytdlp/runYtdlp.ts b/common/ytdlp/runYtdlp.ts @@ -1,3 +1,4 @@ +import { tierChannelMedia } from "../lib/mediaTier-server"; import path from "node:path"; import { mkdir, readdir, readFile } from "node:fs/promises"; import { writeFileAtomic } from "../lib/jsonFile-server"; @@ -162,10 +163,48 @@ export type RunYtdlpOpts = { ) => void | Promise<void>; }; +// The batch modes whose child writes media into `data/<id>/` — yt-dlp's own +// `-x` output, the batch download loops — and so end with the media tier's +// hook. Not store-playlist (writes `playlist`) and not download-missing-subs +// (writes subtitles, which are text). +const MEDIA_WRITING_MODES: ReadonlySet<RunYtdlpOpts["mode"]> = new Set([ + "download-from-playlist", + "download-missing", + "download-one-audio", + "retry-bucket", + "sync", +]); + +// A video dir's mtime is compared with the run's start; a filesystem with +// coarse timestamps rounds down, so the window opens this much earlier. +const TIER_SINCE_SLACK_MS = 2_000; + export async function runYtdlp(opts: RunYtdlpOpts): Promise<void> { if (!opts.channelConfig.url) { throw new Error("Channel has no `url` configured"); } + // A shard save computes a slice and writes it; it downloads nothing. + if (!MEDIA_WRITING_MODES.has(opts.mode) || opts.saveShardOnly) { + return runYtdlpMode(opts); + } + // THE MEDIA TIER'S HOOK FOR A BATCH (release 17): after the child returns — + // done, failed or cancelled, whatever it finished — every video dir this run + // touched has its media moved into channels/<slug>/media when the channel + // has one. The per-video managed downloads tier as they go + // (downloadOneManaged); this catches what a raw yt-dlp run finalised itself. + // Never throws, and costs one lstat on a classic channel. + const startedAt = Date.now(); + try { + await runYtdlpMode(opts); + } finally { + await tierChannelMedia(opts.paths, opts.channelSlug, { + since: startedAt - TIER_SINCE_SLACK_MS, + onLog: opts.onLog, + }); + } +} + +async function runYtdlpMode(opts: RunYtdlpOpts): Promise<void> { switch (opts.mode) { case "store-playlist": await storePlaylist(opts); diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -29,6 +29,7 @@ - **The dashboard and `/jobs` keep answering while a channel's report is regenerated.** Regenerating a report walks every video of the channel inside the editor, and two regenerations of channels with a few thousand videos, running side by side, kept `/`, `/channels` and `/jobs` from loading for over an hour. Regenerations now run one at a time, on their own `refresh-report` queue on `/jobs` — a channel's own **Refresh report** included, which now waits its turn there too: it waits up to 15 seconds, then says where its job is instead ("Queued behind 3 report regenerations — the report updates when it finishes (job …).") and the page catches up when it runs, and a regeneration that fails now shows its reason under the button; a channel whose report is already waiting is not queued a second time, whether the request came from a finished job, **Refresh report** or **Update all reports**, and a change made while a channel's report is being regenerated queues one more regeneration after it rather than being missed; and the walk pauses between batches of videos so pages are served in between. **Update all reports** answers as soon as the regenerations are queued, and the reports land one after another; `pnpm ops refresh-report` answers with the job ids for both a single channel and `{"all":true}`, which `--wait` follows. Needs a rebuild and restart of the editor. - **The operations pages share one count of the lanes' pending work.** Every open operations page asks for the lanes' status every 3 seconds, and each request used to count every lane's pending videos afresh from every channel's report. That count is now made once and handed to every request in the next 3 seconds. Changing a lane's rules, a focus or a channel's priority counts again at once; otherwise a pending count can be up to 3 seconds behind a report that was just rewritten or a video a lane just picked. A lane's hold, its runner and its picks are still read fresh on every request. - **Jobs a stopped editor left "running" are closed when it starts again.** A job that was still running when the editor's process ended (killed, crashed, or shut down before the job had finished unwinding) kept "running" in its record for good, and `/jobs` listed it as archived. On start the editor now marks each one **cancelled**, with "interrupted: the process running it stopped before it finished" as the reason on the job's page, and its end time is the last time its log was written. Nothing is run again; **Retry** works as for any cancelled job. A job that another live process is running, such as `archilyzer run`, is left alone, and the same check now keeps the start-up pass from closing that process's queued jobs. Such leftover jobs never blocked a media move. +- **A channel's text stays on the fast disk when its media moves, so a slow or unplugged media drive no longer holds its transcripts.** A channel's big files — the audio and the raw live-chat replay — can now live in the channel's own `media` folder, on this disk or another, while its transcripts, cues, metadata and every other small file stay in `data/` where they always were; each big file that moves leaves a small link behind, so everything that opens it by name still finds it. New downloads, transcodes and live-chat normalizes put their big files there as they finish, and every cleanup that deletes audio removes the file the link points to, not just the link. What that changes when a media drive is stalled, unplugged or mid-move: **the index and stats builds never wait on it or are held by it** (a live chat whose transcript cues are out of date keeps the cues the last build read until the drive answers), the channel's report still refreshes (its media size reads as unknown until the drive answers), **digests keep running — even during a move of that channel's media** — and so do normalize, the availability checks, the metadata scan and clip eviction. Transcription, downloads, the backfill lane and anything else that opens the audio are held as before. **A channel moved the old way — its whole `data/` on the other drive — is now shown as "Media layout retired" and held by everything, the builds and digests included, until `archilyzer storage migrate-tier <channel>` brings its text home;** every refusal says so. Deleting a video from its page is refused while its channel's media drive is not reachable, so its audio is never left behind on the drive. Needs a rebuild and restart of the editor. ## [0.11.0] - 2026-09-30 - **Transcripts that arrived after a video was first seen are counted.** The stats behind the homepage, the hub and every site's charts were cached per video and refreshed only when the video's metadata changed, so a transcript that came later — a Whisper run days after the download, or a video downloaded after the last index build — never reached them, and a video with YouTube captions alone had no transcription date. Counts and charts were low; the homepage could show a site with 0 transcripts, 0 channels and 0 hours while it served its videos. A stat is now also redone whenever the index re-reads the video, every transcript has a date, and a captioned video is dated by when its captions arrived rather than by a later Normalize run, so its place on "Transcribed over time" can move. **After updating, rebuild and restart the editor before anything else:** until then, **Build stats dataset** runs the old code and would undo the new stats, while a site, hub or homepage build already runs the new code — and the first stats build of any kind re-reads every video once (about 10–30 minutes on a large archive; it can be stopped and picks up where it stopped). Then build the index, the stats, the homepage, the hub, and the sites. diff --git a/editor/app/channels/[slug]/lib/fixIncompleteTranscript.ts b/editor/app/channels/[slug]/lib/fixIncompleteTranscript.ts @@ -9,8 +9,9 @@ // audio on disk is itself short, so re-running whisper on it just reproduces the // short transcript. The fix MUST re-fetch the audio first. +import { removeMediaFile } from "yt-dlp-transcript-common/lib/mediaTier-server"; import path from "node:path"; -import { readdir, rm } from "node:fs/promises"; +import { readdir } from "node:fs/promises"; import type { ChannelConfig } from "yt-dlp-transcript-common/lib/channelConfig"; import { getPaths, type Paths } from "yt-dlp-transcript-common/lib/paths"; import { getSettings } from "yt-dlp-transcript-common/lib/settings"; @@ -94,7 +95,7 @@ export async function fixIncompleteTranscriptOne(opts: { // than seeing it as already present. const entries = await readdir(videoDir).catch(() => [] as string[]); for (const name of entries.filter(isRealAudioFile)) { - await rm(path.join(videoDir, name), { force: true }); + await removeMediaFile(videoDir, name); onLog(`Removed truncated audio ${name}.`); } onLog(`Re-downloading audio for ${videoId}…`); @@ -150,7 +151,8 @@ export async function clearIncompleteTranscriptOne(opts: { name === "transcript.cues.json", ); for (const name of toRemove) { - await rm(path.join(resolved, name), { force: true }); + // Through its link when tiered (release 17): the bytes go too. + await removeMediaFile(resolved, name); } return { removed: toRemove.length }; } diff --git a/editor/app/channels/[slug]/shardActions.ts b/editor/app/channels/[slug]/shardActions.ts @@ -9,7 +9,7 @@ import { type ShardOp, } from "yt-dlp-transcript-common/controller/shard"; import { readChannelConfig } from "yt-dlp-transcript-common/controller/channels"; -import { assertChannelMediaReachable } from "yt-dlp-transcript-common/lib/channelMedia"; +import { assertChannelTextReadable } from "yt-dlp-transcript-common/lib/channelMedia"; import { runYtdlp } from "yt-dlp-transcript-common/ytdlp/runYtdlp"; import { runWhisperBatch } from "yt-dlp-transcript-common/controller/whisperBatch"; import { runAvailabilityCheck } from "yt-dlp-transcript-common/controller/checkAvailability"; @@ -84,8 +84,12 @@ export async function saveShardConfigAction( // // It is at the top rather than in the download branch alone because all three // read the same dirs for the same reason. + // + // THE TEXT GUARD (release 17): the slices are computed over the video dirs + // in `data/` — the text tier, on the corpus disk — so only an unreadable + // text tier (a `legacy` channel) can make them wrong. try { - await assertChannelMediaReachable(paths, slug); + await assertChannelTextReadable(paths, slug); } catch (e) { return { ok: false, error: (e as Error).message }; } diff --git a/editor/app/channels/[slug]/videos/[id]/videoActions.ts b/editor/app/channels/[slug]/videos/[id]/videoActions.ts @@ -36,10 +36,18 @@ import { type KeepVideosField, type KeepVideosResult, } from "yt-dlp-transcript-common/controller/keepVideosMatching"; -import { ChannelMediaUnreachableError } from "yt-dlp-transcript-common/lib/channelMedia"; +import { + ChannelMediaUnreachableError, + inspectChannelMedia, +} from "yt-dlp-transcript-common/lib/channelMedia"; +import { onDrive } from "yt-dlp-transcript-common/lib/storageHealth"; import { setExcludedFromTruncatedCheck } from "yt-dlp-transcript-common/lib/excludeTruncatedCheck-server"; import { pruneFailedTranscriptions } from "yt-dlp-transcript-common/controller/failedTranscriptions"; import { transcodeAudio } from "yt-dlp-transcript-common/controller/transcode"; +import { + removeMediaFile, + removeVideoDirMedia, +} from "yt-dlp-transcript-common/lib/mediaTier-server"; import { transcribeWithWorker } from "yt-dlp-transcript-common/controller/transcribeOne"; import { findVideoSourceUrl } from "yt-dlp-transcript-common/controller/undownloadedVideos"; import { @@ -499,7 +507,9 @@ export async function deleteVideoFileAction( if (!s.isFile()) { return { ok: false, error: `Not a regular file: ${filename}` }; } - await rm(target, { force: true }); + // Through its link when the file is tiered (release 17): the bytes on the + // media tier go too, never orphaned behind a removed link. + await removeMediaFile(videoDir, path.basename(target)); revalidatePath(`/channels/${slug}/videos/${videoId}`); requestChannelSnapshot(getPaths(), slug); return { ok: true }; @@ -565,6 +575,32 @@ export async function deleteOneVideoDir( error: "Refusing to delete: video path resolved outside the data dir", }; } + // The video's media tier first (release 17): every tiered link's bytes and + // its `media/<id>/`, which the recursive rm of the text dir cannot reach. + // Only while the channel's media is reachable, and through the watchdog: on + // an unmounted drive the bytes would be orphaned, on a stalled one the rm + // would block. Refused before anything is touched. + const config = await readChannelConfig(getPaths(), slug); + const media = await inspectChannelMedia(getPaths(), slug, config, { fresh: true }); + if (media.status !== "ok" && media.status !== "in-place") { + return { + ok: false, + error: + `Refusing to delete ${videoId}: its media is not reachable — ` + + `${media.detail ?? media.status}. Nothing has been touched.`, + }; + } + const drive = config?.mediaDir?.trim(); + try { + await (drive + ? onDrive(drive, () => removeVideoDirMedia(resolved)) + : removeVideoDirMedia(resolved)); + } catch (err) { + return { + ok: false, + error: `Refusing to delete ${videoId}: ${(err as Error).message}. The text was not touched.`, + }; + } await rm(resolved, { recursive: true, force: true }); return { ok: true }; } @@ -595,7 +631,7 @@ export async function removeAudioFilesForVideo( wrongFormatOnly: opts.wrongFormatOnly, }); for (const name of toRemove) { - await rm(path.join(resolved, name), { force: true }); + await removeMediaFile(resolved, name); // derefs a tiered link (release 17) } return { ok: true, removed: toRemove.length }; } diff --git a/editor/app/components/MediaLocationBadge.tsx b/editor/app/components/MediaLocationBadge.tsx @@ -54,6 +54,9 @@ const LABELS: Record<ChannelMediaStatus, string | null> = { "in-transition": "Media moving", inconsistent: "Media inconsistent", stalled: "Media not answering", + // Release 17: the retired whole-directory layout (`archilyzer storage + // migrate-tier`). Added by slice T1 so the status union stays exhaustive. + legacy: "Media layout retired", }; // The one-word state, for the compact rendering of a NAMED location: "on @@ -66,6 +69,7 @@ const SHORT_STATUS: Record<ChannelMediaStatus, string | null> = { "in-transition": "moving", inconsistent: "inconsistent", stalled: "not answering", + legacy: "layout retired", }; // Null means "draw nothing" — an in-place channel, or no location at all (a diff --git a/plans/FACTS.md b/plans/FACTS.md @@ -3871,6 +3871,52 @@ actions reject a channel with running or queued jobs before they enqueue. (`channelMedia.ts:40`) is the in-flight marker: present means "in transition" to every guard, its `phase` is what lets an interrupted move resume, `deleteChannel` and `renameChannel` refuse while it exists, and `clearRelocationMarker` (`:160`) removes it and nothing else. +### Release 17 slice T1 — the media tier's model, and which guard holds what (2026-10-01) + +Supersedes, where they differ, the facts in the section above and in "The index build's hold", +"`cues.mjs` carries a reachability twin" and the drive-health notes below; those passages are kept +as they were measured. +- **The model.** A big file — `isTierable(name)`: `audio.<ext>` and `transcript.live_chat.json`, never + `source-media.*` or a partial (`common/lib/mediaTier.ts`, pure; `classifyEntry` is media / text / + scratch by NAME) — may be a RELATIVE link `data/<id>/<name> -> ../../media/<id>/<name>`; + `channels/<slug>/media` is a real dir (tiered in place), ONE absolute link to `<root>/<slug>/media` + recorded as `config.mediaDir` (relocated), or absent (classic). `data/` is always a real dir on the + corpus disk. The hook (`tierMediaFile` & co., `lib/mediaTier-server.ts`) is called after every + media finalisation (transcode, both download-outcome writes, a batch `runYtdlp` mode, a live-chat + normalize that wrote — not the export build's archive pass) and never throws; the bytes are placed by + a hard link (same disk) or an atomic copy before one rename swaps the name, so a crash never leaves the + name missing; the link carries the file's mtime (`lutimes`), and the + live-chat freshness check and the index's sub-track key (`subsMs`) `lstat` a tierable name, so the + index's scan never reaches the media drive; when a tiered live chat's cues are stale the index reads + the raw only while the channel's media is `ok`/`in-place` — one `stat` probe through `onDrive(mediaDir)` + (the watchdog's budget is sized for a stat), then a direct read — and otherwise keeps the cues the last + build held, or, with none, stores `subsMs: null` so the next build retries (`buildIndex.test.ts` (k): a + stalled drive, nothing asked of it; (l): the retry). + An EXDEV tier gives the copy the file's times, so `stat` and `lstat` agree. Every deleter of + a video-dir entry goes through `removeMediaFile` / `removeVideoDirMedia`. +- **`inspectChannelMedia`** now describes the `media` link + `mediaDir`, returns `mediaLink` and + `text: { dir, readable }`, and has a seventh status, **`legacy`**: a `data` link or a recorded + `dataDir` (the retired whole-directory layout), answered without a call to the far drive, text + unreadable. Memo key: slug + `mediaDir` + `dataDir`. +- **Two guards.** `assertChannelMediaReachable` (only `ok`/`in-place` pass) for a kind with + `needsMedia` ("opens or writes the BIG file"); `assertChannelTextReadable` (passes unless `legacy`, + `data/` not a dir, or a marker with `scope: "tier-migration"`) for a kind with the new `needsText` + (14 kinds flipped, pinned in `jobs/jobKinds.test.ts`). `runManagedFunction` asks the one the kind + declares, at enqueue and at start. +- **Who is held by what.** The index and stats builds, the snapshot, `normalizeAllTranscripts`, + `saveShardConfigAction`, clip eviction and the digest batch read the TEXT tier and hold only on an + unreadable text tier; the digest lane holds on `isTextHeld` (= `legacy`) or unreadable text + (`isChannelHeldForLane`, `controller/autoRunner.ts`); transcription, download and backfill lanes, + the backfill batch and `normalizeAllLiveChat` keep the media hold. A media move's marker + (`scope` absent or `"media"`) holds media writers only; `channelWriters(slug, { mediaOnly })` + answers that set. +- **The snapshot's bytes** are three siblings: `totalMediaBytes` (tierable names, through one + `onDrive(mediaDir)` per video for the links; ABSENT, with `totalAudioBytes`, when the media tier + could not be read), `totalTextBytes` (new), `totalClipsBytes` (no longer inside media). +- **`onDrive` by file kind.** `channelMediaStall` keys on `mediaDir` (else the retired `dataDir`); + `channelTextStall` on the retired `dataDir` only. `readChannelStat` and the recency tail reads pass + a drive only for a legacy channel. `rootOfUnknownPath` strips `<root>/<slug>/media` and `/data`. + ## Channel priority (verified 2026-09-11) — one tier per channel, four compiled trees Branch `channel-priority/s5`, off S0's `28bfee3`, merging `s1`–`s4` and closing the twelve @@ -4201,6 +4247,10 @@ corrected twice. ### `cues.mjs` carries a reachability twin of `assertChannelMediaReachable` +**Release 17 slice T1:** it is now the twin of `assertChannelTextReadable` — it refuses a `legacy` +channel (a `data` link or a recorded `dataDir`, mounted or not) and a marker only when its `scope` is +`tier-migration`; relocated MEDIA is read past. The paragraphs below describe the old twin. + `umtool/report-to-video/cues.mjs`'s `checkChannelReachable` replicates `inspectChannelMedia` / `assertChannelMediaReachable` (`common/lib/channelMedia.ts:320-327`, which carries the cross-reference comment) in plain `.mjs`, because umtool's bins run under diff --git a/plans/release-17.md b/plans/release-17.md @@ -1033,4 +1033,233 @@ keeps every `[Unreleased]` bullet (main's, then D0's). Re-gated on the merged tr 2,519/2,519**; **editor unit 109/109**; e2e `dashboard-answers`, `channels`, `jobs`, `channel-storage` (run 7): **25 passed, 0 failed, 6.4 min** (23 with the queue wait; every dashboard answer ≤ 3.7 s). +### Slice T1, as shipped — the classifier, the media link and the two guards (2026-10-01) + +Branch `r17/media-tier-model` off `main` `7f4901f1`, worktree `~/Projects/plans-export-header-first-search` +(editor 4201, test 4211, export 4210 — `pnpm wt list`'s block #12), one Opus implementer, beside D0 and U1. +Scratch files `T1-*` in the job's `tmp`. The plan is "The model (A′)" §§1–4 above plus the `channelWriters` +option and the umtool twin; the mover and every editor surface are T2's. + +**What it does.** +- **The classifier** (`common/lib/mediaTier.ts`, pure): `classifyEntry(name)` is `media` / `text` / `scratch` + by name over `mediaFiles.ts`'s anchored predicates; `isTierable` is narrower (`audio.<ext>` and + `transcript.live_chat.json` — never `source-media.*`, never a partial); `classifyVideoDir`. The table test + pins 32 names, the four the slice table names among them. +- **The hook and the deleter** (`common/lib/mediaTier-server.ts`): `tierMediaFile` moves a finished big file + into `channels/<slug>/media/<id>/` and leaves the relative link `../../media/<id>/<name>` — + `tiered | left | already`; never throws; a classic channel, a dangling or stalled `media` link, a full disk + leave the file real; the per-video `mkdir` is non-recursive; EXDEV copies atomically; the link replaces the + name by a rename (no window with nothing there); a failed same-disk move is undone. `tierVideoDir`, + `tierChannelMedia({ since, createMediaDir })`, `removeMediaFile` (derefs only inside the channel's own + `media/<id>/`), `removeVideoDirMedia`, `channelMediaLink`, `relocatedMediaDir`, `tierLinkTarget`. +- **Every media finalisation calls it:** `transcodeAudio` after its rename (so the app extraction, the + audio-checked download and the video page's Transcode), `downloadOneManaged` before both download-outcome + writes, `runYtdlp`'s five media-writing modes after the child returns (`since` the run's start, in a + `finally`), `normalizeLiveChat` after it writes the cues. **Every deleter derefs:** the three cleanup + sweeps, `purgeSupersededAutoSubs`, `backfillReacquire`, the app extraction's source discard, the audio-checked + source discard, `fixIncompleteTranscript` (both), and the video page's file delete, bulk audio removal and + directory delete (`removeVideoDirMedia` first). Grep gate + `git grep -n "remove(path.join(.*videoDir\|rm(path.join(videoDir" common editor`: **0 hits** (the module's own + `rm`s take a joined path). +- **`config.json`:** `mediaDir` (CHANNEL.md regenerated by `bin/file-schemas-docs.ts`); `dataDir` documented + RETIRED and still parsed. +- **`inspectChannelMedia`** describes `channels/<slug>/media` + `mediaDir`, returns `mediaLink` and + `text: { dir, readable }`, and has a seventh status, **`legacy`** (a `data` link or a recorded `dataDir`, + answered from the corpus disk alone, text unreadable, its detail naming `archilyzer storage migrate-tier + <slug>`). The marker parses an optional `scope` (`media` | `tier-migration`). Memo key: slug + `mediaDir` + + `dataDir`. `assertChannelTextReadable` beside the media guard; `isTextHeld`, `HELD_REASON.legacy`; + `channelMediaStall` keys on `mediaDir` (else the retired `dataDir`), `channelTextStall` on the retired + `dataDir` only; `rootOfUnknownPath` strips `<root>/<slug>/media`. +- **Which guard holds what.** The index and stats builds, the snapshot, `normalizeAllTranscripts`, + `saveShardConfigAction`, clip eviction and the digest batch ask the text tier (`text.readable` / + `assertChannelTextReadable`); the digest lane is held by `isTextHeld` or unreadable text, the other three + by `isMediaHeld` (`isChannelHeldForLane`), and the pick→run marker backstop lets the digest lane through a + media move; the backfill batch and `normalizeAllLiveChat` keep the media guard. Fourteen kinds flip + `needsMedia` → `needsText` (the plan's list), and `runManagedFunction` asks the text guard for them, at + enqueue and at start. `channelWriters(slug, { mediaOnly })` keeps `kindNeedsMedia` jobs and the non-digest + lanes' units (no caller yet; T2's mover passes it). +- **`onDrive` by file kind.** The snapshot reads the text directly; a video's tiered links are statted + together as one `onDrive(mediaDir)` call, only while the media is `ok`/`in-place`; a drive that does not + answer leaves `totalMediaBytes` and `totalAudioBytes` ABSENT and the snapshot is still written. Byte + fields: `totalMediaBytes` (tierable names), `totalTextBytes` (new), `totalClipsBytes` (no longer inside + media). `readChannelStat` and the recency tail reads pass a drive only for a legacy channel. +- **umtool's twin** (`report-to-video/cues.mjs`) is now the twin of the TEXT guard: it refuses a legacy + channel, mounted or not, and a marker only when its `scope` is `tier-migration`; relocated media is read past. + +**Deviations from the plan** (one sentence each): +1. `JobKindMeta.needsText` (+ `kindNeedsText`) is new: a kind flipped off `needsMedia` would otherwise run + unguarded on a legacy channel whose `data/` link dangles and read it as empty. +2. The tier link carries the file's times (`lutimes`) and `normalizeLiveChat`/`isLiveChatCuesFresh` `lstat` + the raw replay: the raw is media now, and the index's freshness check would otherwise reach the media + drive per video. +3. `normalizeLiveChat` tiers only when it wrote the cues (not on "fresh"). As first shipped this was + recorded as "an export build's normalize pass moves nothing", which was not true (`archiveLiveChat` + called it, so a build with stale cues tiered — and read — the raw); since the review `archiveLiveChat` + passes `tier: false`, so a build moves nothing. It still READS a stale raw replay through the link with + no channel guard (T2: treat the corpus-wide live-chat passes as media readers). +4. `normalizeAllLiveChat` (corpus-wide, no slug for `runManagedFunction`) skips a channel whose media is not + reachable — the plan's "`normalize-live-chat` refused". +5. `evictClipWindows` asks the text tier (`clips/` is never tiered) — its kind is in the flip list. +6. The classifier's scratch also takes `*.part`, `*.ytdl` and the hook's own `.<name>.tierlink-<pid>`. +7. `relocatedDataDir` stays exported (the retired shape) for the mover, the re-point, the rename and + `storageActions.ts` until T2 replaces it with `relocatedMediaDir`. +8. `MediaLocationBadge.tsx` (T2's) gained the two `legacy` table entries the plan names (`Media layout + retired`), because its two `Record<ChannelMediaStatus, …>` tables fail tsc without them. +9. The `isChannelHeldForLane` helper is new, so the lane decision is tested over real inspect answers. +10. `markerHoldsText`: a scope-less marker whose target is the retired `<root>/<slug>/data` shape (the old + mover, until T2 writes `scope`) holds the text too, in the TS guard, the digest-lane backstop and the + cues twin — the plan's "refuse only `tier-migration`" would let a digest write into a `data/` the old + mover is copying. +11. The hook writes nothing while a `.relocating.json` stands on the channel (the file stays real), a + backstop for a writer that started before a move. + +**Skipped until T2 rebases them** (28, each `{ skip: T1_SKIP }` with the reason string; they build the retired +layout and expect it to read `ok`): `relocateChannelMedia.test.ts` 13 — "out: copies, links, records the +target, keeps mtimes and reclaims the source", "abort from an onLog hook leaves the source intact, and the +rerun completes", "back: restores a real directory, clears the config and reclaims the target", "out @ swap: +crash before the rename — …", "out @ swap: crash after the config write — …", "out @ reclaim: the rerun +sweeps every parked copy …", "back: an inconsistent channel is refused, …", "back: an unreachable channel is +refused", "out @ swap: a directory timestamp is settled …", "out @ swap: a file the target is missing is +mirrored …", "a resume with a stale extra dir on the destination completes", "reconcile: an extra and a +changed file …", "reconcile: a marker past the copy phase resumes, …"; `renameChannel.test.ts` 1 — "rename +re-points a convention-shaped relocated media dir"; `storageLocations.test.ts` 7 — "channelsOnLocation +buckets ok / unreachable / moving …", "re-point rewrites both channels' links and configs, …", "re-point +refuses a target that has no media …", "a failure on the second channel rolls the first one back", "re-point +refuses a busy channel and names it", "a rerun after a crash finishes the channels that were left", "a +channel killed between its symlink and its config write is resumed, …"; `storageWatch.test.ts` 7 (reason +"release 17 T2 rebases the storage watch on mediaDir"; `storageWatch.ts` still reads `config.dataDir`) — "a +channel whose target is gone is auto-paused, …", "the drive coming back restores the tier it overwrote", +"write: false reports the transition …", "a drive that blips for one pass is never paused", "the restore +needs only one good pass", "one missed probe stalls the location; …", "the stall clears only after two clean +probes in a row". + +**Re-premised (this slice's own guards, not skipped):** `channelMedia.test.ts` (the media link; 7 legacy / +text cases), `buildIndex.test.ts` (the hold now comes from an unreadable text tier — the channel put on the +retired layout with its drive away — plus case (j): an unmounted MEDIA drive does not hold the index; the +write spy no longer counts a symlink's target as a written path), `buildStats.test.ts` ((i) likewise, (i2) +new), `storageStall.test.ts` (a stalled media drive: counts, recency and the snapshot are read with no call on +it; the legacy cases keep the old gates; M3 is now "a tiered file's stat never answers → the snapshot is +written, its media bytes unknown"), `channelSnapshot.test.ts`, `evictClipWindows.test.ts`, +`doctor.test.ts`, `run-operation.test.ts`, and umtool's `cues.test.mjs` (6 cases). + +**Open questions, answered.** +1. **What `import-one` writes:** `importVideoAction` (`editor/app/channels/[slug]/pipelineActions.ts:465`) + runs `downloadOneManaged` for one URL — media, subtitles, metadata, the archive line, then the roster. A + media writer: it keeps `needsMedia`, and its media are tiered by `downloadOneManaged`'s hook. +4. **Who hardcodes `<root>/<slug>/data`:** only the mover family — `relocatedDataDir` and its callers + (`relocateChannelMedia.ts`, `renameChannel.ts`, `storageLocations.ts`'s re-point, `storageActions.ts`'s + marker check), `deleteChannel`'s `<root>/<slug>` reclaim and `lib/storageLocations.ts`'s and + `savedVideoStore.ts`'s comments — all T2's; and `storageHealth.ts`'s `rootOfUnknownPath`, which now takes + both suffixes. Nothing in umtool, mcp, docker or `scripts/`. So `<root>/<slug>/media` beside + `<root>/<slug>/data` is as planned. + +**`isFile()` over a video dir, the census** (a dirent `isFile()` is false for a link): `channelSnapshot.ts`'s +byte loop — rewritten (`lstat`, links statted through the watchdog); `channelSnapshot.ts`'s `dirFileBytes`, +`evictClipWindows.ts:145` and `clipWindow-server.ts:49` — over `clips/`, never tiered, fine; +`downloadOneManaged.ts:430` (`discardPrefetchDir`) — a tiered link keeps the directory, the safe direction; +`relocateDir.ts:70` (`measureTree`) — keeps `isFile()` by the plan. **For T2:** `videos/[id]/page.tsx:67` +(`loadVideoDir`) and `videos/page.tsx:61` hide a tiered link today — on a tiered channel the video page and the +videos list would not show the audio until T2 lands. + +**Found and left.** +- For T2: `views/storage.ts`'s "N of it is fetched clip windows" and `storageLocations.ts`'s + `mediaBytes`/`clipsBytes` rollups still treat clips as part of `totalMediaBytes`; `channelRow.ts` and + `freeUpSelection.ts` read `totalMediaBytes`, now the media tier alone. The old mover still records + `dataDir`, so a Move media on this branch alone produces a `legacy` channel. +- For T3: the migration's links should carry each file's times (`lutimes`, as the hook does), or every + migrated live chat's cues read as stale and the next index build re-parses the raw from the media drive; + `buildIndex` re-parses a stale raw replay directly (no watchdog) and skips the track when it cannot. +- `generateChannelSnapshot` passes a null config (a `config.json` with no valid `handling`) to the guard, + which then reads no relocation — as before this slice. +- `removeVideoDirMedia` cannot clear `media/<id>/` while the media drive is unmounted (its `lstat` fails); + the video page's delete then leaves those bytes behind. + +**Commits** + +| Commit | What | +|---|---| +| `1d228a79` | `common:` the classifier (`mediaTier.ts`) and the hook/deleter (`mediaTier-server.ts`) + tests | +| `8335f151` | `common:` the model — `mediaDir` (CHANNEL.md regenerated), `legacy`, the text guard, `isTextHeld`, the builds/snapshot/normalize/shards/eviction/digest batch on the text tier, the snapshot's three byte fields, `needsText` and the fourteen flips, `mediaOnly`; tests re-premised; the 28 T2 skips | +| `85180cd6` | `common, editor:` every media finalisation tiers (transcode, both outcome writes, the batch modes, live-chat normalize), every deleter derefs, the link carries the file's mtime | +| `05275d86` | `common, umtool:` `isChannelHeldForLane` + its test, `normalizeAllLiveChat`'s media guard, the flips pinned, the text-guard job test, the call-site hook tests, the cues twin as a text guard | +| `d56050ff` | `common, umtool:` `markerHoldsText` (the old mover's scope-less `…/data` marker holds the text), the hook writes nothing under a marker | +| this commit | `plans:` this section, FACTS ("Release 17 slice T1"), the editor changelog | + +#### Gates (logs `$T/T1-*.log`) + +- **tsc** (all workspaces) clean at every commit; last at `d56050ff`. +- **common:** at `05275d86` **2,529 passed, 0 failed, 28 skipped** (2,557; `main`'s 2,528-test run had 49 + failures on this branch's first pass, all re-premised or skipped as above). At `d56050ff` 2,530 passed, + 1 failed, 28 skipped: the failure is `storageHealth.test.ts`'s "M4: a healthy 64-wide walk … at half the + budget" timing case under a machine load of 22–28 (other sessions); the file passes 36/36 twice in + isolation right after. New tests: `mediaTier.test.ts` 35, `mediaTier-server.test.ts` 16, + `mediaTierHooks.test.ts` 4, `channelMedia.test.ts` 23 (rewritten), plus cases in `channelWriters`, + `autoRunner`, `streamCommand`, `jobKinds`, `buildIndex` (j), `buildStats` (i2), `storageStall`, + `channelSnapshot`, `evictClipWindows`, `storageHealth`. +- **Editor unit:** 109/109. **test:scripts:** 304 passed, 2 skipped (306) — a first run had the two + `queue-lock.test.mjs` timing cases fail under load; the rerun is clean. **umtool `cues.test.mjs`:** 23 + passed, 1 skipped (LIVE). +- **Build:** the capped editor build (`systemd-run --scope -p MemoryMax=6G`, `next build`): exit 0, 172 s, + at `05275d86`. +- **e2e** (from the worktree root, `$T/T1-specs.txt`: maybe-missing, video-page, cleanup-holds, + cleanup-actionable, auto-queue, digest, jobs-channel, channel-storage) at `05275d86`: **78 passed, 5 + failed, 8.9 min** (after 33.6 min in the queue). The 5 are all `channel-storage.spec.ts`, all the retired + layout reading `legacy`, as expected until T2 rebases the mover: "relocate a channel's media to another + root, and move it back" (:80), "the /channels bulk move queues one job per channel and skips the rest" + (:250), "a bulk move puts every job on one queue and skips a channel with nothing to move" (:391), "the + Storage panel moves to a location picked by name" (:484), "Sync all skips a channel whose media drive is + not mounted" (:1091 — now says `media legacy: … run archilyzer storage migrate-tier test-youtube`). Not + re-run after `d56050ff` (unit-covered; a marker rule the specs do not reach). +- **Privacy gate:** 0 added lines carry the user or host name (`git diff 7f4901f1`, counts only; the one + file the whole-file grep names is `plans/FACTS.md`, with the same count as on `main`). +- **Numbers tool:** none. + +#### Review (SHIP AFTER FIXES, 12 findings) and the fixes + +Every commit over `main..HEAD` was rewritten (`git filter-branch --msg-filter`, worktree only) so its +trailer is the session's `Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>` + `Claude-Session` +line (finding 12): tip `f3f23792` → `001e9de7`; the shas in the table above are the rewritten ones +(`828dc1a2` → `98ef1772` is the first `plans:` commit). + +| Finding | Fix | +|---|---| +| H1 the index stats (and may read) a tiered live chat through its link | `263c8b17`: the scan's `subsMs` loop `lstat`s a tierable track; a stale-cues fallback on a tiered raw reads through `onDrive(mediaDir)` only while the media is `ok`/`in-place`, else keeps the cues the last build held; `6b795c67`: an EXDEV tier gives the copy the file's times (`stat` and `lstat` agree). New gate `buildIndex.test.ts` (k): a STALLED media drive (the location marked, every call through a link onto it hanging) with a tiered live chat — the index and the stats build complete, nothing held, the live-chat cues kept, no call on the drive. A mutation (a following `stat` in `subsMs`) makes it hang. | +| L2 crash window between the two renames | `6b795c67`: the bytes are placed by hard link (same disk) or atomic copy and renamed into the tier, then ONE rename swaps the name — it always resolves; `tierVideoDir` removes a dead process's stray `.tierlink-<pid>`. Tests: one inode after a same-disk tier; a stray healed. | +| L3 `removeMediaFile` derefs only into `media/<sameId>/` | `6b795c67`: any id under the channel's own `media/` (never out of it); the target's dir is dropped when it empties. Test. | +| L4 the export build tiers the raw live chat | `95c9a82c`: `normalizeLiveChat({ tier: false })` from `archiveLiveChat`; deviation 3 corrected above (the build still READS a stale raw — for T2). | +| L5 a move that starts mid-hook (for T2) | `6b795c67`: the marker is asked again after the bytes land and before the name changes; the tier's copy is removed if one stands. | +| L6 a video-dir delete on an unmounted or stalled media tier | `95c9a82c`: `deleteOneVideoDir` refuses unless the channel's media is `ok`/`in-place`, and removes the tier's side through `onDrive(mediaDir)`, refusing with the drive's sentence when it does not answer — before the text is touched. | +| L7 gates at the tip | re-gated below. | +| N8 an in-place `media/` takes a watchdog slot | `6b795c67`: a real `media/` is answered from the corpus disk. | +| N9 `totalTextBytes` counts containers and partials | `95c9a82c`: documented as everything on the corpus disk but `clips/` and the tier. | +| N10 a shard save runs the tier sweep | `95c9a82c`: skipped when `saveShardOnly`. | +| L11 records | `001e9de7`: FACTS and the changelog bullet say what the index does with a tiered live chat; the bullet says "Needs a rebuild and restart of the editor." and names the refused video delete. | + +**Gates at `001e9de7`** (logs `$T/T1-regate2.log`, `$T/T1-e2e-2.log`): tsc clean; common **2,535 passed, +0 failed, 28 skipped** (2,563); editor unit 109/109; test:scripts 304 passed, 2 skipped; umtool +`cues.test.mjs` 23 passed, 1 skipped; capped editor build exit 0, 101 s; e2e (the same 8 specs) **78 passed, +5 failed, 8.2 min** (after 39 min in the queue) — the same five `channel-storage.spec.ts` cases (:80, :250, +:391, :484, :1091), the old mover's layout reading `legacy`, nothing else red. Privacy: 0 added lines carry +the user or host name. + +**Re-review: SHIP**, with R1 and R2 folded in before the merge. R1 (`15acd7c7`): the index read a whole raw +live chat under the watchdog's one-stat budget, and a timeout was never retried. Ruling: it probes the +drive with one `stat` through `onDrive(mediaDir)` and then reads the raw directly, and a track that is +neither read nor kept is stored with `subsMs: null`, so the next build retries it (case (l), +mutation-checked). R2 (`ba24a313`): the classifier calls the hook's media-side temp +(`.<name>.tiering-<pid>`) scratch, and `tierVideoDir` sweeps a dead process's temp from `media/<id>/` while +the tier answers. Gates at `ba24a313`: `buildIndex.test.ts` + the tier tests 72/72; tsc clean; common 2,538 passed, 0 failed, 28 skipped (2,566). + +**Merge of `main` `173dd42b`** (D0 on U1, export 0.11.1, XP) at `ffbfbd63`: three conflicts, all appends — +`channelWriters.test.ts` keeps both new tests (T1's `mediaOnly`, D0's ghost meta), `editor/CHANGELOG.md` +keeps every `[Unreleased]` bullet (main's, then T1's), and this file keeps U1, XP, D0, then T1 under +`## Record`; `channelSnapshot.ts` (D0's `mapInYieldingChunks` around T1's per-video body), `channelWriters.ts` +and `buildIndex.ts` (XP's per-site posts predicate) merged cleanly. Re-gated on the merged tree +(`$T/T1-regate4.log`, `$T/T1-e2e-3.log`): tsc clean; **common 2,573 passed, 0 failed, 28 skipped** +(2,601); **editor unit 109/109**; **test:scripts 394 passed, 2 skipped** (two runs at a machine load of +32–35 failed `queue-lock.test.mjs`'s timing cases, 1 then 2; that file passed 11/11 three times alone and +the full suite was clean at load 10 — this branch does not touch `scripts/`); capped editor build exit 0, +92 s; **e2e (the 8 specs) 78 passed, 5 failed, 8.2 min** — the same five `channel-storage.spec.ts` cases +(:80, :250, :391, :484, :1091), the old mover's layout reading `legacy`, nothing else red. + ## Rollout diff --git a/umtool/report-to-video/cues.mjs b/umtool/report-to-video/cues.mjs @@ -48,7 +48,7 @@ // option that is reproducible on a machine with no corpus. // Whatever answers, the returned record carries `from` so a caller can record it. -import { readFile, writeFile, mkdir, lstat, readlink, stat } from "node:fs/promises"; +import { readFile, writeFile, mkdir, lstat, readlink } from "node:fs/promises"; import path from "node:path"; import os from "node:os"; import { createHash } from "node:crypto"; @@ -239,31 +239,37 @@ export function createCueSource({ return { ...record, from: "http" }; } - // A TWIN OF THE EDITOR'S GUARD. `assertChannelMediaReachable` + // A TWIN OF THE EDITOR'S TEXT GUARD. `assertChannelTextReadable` // (`common/lib/channelMedia.ts`) is the same check in TypeScript, and that // file carries a pointer back here — change one, change the other. It is // copied rather than imported because umtool's bins run under plain node // with no `tsx` and no build step, and that stays true for now. // - // WHY IT EXISTS. A channel's media can be relocated to another drive: - // `channels/<slug>/data/` becomes an absolute SYMLINK to `<root>/<slug>/data` - // and `config.json` records the target in `dataDir`. An unmounted drive then - // reads as a plain ENOENT, which `load` used to swallow as "no local copy" - // and answer from the published archive instead — silently cutting from a - // snapshot whose cues can differ from the corpus by seconds (see the note at - // the top of this file). Unreachable has to be loud. + // WHY IT EXISTS. This resolver reads TEXT — `transcript.cues.json`. Since + // release 17 a channel's text stays on the corpus disk whatever its media is + // doing (only the big files move, into `channels/<slug>/media`), so an + // unmounted MEDIA drive does not concern it. The one layout whose text is on + // another drive is the RETIRED whole-directory one — `data/` an absolute + // symlink to `<root>/<slug>/data`, `config.json` recording `dataDir` — and + // there an unmounted drive reads as a plain ENOENT, which `load` used to + // swallow as "no local copy" and answer from the published archive instead: + // silently cutting from a snapshot whose cues can differ from the corpus by + // seconds (see the note at the top of this file). So a `legacy` channel is + // refused, loudly, with its way out — mounted or not. And a marker is + // refused only when its `scope` is `tier-migration` (the migration rebuilds + // `data/` itself); a media move leaves the text where it is. // // SCOPE, DELIBERATELY NARROW. This fires only for a channel the local corpus // actually holds. No `channels/` dir at all (a clone with no corpus), or a - // channel this corpus does not mirror, leaves `data/` absent with no - // configured target — which the editor calls "in-place" and passes, and which - // here still falls through to HTTP. That is the supported archive-only case, - // not a failure. + // channel this corpus does not mirror, leaves `data/` absent with nothing + // recorded — which the editor calls readable and passes, and which here still + // falls through to HTTP. That is the supported archive-only case, not a + // failure. const checkedChannels = new Map(); function unreachable(channelSlug, dataDir, detail) { return new CueLookupError( - `channel "${channelSlug}": local media is not reachable — ${detail}`, + `channel "${channelSlug}": its local text is not readable — ${detail}`, { channelSlug, videoId: null, tried: [dataDir] }, ); } @@ -275,84 +281,70 @@ export function createCueSource({ throw unreachable(channelSlug, dataDir, detail); }; - let configured; + let retired; try { const parsed = JSON.parse(await readFile(path.join(channelDir, "config.json"), "utf8")); if (typeof parsed.dataDir === "string" && parsed.dataDir.trim()) { - configured = parsed.dataDir.trim(); + retired = parsed.dataDir.trim(); } } catch { - // No config.json, or one that is not JSON: nothing records a relocation. + // No config.json, or one that is not JSON: nothing records a layout. } - // A relocation in flight (or interrupted) means `data/` is half of two - // places at once. The editor refuses such a channel; so does this. + // A tier migration in flight (or interrupted) is rebuilding `data/`: half + // of two places at once. The editor refuses such a channel; so does this. + // A media move leaves the text alone. let marker = null; try { marker = JSON.parse(await readFile(path.join(channelDir, ".relocating.json"), "utf8")); } catch { /* no marker: the normal case */ } - if (marker && typeof marker.target === "string" && marker.target.trim()) { + // `markerHoldsText` in channelMedia.ts: a tier migration, or a scope-less + // marker aimed at the retired `<root>/<slug>/data` shape (the old mover, + // copying the whole `data/`). + const holdsText = + marker && + typeof marker.target === "string" && + marker.target.trim() && + (marker.scope === "tier-migration" || + (marker.scope !== "media" && + path.basename(marker.target.trim().replace(/\/+$/, "")) === "data")); + if (holdsText) { fail( - `a media relocation (${marker.direction ?? "out"}) is in progress or was ` + - `interrupted at phase "${marker.phase ?? "copy"}" — target ${marker.target}`, + `its media layout is being migrated (phase "${marker.phase ?? "copy"}") — ` + + `wait for archilyzer storage migrate-tier to finish`, ); } - let link; + let link = null; try { link = await lstat(dataDir); } catch { - // No `data/` at all. With no configured target this is a channel that has - // downloaded nothing — or is not mirrored here — and is not an error. - if (!configured) return; - fail( - `config.json records dataDir ${configured} but ${dataDir} does not exist ` + - `— the symlink is missing`, - ); + /* no data/ at all: below */ } - if (link.isSymbolicLink()) { - let target = ""; - try { - target = await readlink(dataDir); - } catch { - /* unreadable link: reported as such below */ - } - if (!configured) { - fail( - `${dataDir} is a symlink to ${target || "(unreadable)"} but config.json ` + - `records no dataDir`, - ); - } - if (path.resolve(target) !== path.resolve(configured)) { - fail( - `${dataDir} points at ${target || "(unreadable)"} but config.json ` + - `records ${configured}`, - ); - } - // The link points at a DEEP path (<root>/<slug>/data), so an unmounted - // root gives ENOENT here and an empty mountpoint can never be mistaken - // for the media. - let st = null; - try { - st = await stat(configured); - } catch { - /* reported below */ + // THE RETIRED LAYOUT: a `data` link, or a recorded `dataDir`. Refused + // whether or not its drive is mounted, never followed. + if ((link && link.isSymbolicLink()) || retired) { + let target = retired ?? ""; + if (link && link.isSymbolicLink() && !target) { + try { + target = await readlink(dataDir); + } catch { + /* unreadable link: named as such */ + } } - if (!st) fail(`${configured} does not exist (drive not mounted?)`); - if (!st.isDirectory()) fail(`${configured} exists but is not a directory`); - return; - } - - if (!link.isDirectory()) fail(`${dataDir} is neither a directory nor a symlink`); - if (configured) { fail( - `config.json records dataDir ${configured} but ${dataDir} is a real ` + - `directory — the media was never moved, or was moved back by hand`, + `its media layout is the retired whole-directory one${target ? ` (${target})` : ""} — ` + + `run archilyzer storage migrate-tier ${channelSlug}`, ); } + + // No `data/`: a channel that has downloaded nothing — or is not mirrored + // here — and not an error. + if (!link) return; + if (!link.isDirectory()) fail(`${dataDir} is neither a directory nor a symlink`); } function assertChannelReachable(channelSlug) { diff --git a/umtool/report-to-video/cues.test.mjs b/umtool/report-to-video/cues.test.mjs @@ -148,28 +148,36 @@ test("a local corpus is preferred over the network", async () => { } }); -// --- a channel the corpus holds but cannot reach ---------------------------- +// --- a channel the corpus holds but whose text it cannot read --------------- // -// The bug these cover: `data/` may be a symlink to another drive, with the -// target recorded in the channel's `config.json` as `dataDir`. An unmounted -// drive reads as a plain ENOENT, which used to fall through to the archive — -// so a relocated channel silently cut from a snapshot's cues, which can differ -// from the corpus's by seconds. +// The bug these cover: on the RETIRED layout `data/` is a symlink to another +// drive, with the target recorded in the channel's `config.json` as `dataDir`. +// An unmounted drive reads as a plain ENOENT, which used to fall through to the +// archive — so a relocated channel silently cut from a snapshot's cues, which +// can differ from the corpus's by seconds. Since release 17 such a channel is +// `legacy` and refused mounted or not (its way out is `migrate-tier`); a +// channel whose MEDIA alone is relocated (`media/` a link, `mediaDir`) keeps +// its text on the corpus disk and is read as usual, drive or no drive. const SCRATCH = process.env.CUES_TEST_DIR ?? tmpdir(); // A mirrored channel: <root>/channels/<slug>/ with a config.json, and `data/` -// however the caller wants it. -async function corpusWith(slug, { dataDir, data } = {}) { +// (and `media`) however the caller wants it. +async function corpusWith(slug, { dataDir, data, mediaDir } = {}) { const root = await mkdtemp(path.join(SCRATCH, "cues-reach-")); const channelDir = path.join(root, slug); await mkdir(channelDir, { recursive: true }); await writeFile( path.join(channelDir, "config.json"), - JSON.stringify(dataDir ? { url: "https://x/", dataDir } : { url: "https://x/" }), + JSON.stringify({ + url: "https://x/", + ...(dataDir ? { dataDir } : {}), + ...(mediaDir ? { mediaDir } : {}), + }), ); if (data === "symlink") await symlink(dataDir, path.join(channelDir, "data")); if (data === "dir") await mkdir(path.join(channelDir, "data"), { recursive: true }); + if (mediaDir) await symlink(mediaDir, path.join(channelDir, "media")); return root; } @@ -179,66 +187,71 @@ async function writeCues(dir, videoId, record) { await writeFile(path.join(vdir, "transcript.cues.json"), JSON.stringify(record)); } -test("a relocated channel whose drive is not mounted throws, and never fetches", async () => { +function sourceOver(root, seen) { + return createCueSource({ + channelsDir: root, + siteOrigin: ORIGIN, + cacheDir: null, + fetchImpl: stubFetch(ROUTES, seen), + }); +} + +test("a legacy channel whose drive is not mounted throws, names the way out, and never fetches", async () => { const missing = path.join(SCRATCH, "cues-not-mounted-" + process.pid, "chan", "data"); const root = await corpusWith("chan", { dataDir: missing, data: "symlink" }); const seen = []; try { - const src = createCueSource({ - channelsDir: root, - siteOrigin: ORIGIN, - cacheDir: null, - fetchImpl: stubFetch(ROUTES, seen), - }); - await assert.rejects(() => src.load("chan", "vid1"), (err) => { + await assert.rejects(() => sourceOver(root, seen).load("chan", "vid1"), (err) => { assert.equal(err.name, "CueLookupError"); assert.match(err.message, /chan/); - assert.match(err.message, /drive not mounted/); - assert.ok(err.message.includes(missing), "the error names the unreachable path"); + assert.match(err.message, /retired whole-directory one/); + assert.match(err.message, /archilyzer storage migrate-tier chan/); + assert.ok(err.message.includes(missing), "the error names the retired target"); return true; }); - assert.deepEqual(seen, [], "an unreachable channel must not fall through to the archive"); + assert.deepEqual(seen, [], "a legacy channel must not fall through to the archive"); } finally { await rm(root, { recursive: true, force: true }); } }); -test("a dataDir with no symlink at all is refused too", async () => { - const root = await corpusWith("chan", { dataDir: path.join(SCRATCH, "elsewhere") }); +test("a legacy channel is refused even with its drive mounted — never followed", async () => { + const elsewhere = await mkdtemp(path.join(SCRATCH, "cues-drive-")); + const target = path.join(elsewhere, "chan", "data"); + await mkdir(target, { recursive: true }); + await writeCues(target, "vid1", { ...RECORD, title: "RELOCATED COPY" }); + const root = await corpusWith("chan", { dataDir: target, data: "symlink" }); const seen = []; try { - const src = createCueSource({ - channelsDir: root, - siteOrigin: ORIGIN, - cacheDir: null, - fetchImpl: stubFetch(ROUTES, seen), - }); - await assert.rejects(() => src.load("chan", "vid1"), (err) => { - assert.match(err.message, /symlink is missing/); - return true; - }); + await assert.rejects(() => sourceOver(root, seen).load("chan", "vid1"), /migrate-tier chan/); + assert.deepEqual(seen, []); + } finally { + await rm(root, { recursive: true, force: true }); + await rm(elsewhere, { recursive: true, force: true }); + } +}); + +test("a recorded dataDir with a real data/ is legacy too", async () => { + const root = await corpusWith("chan", { dataDir: path.join(SCRATCH, "elsewhere"), data: "dir" }); + const seen = []; + try { + await assert.rejects(() => sourceOver(root, seen).load("chan", "vid1"), /retired whole-directory/); assert.deepEqual(seen, []); } finally { await rm(root, { recursive: true, force: true }); } }); -test("a relocation in flight is refused rather than half-read", async () => { +test("a tier migration in flight is refused rather than half-read", async () => { const root = await corpusWith("chan", { data: "dir" }); await writeFile( path.join(root, "chan", ".relocating.json"), - JSON.stringify({ target: "/mnt/big/chan/data", direction: "out", phase: "copy" }), + JSON.stringify({ target: "/mnt/big/chan/media", direction: "out", phase: "copy", scope: "tier-migration" }), ); const seen = []; try { - const src = createCueSource({ - channelsDir: root, - siteOrigin: ORIGIN, - cacheDir: null, - fetchImpl: stubFetch(ROUTES, seen), - }); - await assert.rejects(() => src.load("chan", "vid1"), (err) => { - assert.match(err.message, /relocation \(out\) is in progress/); + await assert.rejects(() => sourceOver(root, seen).load("chan", "vid1"), (err) => { + assert.match(err.message, /being migrated \(phase "copy"\)/); return true; }); assert.deepEqual(seen, []); @@ -247,27 +260,36 @@ test("a relocation in flight is refused rather than half-read", async () => { } }); -test("a reachable dataDir is read from, wherever it points", async () => { - const elsewhere = await mkdtemp(path.join(SCRATCH, "cues-drive-")); - const target = path.join(elsewhere, "chan", "data"); - await mkdir(target, { recursive: true }); - await writeCues(target, "vid1", { ...RECORD, title: "RELOCATED COPY" }); - const root = await corpusWith("chan", { dataDir: target, data: "symlink" }); +test("a media move in flight is not refused: the text stays where it is", async () => { + const root = await corpusWith("chan", { data: "dir" }); + await writeCues(path.join(root, "chan", "data"), "vid1", { ...RECORD, title: "LOCAL DURING MOVE" }); + await writeFile( + path.join(root, "chan", ".relocating.json"), + JSON.stringify({ target: "/mnt/big/chan/media", direction: "out", phase: "copy", scope: "media" }), + ); const seen = []; try { - const src = createCueSource({ - channelsDir: root, - siteOrigin: ORIGIN, - cacheDir: null, - fetchImpl: stubFetch(ROUTES, seen), - }); - const got = await src.load("chan", "vid1"); + const got = await sourceOver(root, seen).load("chan", "vid1"); assert.equal(got.from, "local"); - assert.equal(got.title, "RELOCATED COPY"); - assert.deepEqual(seen, [], "a reachable relocation is still a local read"); + assert.equal(got.title, "LOCAL DURING MOVE"); + assert.deepEqual(seen, []); + } finally { + await rm(root, { recursive: true, force: true }); + } +}); + +test("relocated MEDIA on an unmounted drive: the text is read locally as usual", async () => { + const missing = path.join(SCRATCH, "cues-media-not-mounted-" + process.pid, "chan", "media"); + const root = await corpusWith("chan", { data: "dir", mediaDir: missing }); + await writeCues(path.join(root, "chan", "data"), "vid1", { ...RECORD, title: "TEXT ON THE SSD" }); + const seen = []; + try { + const got = await sourceOver(root, seen).load("chan", "vid1"); + assert.equal(got.from, "local"); + assert.equal(got.title, "TEXT ON THE SSD"); + assert.deepEqual(seen, [], "a local read, whatever the media drive is doing"); } finally { await rm(root, { recursive: true, force: true }); - await rm(elsewhere, { recursive: true, force: true }); } });