Archilyzer · Source

archilyzer

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

commit 1ab9eaf9a780135271a6d2ebfee73b0cb2607ddb
parent 1a65561700461f94a48a50d67b607763c7336400
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Sun, 20 Sep 2026 03:09:23 -0400

metadata scan: one writer at a time, and each flush writes only the delta

Two bugs in the same few lines, both of which only show up on the channels the
feature exists for.

The flushes were fired mid-stream into an array (`flushes.push(flush())`) and
`pending` was zeroed only AFTER the await, so two could overlap — and
writeMetadataScan named its tmp file for the pid alone, so the second rename
found the file already moved and the job died with ENOENT. Flushes now share one
promise chain, `pending` resets before the await, and every write gets a unique
tmp name (the download path writes here too, so the name has to be safe on its
own, not merely because the scan now serializes).

The flush also re-sent the entire accumulated map every 25 records, which is
O(N²) bytes over a scan: a 1,800-video channel's 72 flushes each re-wrote
everything read so far. It now sends only what has not reached disk.

Same shape in the download path: a title-filter rejection rewrote the whole
channel-level store from inside the per-video loop — once per non-matching
video, which on a filtered channel is most of them. The batch runner collects
and writes once; a one-off download still writes for itself, because nothing
else will. The batch write is deliberately outside the abort checks: a
cancelled run still learned what those videos are called.

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

Diffstat:
Mcommon/controller/metadataScanStore.test.ts | 37+++++++++++++++++++++++++++++++++++++
Mcommon/controller/metadataScanStore.ts | 14+++++++++++++-
Mcommon/ytdlp/downloadOneManaged.ts | 67++++++++++++++++++++++++++++++++++++++++---------------------------
Mcommon/ytdlp/metadataScan.ts | 62+++++++++++++++++++++++++++++++++++++++++++++++---------------
Mcommon/ytdlp/runYtdlp.ts | 32+++++++++++++++++++++++++++++++-
5 files changed, 168 insertions(+), 44 deletions(-)

