Archilyzer · Source

archilyzer

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

commit 21ccd05e1c72bdc66f8e44cb54c96e3b7c038fd0
parent b4e0215e07c14fe8a647b7b5f1a1f4d9142d7bbf
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Fri, 11 Sep 2026 20:36:16 -0400

relocate: ten ticked rows started ten rsyncs onto one drive

THE QUEUE KEY IS THE ONLY THING THAT DECIDES WHETHER TWO MOVES RUN AT ONCE.
registry.ts submits every non-empty key at concurrency 1 and caps NOTHING across
keys, so `channelQueueKey(slug)` serializes a channel against its own downloads
and transcriptions and against nothing else. Selecting a page of channels on
/channels and pressing Move therefore started one rsync per channel,
simultaneously, all writing to the same destination volume.

That is wrong twice. It is slower — N interleaved sequential writes to one
platter is the worst access pattern the drive has, and the platter is the whole
point. And it breaks the space check: each job runs its own preflight when it
STARTS, so jobs that start together each measure a root none of the others has
written to yet, every one of them is credited room that is already spoken for,
and they jointly overrun it. The failure is recoverable (ENOSPC aborts the copy,
the source is untouched until a verify passes, the partial is resumable) but it
costs hours of copying to learn something one queue slot would have known.

Both comments asserted the opposite. bulkStorageActions.ts said the queue
"already serializes" and that a root filling partway "refuses the remainder
cleanly, one job at a time"; ChannelStorageBulkBar.tsx said the same to the
operator. Neither was true of a per-channel key, and the bar's no-preview-gate
argument rests entirely on it.

relocationQueueKey() takes no slug, deliberately, and lives in
common/lib/queueKeys.ts beside BACKFILL_QUEUE and the two digest keys — the same
question they answer. The one enqueue both actions already share sets it, so
there is no per-caller key to drift. Serializing against the channel's OWN jobs
is not lost: both actions refuse a channel with running or queued jobs before
they enqueue, which refuses rather than waits — the right answer for a move the
operator can see is already in flight.

THE NIT, same file: a channel that has downloaded nothing has no `data/`, and
inspect() calls that `in-place` — correctly, it is not relocated. The bulk path
queued it anyway, and the job's only act was to throw "has no data/ to move"
from the copy phase, after a job record, a log and a queue slot. On a fresh
corpus that is most of a page. It is an up-front skip now, with the reason in
the bar's title attribute like every other skip.

common/lib/queueKeys.test.ts pins the claim as a function of the slug on both
sides, so it still reads as the claim — and fails — if someone gives
relocationQueueKey a slug parameter. The e2e case asserts it where an operator
would see it: two channels moved, two rows in the /jobs table, both on
`relocate` and neither on `channel:`.

common 971/971. channel-storage + disk-space + scheduler --repeat-each 3: 54/54.

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

Diffstat:
Mcommon/jobs/jobKinds.ts | 11++++++++---
Acommon/lib/queueKeys.test.ts | 52++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/lib/queueKeys.ts | 29+++++++++++++++++++++++++++++
Meditor/app/channels/bulkStorageActions.ts | 54++++++++++++++++++++++++++++++++++++++++++------------
Meditor/app/channels/components/ChannelStorageBulkBar.tsx | 6++++--
Meditor/app/channels/lib/relocationJob.ts | 11+++++++++--
Meditor/e2e/channel-storage.spec.ts | 93+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
7 files changed, 237 insertions(+), 19 deletions(-)

