Archilyzer · Source

archilyzer

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

commit 1058f75cbcc62ea6c3afce5278d2bf6d0eee1387
parent 1ddbd1c64b98f204f5186b72b2326a9e3d3b2323
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Fri, 11 Sep 2026 11:09:21 -0400

common: four guards stand between an unmounted drive and a re-download

Once `data/` can be a symlink, ENOENT stops meaning "this channel has
downloaded nothing" and starts also meaning "the drive is not mounted". Three
places swallow that ENOENT as the former, and the auto-download lane reading it
that way is an instruction to re-fetch a whole channel onto the volume that was
too full to hold it. The post-Phase-1 dispatch makes the fix four sites, not the
six action files an earlier survey listed.

1. `runManagedFunction` — the funnel every job-shaped action goes through, 35 of
   its 53 call sites already passing a channelSlug. Refused BEFORE the job
   record exists: no queued job, no log, no sidecar, just `{ ok: false }` with
   the reason, which every caller already renders.
2. `autoRunner.buildChannelWork` — the four lane runners' only channel-level
   chokepoint, and the one that covers the direct `runOperationUnit` call a
   batch-level guard misses. A SKIP IS A SKIP: the lane keeps running every
   other channel. It is not a lane stop and it is not a hold — "a zero limit is
   a hold, never a stop" is a different mechanism and is untouched here. Logged
   once per state change, not once per tick, and the recovery is logged too.