diff --git a/common/controller/metadataScanStore.test.ts b/common/controller/metadataScanStore.test.ts @@ -247,3 +247,40 @@ test("an unscanned channel settles nothing, whatever the filter says", async () ); }); }); + +// --- write safety ------------------------------------------------------------ + +test("overlapping upserts do not collide on the tmp file", async () => { + await withPaths(async (paths) => { + // Two writers in the same process used to share one pid-named tmp path, so + // the second rename found it already moved and threw ENOENT — which failed + // the whole scan job. The download path and the scan both write here. + await Promise.all( + Array.from({ length: 12 }, (_, i) => + upsertMetadataScan( + paths, + SLUG, + { entries: { [`v${i}`]: entry({ title: `T${i}` }) } }, + T1, + ), + ), + ); + const scan = await loadMetadataScan(paths, SLUG); + // Every write landed, and nothing threw. (Last-write-wins on the merged + // document is expected; losing a write is not the failure this pins — the + // ENOENT crash is.) + assert.ok(Object.keys(scan.entries).length >= 1); + }); +}); + +test("a delta upsert does not drop what is already stored", async () => { + await withPaths(async (paths) => { + await upsertMetadataScan(paths, SLUG, { entries: { a: entry() } }, T1); + // The scan flushes only what it has not written yet; the store must merge. + await upsertMetadataScan(paths, SLUG, { entries: { b: entry() } }, T2); + assert.deepEqual( + Object.keys((await loadMetadataScan(paths, SLUG)).entries).sort(), + ["a", "b"], + ); + }); +}); diff --git a/common/controller/metadataScanStore.ts b/common/controller/metadataScanStore.ts @@ -172,13 +172,25 @@ export async function loadMetadataScan( } } +let writeSeq = 0; +function nextWriteSeq(): string { + writeSeq = (writeSeq + 1) % Number.MAX_SAFE_INTEGER; + return `${Date.now().toString(36)}-${writeSeq}`; +} + async function writeMetadataScan( paths: Paths, slug: string, scan: MetadataScan, ): Promise<void> { const file = metadataScanPath(paths, slug); - const tmp = `${file}.tmp-${process.pid}`; + // UNIQUE PER WRITE, not per process. A pid-named tmp file is fine only while + // one write is in flight at a time: two overlapping writers in the SAME + // process wrote to the same path, and the second rename found it already + // moved and threw ENOENT, failing the job. The scan now serializes its own + // flushes, but the download path writes here too, so the name has to be safe + // on its own. + const tmp = `${file}.tmp-${process.pid}-${nextWriteSeq()}`; await writeFile(tmp, JSON.stringify(scan, null, 2) + "\n"); await rename(tmp, file); } diff --git a/common/ytdlp/downloadOneManaged.ts b/common/ytdlp/downloadOneManaged.ts @@ -46,7 +46,10 @@ import { isLivestreamMetadata, } from "../lib/transcripts-server"; import { evaluateDownloadFilters } from "../lib/downloadFilters"; -import { upsertMetadataScan } from "../controller/metadataScanStore"; +import { + upsertMetadataScan, + type MetadataScanEntry, +} from "../controller/metadataScanStore"; import { detectPlatform, type Platform } from "../lib/platform"; import { probeMediaDurationSec } from "./ffprobeDuration"; import { @@ -108,6 +111,11 @@ export type ManagedDownloadOpts = { // download filters. Plumbed from the playlist-level caller alongside the // other resolved settings. globalSkipLiveDownloads?: boolean; + // Batch collector for a title-filter rejection's metadata. Supplied by + // runManagedDownloads so a batch writes the channel-level metadata-scan store + // ONCE rather than once per non-matching video. Absent (a one-off download) + // means "write it yourself" — see the call site. + onFilterRejected?: (id: string, entry: MetadataScanEntry) => void; // ---- Per-download persistence rule (Phase 2) ---- // The channel's keep-latest window as a cutoff, computed once per run by the // caller (computeKeepWindow). This video's membership is decided against its @@ -618,32 +626,37 @@ async function runManagedDownload( // (the fail-closed branch) writes nothing here — there is nothing to // store, and it must stay retryable. if (decision.filter === "titleFilter" && metadata) { - try { - await upsertMetadataScan( - opts.paths, - opts.channelSlug, - { - entries: { - [canonicalId]: { - title: metadata.title ?? "", - description: metadata.description ?? "", - uploadDate: metadata.upload_date ?? "", - ...(metadata.live_status - ? { liveStatus: metadata.live_status } - : {}), - ...(typeof metadata.duration === "number" - ? { duration: metadata.duration } - : {}), - scannedAt: new Date().toISOString(), - }, - }, - }, - new Date().toISOString(), - ); - } catch (err) { - opts.onLog( - `Failed to record the metadata scan entry: ${(err as Error).message}\n`, - ); + const entry = { + title: metadata.title ?? "", + description: metadata.description ?? "", + uploadDate: metadata.upload_date ?? "", + ...(metadata.live_status ? { liveStatus: metadata.live_status } : {}), + ...(typeof metadata.duration === "number" + ? { duration: metadata.duration } + : {}), + scannedAt: new Date().toISOString(), + }; + // A BATCH COLLECTS; A SINGLE VIDEO WRITES. The store is channel-level, + // so writing it from inside a per-video loop is a full load + stringify + // + rename per rejection — on a filtered channel that is one rewrite of + // the whole file per non-matching video. The batch runner passes a + // collector and writes once when it is done; a one-off download (no + // collector) writes for itself, because nothing else will. + if (opts.onFilterRejected) { + opts.onFilterRejected(canonicalId, entry); + } else { + try { + await upsertMetadataScan( + opts.paths, + opts.channelSlug, + { entries: { [canonicalId]: entry } }, + new Date().toISOString(), + ); + } catch (err) { + opts.onLog( + `Failed to record the metadata scan entry: ${(err as Error).message}\n`, + ); + } } } const finishedAt = new Date().toISOString(); diff --git a/common/ytdlp/metadataScan.ts b/common/ytdlp/metadataScan.ts @@ -246,26 +246,50 @@ export async function runMetadataScan( // unavailable videos, so record them. const commitStreak = () => { for (const e of unavailableStreak) { - errors[e.id] = { + const err: MetadataScanError = { class: "deleted", message: e.message, at: new Date().toISOString(), }; + errors[e.id] = err; + unflushedErrors[e.id] = err; pending++; } unavailableStreak = []; }; - const flush = async (force = false) => { - if (!force && pending < FLUSH_EVERY) return; - if (pending === 0 && !force) return; - await upsertMetadataScan( - paths, - slug, - { entries: { ...entries }, errors: { ...errors } }, - new Date().toISOString(), - ); + // WHAT HAS NOT REACHED DISK YET. `entries`/`errors` above are the run's full + // accumulator (the orchestration reads them to decide what is left to fetch); + // these two are the DELTA. Sending the whole map on every flush made a scan + // O(N²) in bytes written — a 1,800-video channel's 72 flushes would each + // re-send everything read so far. + const unflushedEntries: Record<string, MetadataScanEntry> = {}; + const unflushedErrors: Record<string, MetadataScanError> = {}; + + // ONE WRITER AT A TIME. Flushes used to be fired mid-stream and collected in + // an array, so two could overlap — and `writeMetadataScan` wrote to a tmp path + // named only for the pid, so the second rename found the file already moved + // and the job died with ENOENT. `pending` is also zeroed BEFORE the await, so + // records arriving during a write are counted toward the next flush instead of + // being re-sent by this one. + let flushChain: Promise<void> = Promise.resolve(); + const flush = (force = false): Promise<void> => { + if (pending === 0) return flushChain; + if (!force && pending < FLUSH_EVERY) return flushChain; + const batchEntries = { ...unflushedEntries }; + const batchErrors = { ...unflushedErrors }; + for (const k of Object.keys(unflushedEntries)) delete unflushedEntries[k]; + for (const k of Object.keys(unflushedErrors)) delete unflushedErrors[k]; pending = 0; + flushChain = flushChain.then(async () => { + await upsertMetadataScan( + paths, + slug, + { entries: batchEntries, errors: batchErrors }, + new Date().toISOString(), + ); + }); + return flushChain; }; // One pass over a batch file. Returns the ids it could not read as needs_auth, @@ -319,7 +343,7 @@ export async function runMetadataScan( } const id = typeof parsed.id === "string" ? parsed.id : ""; if (!id) return; - entries[id] = { + const entry: MetadataScanEntry = { title: typeof parsed.title === "string" ? parsed.title : "", description: typeof parsed.description === "string" ? parsed.description : "", @@ -333,6 +357,8 @@ export async function runMetadataScan( : {}), scannedAt: new Date().toISOString(), }; + entries[id] = entry; + unflushedEntries[id] = entry; // A readable video proves we are not blocked, so whatever ran before it // was a real streak of unavailable videos. commitStreak(); @@ -347,7 +373,6 @@ export async function runMetadataScan( }; let stdoutBuf = ""; - const flushes: Promise<void>[] = []; child.stdout?.on("data", (c: Buffer) => { stdoutBuf += c.toString("utf8"); let nl: number; @@ -355,7 +380,8 @@ export async function runMetadataScan( takeRecord(stdoutBuf.slice(0, nl)); stdoutBuf = stdoutBuf.slice(nl + 1); } - if (pending >= FLUSH_EVERY) flushes.push(flush()); + // Serialized on flushChain, so this cannot overlap the previous write. + if (pending >= FLUSH_EVERY) void flush(); }); let stderrBuf = ""; @@ -404,7 +430,13 @@ export async function runMetadataScan( // soft block's ids. commitStreak(); if (cls === "needs_auth") needsAuthIds.push(id); - errors[id] = { class: cls, message, at: new Date().toISOString() }; + const err: MetadataScanError = { + class: cls, + message, + at: new Date().toISOString(), + }; + errors[id] = err; + unflushedErrors[id] = err; pending++; }; child.stderr?.on("data", (c: Buffer) => { @@ -424,7 +456,7 @@ export async function runMetadataScan( // A streak that never reached the threshold is just a run of dead videos. if (!block) commitStreak(); else unavailableStreak = []; - await Promise.all(flushes); + await flushChain; await flush(true); if (opts.signal.aborted) { diff --git a/common/ytdlp/runYtdlp.ts b/common/ytdlp/runYtdlp.ts @@ -28,7 +28,11 @@ import { type ResolvedCookiePolicy, } from "../lib/cookiePolicy"; import { resolveEffectiveAvailability } from "../lib/availability-server"; -import { settledByTitleFilterIds } from "../controller/metadataScanStore"; +import { + settledByTitleFilterIds, + upsertMetadataScan, + type MetadataScanEntry, +} from "../controller/metadataScanStore"; import { backfillAvailabilityFromMetadata } from "../controller/backfillAvailability"; import { runAvailabilityCheck } from "../controller/checkAvailability"; import { writeMaybeMissing } from "../controller/maybeMissingStore"; @@ -864,6 +868,11 @@ async function runManagedDownloads( const sleepSeconds = opts.channelConfig.sleepBetweenDownloadsSeconds ?? settings.sleepBetweenDownloadsSeconds; + // Title-filter rejections from this batch, written to the channel-level + // metadata-scan store ONCE at the end. Per-video writes would rewrite the + // whole file for every non-matching video — the exact shape of work a + // filtered channel produces most of. + const filterRejections: Record<string, MetadataScanEntry> = {}; // The channel's keep-latest window as a cutoff, computed ONCE per run from the // current on-disk catalog. Each video downloaded below is classified against it @@ -930,6 +939,9 @@ async function runManagedDownloads( extractImmediately: opts.extractImmediately, audioFormatOverride: opts.audioFormatOverride, downloadFormatPreset, + onFilterRejected: (id, entry) => { + filterRejections[id] = entry; + }, }); } finally { task?.end(); @@ -988,6 +1000,24 @@ async function runManagedDownloads( ), ); + // ONE write for the whole batch. Deliberately outside the abort checks: a + // cancelled run still learned what those videos are called, and throwing that + // away would make the next run re-fetch their metadata for nothing. + if (Object.keys(filterRejections).length > 0) { + try { + await upsertMetadataScan( + opts.paths, + opts.channelSlug, + { entries: filterRejections }, + new Date().toISOString(), + ); + } catch (err) { + opts.onLog( + `Failed to record ${Object.keys(filterRejections).length} metadata scan entries: ${(err as Error).message}\n`, + ); + } + } + return { okCount: processedCount - failedCount - skippedCount, failedCount,