diff --git a/common/jobs/jobKinds.ts b/common/jobs/jobKinds.ts @@ -370,9 +370,14 @@ const JOB_KINDS: Record<string, JobKindMeta> = { // honor. Not replayable: a replay carries no direction and no root, and // re-running a move against a channel that has since moved is not a retry. // - // Queue key is channelQueueKey(slug) (set by the action), so a relocation - // serializes with the channel's own bookkeeping jobs instead of running while - // one of them writes into the dir being copied. + // Queue key is relocationQueueKey() (set by the one enqueue both actions + // share), so EVERY relocation in the process serializes against every other + // one: the registry caps a key at concurrency 1 and caps nothing across keys, + // and a bulk move's jobs all write to the same destination volume and all run + // their space check when they START. Serializing against the channel's own + // bookkeeping jobs is not what the key buys — both actions refuse a channel + // that has running or queued jobs before they enqueue, which refuses rather + // than waits. "relocate-channel-media": { kind: "relocate-channel-media", label: "Relocate channel media", diff --git a/common/lib/queueKeys.test.ts b/common/lib/queueKeys.test.ts @@ -0,0 +1,52 @@ +import test from "node:test"; +import assert from "node:assert/strict"; + +import { + BACKFILL_QUEUE, + DIGEST_LOCAL_QUEUE, + DIGEST_REMOTE_QUEUE, + channelQueueKey, + relocationQueueKey, + resolveQueueKey, +} from "./queueKeys"; + +// THE ONE PROPERTY THAT IS NOT OBVIOUS FROM READING THE CALLER. +// +// registry.ts submits every non-empty queueKey at concurrency 1 and caps +// nothing across keys, so "do these two jobs serialize?" is answered entirely by +// whether their keys are EQUAL. For a media relocation the answer has to be yes +// even across channels — N moves selected on /channels write to one destination +// volume, and each job's space check runs when it starts, so concurrent starts +// each measure a root the others have not written to yet and jointly overrun it. +test("two channels' relocations land on one queue, and their bookkeeping does not", () => { + // THE CLAIM: two channels, one key. It is stated as a function of the slug on + // both sides so the test still reads as the claim if someone gives + // relocationQueueKey a slug parameter — at which point it fails, which is the + // point. channelQueueKey is the control: that one SHOULD differ per channel. + const relocationKeyFor = (_slug: string) => relocationQueueKey(); + assert.equal(relocationKeyFor("omnimirror"), relocationKeyFor("rekietalaw")); + assert.notEqual(channelQueueKey("omnimirror"), channelQueueKey("rekietalaw")); + // A relocation does not sit on the channel's bookkeeping queue either — that + // key is per-channel, so using it is the bug above by another route. + assert.notEqual(relocationQueueKey(), channelQueueKey("omnimirror")); + // Non-empty, or the registry would run it immediately and untracked. + assert.notEqual(relocationQueueKey(), ""); +}); + +// The lane keys are distinct from each other and from the relocation queue, so +// a move never blocks (or is blocked by) a sweep. +test("the named lane queues stay distinct", () => { + const keys = [ + BACKFILL_QUEUE, + DIGEST_LOCAL_QUEUE, + DIGEST_REMOTE_QUEUE, + relocationQueueKey(), + ]; + assert.equal(new Set(keys).size, keys.length); +}); + +test("resolveQueueKey: undefined keeps the default, a string overrides, blank is immediate", () => { + assert.equal(resolveQueueKey("channel:a", undefined), "channel:a"); + assert.equal(resolveQueueKey("channel:a", " relocate "), "relocate"); + assert.equal(resolveQueueKey("channel:a", ""), ""); +}); diff --git a/common/lib/queueKeys.ts b/common/lib/queueKeys.ts @@ -35,6 +35,35 @@ export function channelQueueKey(slug: string): string { return `channel:${slug}`; } +// ONE QUEUE FOR EVERY MEDIA RELOCATION IN THE PROCESS, and it takes no slug on +// purpose. +// +// The obvious key here is channelQueueKey(slug) — a move is a channel-local +// operation and it does have to serialize against that channel's own downloads +// and transcriptions. But registry.ts submits every non-empty key at +// concurrency 1 and has no cap ACROSS keys, so a per-channel key serializes a +// channel only against itself: ticking ten rows on /channels and pressing Move +// starts ten rsyncs at once, all writing to the SAME destination volume. That +// is wrong twice over. It is slower — ten interleaved sequential writes to one +// spinning platter is the worst access pattern the drive has — and it breaks +// the space check, which each job runs when it STARTS: ten jobs that all start +// together each measure a root none of the others has written to yet, all nine +// of the later ones are credited room that is already spoken for, and they +// jointly overrun it. The failure is recoverable (ENOSPC aborts the copy and the +// source is untouched until a verify passes) but it costs hours of copying to +// learn something one queue slot would have known. +// +// So a relocation runs behind every other relocation, and the run-time space +// check then measures a root that no other move is writing to. Serializing +// against the channel's OWN jobs is not lost: both callers refuse a channel that +// has running or queued jobs before they enqueue, which is a stronger rule than +// the queue's (it refuses rather than waits, for a move the operator can see is +// already in flight). +export const RELOCATION_QUEUE = "relocate"; +export function relocationQueueKey(): string { + return RELOCATION_QUEUE; +} + // Queue a network-bound download/pipeline job lands on: the channel's platform // queue, falling back to a per-domain queue for unrecognized hosts. export function downloadQueueKey(config: ChannelConfig): string { diff --git a/editor/app/channels/bulkStorageActions.ts b/editor/app/channels/bulkStorageActions.ts @@ -2,25 +2,33 @@ // MOVE SEVERAL CHANNELS' MEDIA TO THE COLD ROOT, FROM /channels. // -// ONE JOB PER CHANNEL, on that channel's own `channelQueueKey`. There is no -// batch controller and there is deliberately not going to be one: the job queue -// already serializes a channel against its own downloads and transcriptions, -// which is the ordering that matters, and a cross-channel batch would have to -// reinvent cancellation, resume and the per-channel log that already exist. +// ONE JOB PER CHANNEL, ALL ON ONE QUEUE — `relocationQueueKey()`, the shared key +// the per-channel Storage panel uses too. There is no batch controller and there +// is deliberately not going to be one: the job queue is the batch, and a +// cross-channel controller would have to reinvent cancellation, resume and the +// per-channel log that already exist. +// +// THE SHARED KEY IS LOAD-BEARING, not tidiness. registry.ts caps a queue key at +// concurrency 1 and caps NOTHING across keys, so a per-channel key would start +// every selected channel's rsync at once, onto one destination volume. See +// relocationQueueKey() for why that breaks the space check as well as the disk. // // EVERY CHECK RUNS TWICE, AND THAT IS THE DESIGN. What this action refuses at -// ENQUEUE time is what it can see now — a social channel, a channel already -// relocated, one mid-transition, one with live jobs. What it cannot see is the -// state of the destination twenty minutes from now, so each job runs its own -// `previewRelocation`-equivalent preflight (space, writability, containment, -// movable state) when it STARTS. A root that fills up partway through a -// selection therefore refuses the remainder cleanly, one job at a time, instead -// of the whole bulk failing at the point the first one is enqueued. +// ENQUEUE time is what it can see now — a social channel, a channel with no +// media to move, one already relocated, one mid-transition, one with live jobs. +// What it cannot see is the state of the destination twenty minutes from now, so +// each job runs its own `previewRelocation`-equivalent preflight (space, +// writability, containment, movable state) when it STARTS. Because the jobs run +// one at a time, that check measures a root no other move is writing to — so a +// root that fills up partway through a selection refuses the remainder cleanly, +// one job at a time, instead of the whole bulk failing at the point the first +// one is enqueued. // // A skip is never a failure: the result names every slug that did not queue and // why, and the bar renders both numbers. import path from "node:path"; +import { stat } from "node:fs/promises"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { getSettings } from "yt-dlp-transcript-common/lib/settings"; import { getRegistry } from "yt-dlp-transcript-common/jobs/registry"; @@ -35,6 +43,18 @@ import { queueForSlugs, type QueueOutcome } from "./lib/queueForSlugs"; // so the bar renders it the same way they do. export type BulkRelocateResult = QueueOutcome; +// Follows the link, deliberately: for a relocated channel `data/` is a symlink +// and what matters is whether the thing it points at is there. (This path only +// reaches it for a channel inspect() already called in-place, so in practice it +// is a real directory or nothing.) +async function isDirectory(p: string): Promise<boolean> { + try { + return (await stat(p)).isDirectory(); + } catch { + return false; + } +} + // The rename's guard, per channel — the copy in [slug]/storageActions.ts, for // the reason stated there: a download writing into `data/` while its bytes are // being copied out either fails the verify (safe) or is lost (not). @@ -88,6 +108,16 @@ export async function bulkRelocateChannelMediaAction( const jobs = activeJobsRefusal(slug); if (jobs) return jobs; const media = await inspectChannelMedia(paths, slug, config); + // NOTHING TO MOVE IS A SKIP, NOT A JOB. A channel that has downloaded + // nothing has no `data/` at all, and inspect() calls that `in-place` — + // correctly, since it is not relocated. Without this the bulk path queues + // a job whose only act is to throw the same sentence from + // relocateChannelMedia.ts's copy phase, after a job record, a log and a + // queue slot. Selecting a whole page of channels is the ordinary way this + // is used, and on a fresh corpus most of them are this case. + if (!media.relocated && !(await isDirectory(media.dataDir))) { + return "nothing to move — no media has been downloaded for it yet"; + } if (media.marker) { return `a relocation (${media.marker.direction}) to ${media.marker.target} is already in flight`; } diff --git a/editor/app/channels/components/ChannelStorageBulkBar.tsx b/editor/app/channels/components/ChannelStorageBulkBar.tsx @@ -12,8 +12,10 @@ // preview measures ONE channel's tree against one volume, and the useful answer // for a batch ("will all of these fit") is not the sum of the previews — the // root fills up as the jobs run. Each job therefore re-checks space when it -// STARTS, so a root that fills partway refuses the remainder one job at a time -// with the reason in that job's log. What the bar owes the operator is the +// STARTS, and because every relocation in the process shares one queue key the +// jobs run ONE AT A TIME — so that check measures a root no other move is +// writing to, and a root that fills partway refuses the remainder one job at a +// time with the reason in that job's log. What the bar owes the operator is the // count and the skips, not a number that would be stale before the second job. import { useState, useTransition } from "react"; diff --git a/editor/app/channels/lib/relocationJob.ts b/editor/app/channels/lib/relocationJob.ts @@ -1,6 +1,6 @@ import { revalidatePath } from "next/cache"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; -import { channelQueueKey } from "yt-dlp-transcript-common/lib/queueKeys"; +import { relocationQueueKey } from "yt-dlp-transcript-common/lib/queueKeys"; import { runManagedFunction, type StreamActionResult, @@ -25,6 +25,13 @@ import type { RelocationDirection } from "yt-dlp-transcript-common/lib/channelMe // the job record's shape — kind, queue key, channel slug, the summary line and // the three revalidations — because those are what must not drift between a // single move and a bulk one. +// +// THE QUEUE KEY IS THE SHARED ONE AND TAKES NO SLUG, which is the whole reason +// this helper exists rather than each caller building its own record. See +// relocationQueueKey() in common/lib/queueKeys.ts: a per-channel key would +// serialize a channel only against itself, so a bulk move of ten channels would +// start ten rsyncs onto one destination volume whose run-time space checks would +// then each be credited room the others had already claimed. export async function enqueueRelocation(opts: { slug: string; direction: RelocationDirection; @@ -34,7 +41,7 @@ export async function enqueueRelocation(opts: { const paths = getPaths(); return runManagedFunction({ kind: "relocate-channel-media", - queueKey: channelQueueKey(slug), + queueKey: relocationQueueKey(), paths, channelSlug: slug, fn: async (onLog, signal) => { diff --git a/editor/e2e/channel-storage.spec.ts b/editor/e2e/channel-storage.spec.ts @@ -298,3 +298,96 @@ test("the /channels bulk move queues one job per channel and skips the rest", as await pathExists(`test-transcripts/channels/${PRE}/.relocating.json`), ).toBe(false); }); + +// ONE QUEUE FOR EVERY MOVE, AND A CHANNEL WITH NOTHING TO MOVE NEVER BECOMES A +// JOB. +// +// The queue key is the ONLY thing that decides whether two relocations run at +// once: registry.ts caps a key at concurrency 1 and caps nothing across keys. +// A per-channel key therefore starts every selected channel's rsync together, +// onto one destination volume, and each job's space check runs when it STARTS — +// so concurrent starts are each credited room the others have already claimed. +// The shared key is what the bulk bar's "refuses the remainder one job at a +// time" sentence actually rests on, so it is asserted rather than described: +// both jobs land on `relocate`, read off the /jobs table's own queue column. +// +// The empty channel is the second half: `inspect()` calls a channel that has +// downloaded nothing `in-place` — correctly, it is not relocated — so without an +// up-front skip it becomes a job record, a log and a queue slot whose only act +// is to throw "has no data/ to move". On a fresh corpus most of a page is that. +test("a bulk move puts every job on one queue and skips a channel with nothing to move", async ({ + page, +}, testInfo) => { + test.setTimeout(120_000); + await resetData("one-youtube-channel-with-data"); + const root = testInfo.outputPath("queue-root"); + await mkdir(root, { recursive: true }); + await writeSettings({ + adminTitle: "Test Admin", + storage: { mediaRoot: root }, + }); + + // A second channel with real media, so the selection queues TWO moves — one + // queue key is only a claim about two of them. + const SECOND = "second-mover"; + await writeChannelConfig(SECOND, { + handling: "youtube", + name: "Second Mover", + url: "https://www.youtube.com/@second/videos", + }); + await mkdir( + resolvePath(`test-transcripts/channels/${SECOND}/data/20240103_second12345`), + { recursive: true }, + ); + await writeFile( + resolvePath( + `test-transcripts/channels/${SECOND}/data/20240103_second12345/transcript.en.vtt`, + ), + "WEBVTT\n\n00:00.000 --> 00:01.000\nhello\n", + ); + + // ...and a third with a config and no data/ at all: the skip. + const EMPTY = "no-media-yet"; + await writeChannelConfig(EMPTY, { + handling: "youtube", + name: "No Media Yet", + url: "https://www.youtube.com/@empty/videos", + }); + + await page.goto("/channels"); + await page.getByLabel("select all channels").check(); + await page.getByLabel("move media for selected channels").click(); + + const result = page.getByLabel("bulk media move result"); + await expect(result).toContainText("Queued 2 · skipped 1", { + timeout: 30_000, + }); + await expect(result).toHaveAttribute( + "title", + `${EMPTY}: nothing to move — no media has been downloaded for it yet`, + ); + + // Both moves land, through the badge, and then both jobs are on ONE queue. + await expect + .poll( + async () => { + await page.reload(); + return page.getByLabel(/^media location: Media relocated/).count(); + }, + { timeout: 90_000 }, + ) + .toBe(2); + + await page.goto("/jobs"); + const rows = page + .getByRole("row") + .filter({ hasText: "Relocate channel media" }); + await expect(rows).toHaveCount(2); + // The queue column: `relocate` for both. The discriminating half is the + // negative — a per-channel key prints `channel:<slug>` there, and that is + // exactly the shape that would let the two run at once. + for (const i of [0, 1]) { + await expect(rows.nth(i)).toContainText("relocate"); + await expect(rows.nth(i)).not.toContainText("channel:"); + } +});