3. `generateChannelSnapshot` — one guard protecting everything downstream of a
   snapshot (/cleanup, the bands, all four lanes' work lists). The scheduler
   already keeps the last good snapshot.json on a failed refresh, so throwing is
   what preserves the truth rather than overwriting it with zeroes.
4. `runOperationBatch` — the same swallow, deciding both operation lanes'
   candidate lists. All three callers are already behind guard 1; this is for
   the fourth that will not be.

The kind decides at guard 1, via a declarative `needsMedia` on JobKindMeta.
ABSENT MEANS FALSE: the guard is opt-in, so an unlisted kind keeps exactly its
current behavior, a bookkeeping kind (normalize, availability check, store
playlist, clear markers) is never refused for a drive it never reads, and
`refresh-report` — which is not in the table at all — still runs, so guard 3
reports the reason in that job's own log rather than being silently pre-empted.
28 of the 35 kinds declare it.

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

Diffstat:
Mcommon/controller/autoRunner.ts | 53+++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/controller/channelSnapshot.test.ts | 33+++++++++++++++++++++++++++++++++
Mcommon/controller/channelSnapshot.ts | 10++++++++++
Mcommon/controller/operationBatch.ts | 7+++++++
Mcommon/jobs/jobKinds.ts | 50++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/streamCommand.ts | 27+++++++++++++++++++++++++++
6 files changed, 180 insertions(+), 0 deletions(-)

diff --git a/common/controller/autoRunner.ts b/common/controller/autoRunner.ts @@ -78,6 +78,10 @@ import { type DownloadOutcomeStatus } from "../lib/downloadOutcome"; import { downloadQueueKey } from "../lib/queueKeys"; import { isGateHeld } from "../lib/pauseGates"; import { + inspectChannelMedia, + type ChannelMediaStatus, +} from "../lib/channelMedia"; +import { listChannelConfigs, readChannelConfig, readChannelSnapshotShared, @@ -293,6 +297,37 @@ async function listChannelMeta(paths: Paths): Promise<ChannelMeta[]> { })); } +// ONCE PER STATE CHANGE, not once per tick. buildChannelWork runs on every +// scheduling tick of all four lanes and on a three-second status poll; logging a +// skip each time would write the same line thousands of times an hour and bury +// the one that matters. Keyed by slug, so mounting the drive logs the recovery +// too — an operator watching the log sees the channel leave and come back. +const mediaSkipLogged = new Map<string, ChannelMediaStatus>(); + +function noteSkippedForMedia( + slug: string, + status: ChannelMediaStatus, + detail?: string, +): void { + if (mediaSkipLogged.get(slug) === status) return; + mediaSkipLogged.set(slug, status); + console.log( + `[auto] skipping ${slug}: media ${status}${detail ? ` — ${detail}` : ""}`, + ); +} + +function noteMediaReachable(slug: string): void { + if (!mediaSkipLogged.has(slug)) return; + mediaSkipLogged.delete(slug); + console.log(`[auto] ${slug}: media reachable again`); +} + +// For tests and for /api/test/invalidate-cache: the log-once memory is +// process-local state, not a decision, and nothing downstream reads it. +export function resetChannelMediaSkipLog(): void { + mediaSkipLogged.clear(); +} + // Read each channel's snapshot and project the buckets this runner kind cares // about into ChannelWork, plus a videoId -> owning channel map (a platform/all // leaf spans channels, so the pick needs the owner to locate the video dir). @@ -323,9 +358,27 @@ async function buildChannelWork( const snaps = await mapConcurrent(meta, SNAPSHOT_READ_CONCURRENCY, (m) => readChannelSnapshotShared(paths, m.slug), ); + // GUARD 2 OF FOUR (see plans/relocate-channel-media.md). This is the runners' + // ONLY channel-level chokepoint, and it is the one that matters most: the + // download lane reading an unmounted channel's snapshot as "everything + // undownloaded" is an instruction to re-fetch the whole channel onto the disk + // that was too full to hold it. Three syscalls per channel per tick, run + // concurrently alongside the snapshot reads. + const media = await mapConcurrent(meta, SNAPSHOT_READ_CONCURRENCY, (m) => + inspectChannelMedia(paths, m.slug), + ); for (const [i, { slug, platform }] of meta.entries()) { const snap = snaps[i]; if (!snap) continue; + // A SKIP IS A SKIP. The lane keeps running every other channel: this is + // never a lane stop and it is not a hold — "a zero limit is a hold, never a + // stop" (pauseGates.ts) is a different mechanism and is untouched by it. + const location = media[i]; + if (location && location.status !== "ok" && location.status !== "in-place") { + noteSkippedForMedia(slug, location.status, location.detail); + continue; + } + noteMediaReachable(slug); const buckets: Record<string, string[]> = {}; // Project every bucket a leaf could be pointed at — including the opt-in // auto-caption ones. A leaf that names a bucket the runner never projected diff --git a/common/controller/channelSnapshot.test.ts b/common/controller/channelSnapshot.test.ts @@ -1,5 +1,8 @@ import { test } from "node:test"; import assert from "node:assert/strict"; +import { mkdir, mkdtemp, rm, symlink, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import path from "node:path"; import { attributeAudioHold, diarizationWillNeverClear, @@ -7,7 +10,10 @@ import { emptyHeldAudio, foldBackfillEntry, foldBucketLaneEntry, + generateChannelSnapshot, } from "./channelSnapshot"; +import { ChannelMediaUnreachableError } from "../lib/channelMedia"; +import type { Paths } from "../lib/paths"; import { emptyOperationCounts, presentOperationWork, @@ -325,3 +331,30 @@ test("only the reachable diarization states will ever clear a hold", () => { assert.equal(diarizationWillNeverClear(state), true); } }); + +test("generateChannelSnapshot refuses an unreachable 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 + // 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. + const dir = await mkdtemp(path.join(tmpdir(), "ttb-snap-media-")); + try { + const paths = { channelsDir: path.join(dir, "channels") } as Paths; + const channelDir = path.join(paths.channelsDir, "alpha"); + await mkdir(channelDir, { recursive: true }); + const target = path.join(dir, "platter", "alpha", "data"); + await writeFile( + path.join(channelDir, "config.json"), + JSON.stringify({ url: "https://example.com/c", dataDir: target }), + ); + // A link with no target: an unmounted drive, exactly. + await symlink(target, path.join(channelDir, "data")); + + await assert.rejects( + () => generateChannelSnapshot(paths, "alpha"), + ChannelMediaUnreachableError, + ); + } finally { + await rm(dir, { recursive: true, force: true }); + } +}); diff --git a/common/controller/channelSnapshot.ts b/common/controller/channelSnapshot.ts @@ -18,6 +18,7 @@ import { type Availability, } from "../lib/availability"; import { resolveCookiePolicy } from "../lib/cookiePolicy"; +import { assertChannelMediaReachable } from "../lib/channelMedia"; import { getSettings } from "../lib/settings"; import { loadAvailability, @@ -603,6 +604,15 @@ export async function generateChannelSnapshot( const archivePath = path.join(channelDir, "archive"); const playlistPath = path.join(channelDir, "playlist"); + // GUARD 3 OF FOUR (see plans/relocate-channel-media.md). The readdir below + // swallows ENOENT as "this channel has no videos", so a relocated channel + // whose drive is not mounted would generate a snapshot saying every video is + // undownloaded and every transcript is missing — and everything downstream + // (/cleanup, the channel bands, all four lanes' work lists) reads that + // snapshot as the truth. Throw instead: the scheduler keeps the last good + // snapshot.json on a failed refresh, which is exactly the right outcome. + await assertChannelMediaReachable(paths, slug); + // 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). try { diff --git a/common/controller/operationBatch.ts b/common/controller/operationBatch.ts @@ -38,6 +38,7 @@ import { open, readFile, readdir } from "node:fs/promises"; import type { Paths } from "../lib/paths"; import { getSettings, type SiteSettings } from "../lib/settings"; import { isGateHeld } from "../lib/pauseGates"; +import { assertChannelMediaReachable } from "../lib/channelMedia"; import { runPool } from "../jobs/concurrentRunner"; import type { TaskTracker } from "../jobs/taskHooks"; import type { JobProgress } from "../jobs/registry"; @@ -1576,6 +1577,12 @@ export async function runOperationBatch( ? await resolveDigestFor(run, opts.channelSlug) : null; + // GUARD 4 OF FOUR (see plans/relocate-channel-media.md). Same swallow as the + // snapshot's, deciding the candidate list for both operation lanes. All three + // 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); 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/jobs/jobKinds.ts b/common/jobs/jobKinds.ts @@ -37,6 +37,21 @@ 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. + // + // 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 + // so that `refresh-report` — which is not in this table at all — still runs + // and lets the snapshot generator itself report the reason it refused. A + // bookkeeping kind (normalize, availability check, store playlist, clear + // markers) never declares it: it must not be refused for a drive it never + // reads. + needsMedia?: boolean; }; // One entry per kind known to the system. `label` is included only where the @@ -48,6 +63,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: true, replayable: false, queueKeyStrategy: "parallel", + needsMedia: true, }, "auto-download": { kind: "auto-download", @@ -55,6 +71,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: true, replayable: false, queueKeyStrategy: "parallel", + needsMedia: true, }, "auto-download-unit": { kind: "auto-download-unit", @@ -62,6 +79,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: false, replayable: false, queueKeyStrategy: "platform", + needsMedia: true, }, // The two OPERATION lanes' runners. Same shape as the two above — one // long-lived job per lane on queueKey "", drainable, never replayable — with @@ -77,6 +95,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: true, replayable: false, queueKeyStrategy: "parallel", + needsMedia: true, }, "auto-backfill": { kind: "auto-backfill", @@ -84,6 +103,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: true, replayable: false, queueKeyStrategy: "parallel", + needsMedia: true, }, "whisper-all": { kind: "whisper-all", @@ -91,6 +111,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: true, replayable: true, queueKeyStrategy: "custom", + needsMedia: true, }, "whisper-bucket-downloaded-no-transcript": { kind: "whisper-bucket-downloaded-no-transcript", @@ -98,6 +119,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: true, replayable: true, queueKeyStrategy: "custom", + needsMedia: true, }, // Replace-auto-captions lane, transcribe half: whisper over videos whose only // transcript is a YouTube ASR VTT (the downloadedAutoSubsOnly bucket). Same @@ -109,6 +131,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: true, replayable: true, queueKeyStrategy: "custom", + needsMedia: true, }, // Delete the superseded English ASR VTTs kept as backups next to a finished // whisper transcript. Manual only — never auto-queued — and the single @@ -120,6 +143,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: false, replayable: true, queueKeyStrategy: "custom", + needsMedia: 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 @@ -131,6 +155,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: true, replayable: true, queueKeyStrategy: "custom", + needsMedia: 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 @@ -141,6 +166,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: true, replayable: true, queueKeyStrategy: "custom", + needsMedia: 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 @@ -151,6 +177,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: false, replayable: false, queueKeyStrategy: "parallel", + needsMedia: 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 @@ -176,6 +203,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: true, replayable: true, queueKeyStrategy: "custom", + needsMedia: true, }, "download-from-playlist": { kind: "download-from-playlist", @@ -183,6 +211,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: true, replayable: true, queueKeyStrategy: "platform", + needsMedia: true, }, "download-missing": { kind: "download-missing", @@ -190,6 +219,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: true, replayable: true, queueKeyStrategy: "platform", + needsMedia: true, }, "download-missing-subs": { kind: "download-missing-subs", @@ -197,6 +227,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: true, replayable: true, queueKeyStrategy: "platform", + needsMedia: true, }, "import-one": { kind: "import-one", @@ -204,6 +235,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: false, replayable: false, queueKeyStrategy: "custom", + needsMedia: true, }, "redownload-archive": { kind: "redownload-archive", @@ -211,6 +243,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: false, replayable: false, queueKeyStrategy: "custom", + needsMedia: true, }, "retry-bucket": { kind: "retry-bucket", @@ -218,6 +251,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: true, replayable: true, queueKeyStrategy: "platform", + needsMedia: true, }, "clean-audio-transcribed": { kind: "clean-audio-transcribed", @@ -225,6 +259,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: false, replayable: true, queueKeyStrategy: "custom", + needsMedia: true, }, // Speaker-diarization backfill over a channel's retained audio. The capture // lane's catch-all: it picks up everything the post-transcribe hook missed @@ -239,6 +274,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: true, replayable: true, queueKeyStrategy: "custom", + needsMedia: true, }, // One channel through the backfill lane. Drainable (the batch stops taking new // videos and lets the in-flight one finish) and replayable, because it @@ -250,6 +286,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: true, replayable: true, queueKeyStrategy: "custom", + needsMedia: true, }, // The corpus-wide corrupt-media scan and its per-channel twin. Registered // properly, unlike check-availability / refresh-report / detect-duplicates, @@ -265,6 +302,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { // escape hatch refresh-report and detect-duplicates use. queueKeyStrategy: "parallel", defaultTier: "background", + needsMedia: true, }, "scan-media-channel": { kind: "scan-media-channel", @@ -273,6 +311,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { replayable: false, queueKeyStrategy: "parallel", defaultTier: "background", + needsMedia: true, }, "check-kept-deleted": { kind: "check-kept-deleted", @@ -280,6 +319,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: false, replayable: true, queueKeyStrategy: "custom", + needsMedia: true, }, "persist-kept": { kind: "persist-kept", @@ -287,6 +327,7 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: false, replayable: true, queueKeyStrategy: "custom", + needsMedia: true, }, "backup-saved-videos": { kind: "backup-saved-videos", @@ -349,12 +390,14 @@ const JOB_KINDS: Record<string, JobKindMeta> = { drainable: false, replayable: true, queueKeyStrategy: "custom", + needsMedia: true, }, "remove-wrong-format-audio": { kind: "remove-wrong-format-audio", drainable: false, replayable: true, queueKeyStrategy: "custom", + needsMedia: true, }, }; @@ -371,3 +414,10 @@ export function jobKindLabel(kind: string): string { 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. +export function kindNeedsMedia(kind: string): boolean { + return JOB_KINDS[kind]?.needsMedia ?? false; +} diff --git a/common/jobs/streamCommand.ts b/common/jobs/streamCommand.ts @@ -19,6 +19,8 @@ import { import { writeJobMeta } from "./jobMeta"; import { maybePruneJobLogs } from "./listJobs"; import type { JobSpec } from "./jobSpec"; +import { kindNeedsMedia } from "./jobKinds"; +import { assertChannelMediaReachable } from "../lib/channelMedia"; // Mark a job's channel report dirty so the debounced scheduler regenerates the // snapshot — called both on each completed sub-operation and on the job's @@ -261,9 +263,34 @@ export async function runManagedCommand( return { ok: true, jobId: id, stream, done }; } +// GUARD 1 OF FOUR (see plans/relocate-channel-media.md). runManagedFunction is +// the funnel every job-shaped action goes through, and 35 of its 53 call sites +// already pass a channelSlug — so one check here covers every per-channel media +// action without touching any of them. A refusal happens BEFORE the job record +// is made: no queued job, no log, no sidecar, just `{ ok: false }` carrying the +// reason, which every caller already renders. +// +// The kind decides. `needsMedia` is declarative on JobKindMeta and absent means +// false, so a bookkeeping kind is never refused for a drive it does not read, +// and the relocate job itself — the thing that FIXES an unreachable channel — +// must never declare it. +async function refuseForUnreachableMedia( + opts: CommonOpts, +): Promise<string | null> { + if (!opts.channelSlug || !kindNeedsMedia(opts.kind)) return null; + try { + await assertChannelMediaReachable(opts.paths, opts.channelSlug); + return null; + } catch (err) { + return (err as Error).message; + } +} + export async function runManagedFunction( opts: RunManagedFunctionOpts, ): Promise<StreamActionResult> { + const refusal = await refuseForUnreachableMedia(opts); + if (refusal) return { ok: false, error: refusal }; const registry = getRegistry(); await ensureJobsDir(opts.paths); const { id, logPath, record } = makeJob(