Archilyzer · Source

archilyzer

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

commit af19fc627135cb50f6827a174d857eeed54c7c5a
parent 37d9a7de8446734a9529093c9fc473192869ac8c
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Thu, 24 Sep 2026 20:03:20 -0400

Merge one-core/phase-3-w — one write idiom: writeFileAtomic under writeJsonAtomic, 26 writers folded, the two write counters deleted

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

Diffstat:
Mcommon/bin/migrate-channel-priority.ts | 7+++----
Mcommon/controller/backupSavedVideos.ts | 7+++----
Acommon/controller/compactJsonWriters.test.ts | 98+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/controller/duplicateShorts.ts | 14+++++---------
Mcommon/controller/failedTranscriptions.ts | 14+++++++-------
Mcommon/controller/maybeMissingStore.ts | 10+++++-----
Mcommon/controller/metadataScanStore.ts | 28++++++++++------------------
Mcommon/controller/normalizeLiveChat.ts | 8++++----
Mcommon/controller/normalizeTranscript.ts | 8++++----
Mcommon/controller/relocateDir.ts | 7++-----
Mcommon/controller/rosterStore.ts | 13++++++-------
Mcommon/controller/scanCorruptMedia.ts | 14+++++---------
Mcommon/controller/shard.ts | 8+++-----
Mcommon/controller/storageWatch.ts | 4++++
Mcommon/jobs/autoQueueState.test.ts | 28+++++++++++++++++++++++++++-
Mcommon/jobs/autoQueueState.ts | 15++++++---------
Mcommon/jobs/syncSchedulerState.ts | 9+++------
Mcommon/jobs/workerDefaults.ts | 9++-------
Mcommon/lib/homepage.ts | 8+++-----
Mcommon/lib/jsonFile-server.test.ts | 81+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/lib/jsonFile-server.ts | 101++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------------------
Mcommon/lib/savedVideo-server.ts | 7++-----
Mcommon/lib/widgetPresets.ts | 9++-------
Mcommon/social/xSessionBroker.ts | 14++++++++------
Mcommon/ytdlp/runYtdlp.ts | 7+++----
Meditor/CHANGELOG.md | 1+
Meditor/app/channels/[slug]/videos/[id]/videoActions.ts | 13+++++--------
Meditor/app/sites/lib/cutReleaseAction.ts | 5++---
Mplans/one-core-phase-3.md | 119+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aplans/tools/phase3-writers-numbers.ts | 439+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
30 files changed, 936 insertions(+), 169 deletions(-)

diff --git a/common/bin/migrate-channel-priority.ts b/common/bin/migrate-channel-priority.ts @@ -1,6 +1,7 @@ #!/usr/bin/env tsx import path from "node:path"; -import { readdir, readFile, rename, writeFile } from "node:fs/promises"; +import { readdir, readFile, writeFile } from "node:fs/promises"; +import { writeJsonAtomic } from "../lib/jsonFile-server"; import { getPaths } from "../lib/paths"; import { getSettings, writeSettings } from "../lib/settings"; import { siteChannelIndex } from "../lib/site"; @@ -203,9 +204,7 @@ async function clearExcludeFromSync(slug: string): Promise<boolean> { } if (!("excludeFromSync" in parsed)) return false; delete parsed.excludeFromSync; - const tmp = `${file}.tmp-${process.pid}`; - await writeFile(tmp, JSON.stringify(parsed, null, 2) + "\n"); - await rename(tmp, file); + await writeJsonAtomic(file, parsed); return true; } diff --git a/common/controller/backupSavedVideos.ts b/common/controller/backupSavedVideos.ts @@ -1,7 +1,8 @@ import path from "node:path"; import { createHash } from "node:crypto"; import { createReadStream } from "node:fs"; -import { mkdir, readFile, readdir, rename, stat, writeFile } from "node:fs/promises"; +import { mkdir, readFile, readdir, stat } from "node:fs/promises"; +import { writeJsonAtomic } from "../lib/jsonFile-server"; import { execa } from "execa"; import type { Paths } from "../lib/paths"; import { @@ -115,9 +116,7 @@ export async function backupSavedVideos({ entries: manifestEntries, }; const manifestPath = path.join(destRoot, BACKUP_MANIFEST_FILENAME); - const tmp = `${manifestPath}.tmp-${process.pid}`; - await writeFile(tmp, JSON.stringify(manifest, null, 2) + "\n"); - await rename(tmp, manifestPath); + await writeJsonAtomic(manifestPath, manifest); log( `Backup complete: ${backedUp}/${saved.length} container(s), ${bytes} bytes. Manifest at ${manifestPath}`, ); diff --git a/common/controller/compactJsonWriters.test.ts b/common/controller/compactJsonWriters.test.ts @@ -0,0 +1,98 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import path from "node:path"; +import { mkdir, mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import type { Paths } from "../lib/paths"; +import { + CUES_JSON_FILENAME, + LIVE_CHAT_CUES_FILENAME, + LIVE_CHAT_FILENAME, + META_FILENAME, + VTT_FILENAME, +} from "../lib/videoStatus"; +import { DUPLICATES_FILENAME } from "../lib/duplicates"; +import { MEDIA_SCAN_FILENAME } from "../lib/mediaScan"; +import { normalizeTranscript } from "./normalizeTranscript"; +import { normalizeLiveChat } from "./normalizeLiveChat"; +import { detectDuplicateShorts } from "./duplicateShorts"; +import { scanCorruptMedia } from "./scanCorruptMedia"; + +// THE FOUR COMPACT WRITERS keep their historical bytes after the fold onto +// writeJsonAtomic (release 4 slice W): each wrote `JSON.stringify(x)` — no +// indent, NO trailing newline — and a reader-side diff (the phase-3 numbers +// tools, a site's published bytes) would move if that changed. Each test pins +// that the file on disk is exactly `JSON.stringify` of its own parse. + +async function withDir(fn: (dir: string) => Promise<void>): Promise<void> { + const dir = await mkdtemp(path.join(tmpdir(), "compact-writers-")); + try { + await fn(dir); + } finally { + await rm(dir, { recursive: true, force: true }); + } +} + +function assertCompact(text: string, what: string): void { + assert.ok(!text.endsWith("\n"), `${what}: no trailing newline`); + assert.equal(text, JSON.stringify(JSON.parse(text)), `${what}: compact`); +} + +async function videoDir(dir: string): Promise<string> { + const v = path.join(dir, "vid1"); + await mkdir(v, { recursive: true }); + await writeFile( + path.join(v, META_FILENAME), + JSON.stringify({ id: "vid1", title: "A video", duration: 60 }), + ); + return v; +} + +test("normalizeTranscript writes transcript.cues.json compact, no newline", async () => { + await withDir(async (dir) => { + const v = await videoDir(dir); + await writeFile( + path.join(v, VTT_FILENAME), + "WEBVTT\n\n00:00:00.000 --> 00:00:05.000\nhello\n", + ); + const out = await normalizeTranscript({ videoDir: v, channelSlug: "c" }); + assert.equal(out.status, "wrote"); + assertCompact(await readFile(path.join(v, CUES_JSON_FILENAME), "utf8"), "cues"); + }); +}); + +test("normalizeLiveChat writes the live-chat cues compact, no newline", async () => { + await withDir(async (dir) => { + const v = await videoDir(dir); + await writeFile(path.join(v, LIVE_CHAT_FILENAME), ""); + const out = await normalizeLiveChat({ videoDir: v, channelSlug: "c" }); + assert.equal(out.status, "wrote"); + assertCompact( + await readFile(path.join(v, LIVE_CHAT_CUES_FILENAME), "utf8"), + "live chat cues", + ); + }); +}); + +function emptyCorpus(dir: string): Paths { + return { + transcriptsDir: dir, + channelsDir: path.join(dir, "channels"), + } as Paths; +} + +test("detectDuplicateShorts writes duplicates.json compact, no newline", async () => { + await withDir(async (dir) => { + await mkdir(path.join(dir, "channels"), { recursive: true }); + await detectDuplicateShorts({ paths: emptyCorpus(dir), onLog: () => {} }); + assertCompact(await readFile(path.join(dir, DUPLICATES_FILENAME), "utf8"), "duplicates"); + }); +}); + +test("scanCorruptMedia writes media-scan.json compact, no newline", async () => { + await withDir(async (dir) => { + await mkdir(path.join(dir, "channels"), { recursive: true }); + await scanCorruptMedia({ paths: emptyCorpus(dir), onLog: () => {} }); + assertCompact(await readFile(path.join(dir, MEDIA_SCAN_FILENAME), "utf8"), "media scan"); + }); +}); diff --git a/common/controller/duplicateShorts.ts b/common/controller/duplicateShorts.ts @@ -38,7 +38,8 @@ // global transcripts/duplicates.json. Flag-only: nothing is merged or deleted. import path from "node:path"; -import { mkdir, rename, writeFile, readFile } from "node:fs/promises"; +import { readFile } from "node:fs/promises"; +import { writeJsonAtomic } from "../lib/jsonFile-server"; import { createHash } from "node:crypto"; import { open } from "lmdb"; import pLimit from "p-limit"; @@ -636,10 +637,8 @@ export async function detectDuplicateShorts( }; const outPath = path.join(opts.paths.transcriptsDir, DUPLICATES_FILENAME); - await mkdir(path.dirname(outPath), { recursive: true }); - const tmp = `${outPath}.tmp-${process.pid}`; - await writeFile(tmp, JSON.stringify(report)); - await rename(tmp, outPath); + // Compact, no trailing newline: the report's historical bytes. + await writeJsonAtomic(outPath, report, { indent: 0, newline: false, mkdir: true }); noteRss(); const suspectClusters = clusters.filter((c) => c.needsReview).length; log( @@ -730,10 +729,7 @@ export async function updateDuplicateOverride( } const out = sanitizeDuplicateOverrides(current); const file = duplicateOverridesPath(paths); - await mkdir(path.dirname(file), { recursive: true }); - const tmp = `${file}.tmp-${process.pid}`; - await writeFile(tmp, JSON.stringify(out, null, 2) + "\n"); - await rename(tmp, file); + await writeJsonAtomic(file, out, { mkdir: true }); return out; } diff --git a/common/controller/failedTranscriptions.ts b/common/controller/failedTranscriptions.ts @@ -1,5 +1,6 @@ import path from "node:path"; -import { readFile, rename, writeFile } from "node:fs/promises"; +import { readFile } from "node:fs/promises"; +import { writeFileAtomic } from "../lib/jsonFile-server"; import type { Paths } from "../lib/paths"; export function failedTranscriptionsFile(paths: Paths, slug: string): string { @@ -30,9 +31,10 @@ export async function pruneFailedTranscriptions( } const original = raw.split("\n").filter(Boolean); const filtered = original.filter((id) => !idsToRemove.has(id)); - const tmp = `${failureListFile}.tmp-${process.pid}`; - await writeFile(tmp, filtered.length ? filtered.join("\n") + "\n" : ""); - await rename(tmp, failureListFile); + await writeFileAtomic( + failureListFile, + filtered.length ? filtered.join("\n") + "\n" : "", + ); return { remaining: filtered.length, pruned: original.length - filtered.length, @@ -51,8 +53,6 @@ export async function clearFailedTranscriptions( } catch { return { cleared: 0 }; } - const tmp = `${file}.tmp-${process.pid}`; - await writeFile(tmp, ""); - await rename(tmp, file); + await writeFileAtomic(file, ""); return { cleared }; } diff --git a/common/controller/maybeMissingStore.ts b/common/controller/maybeMissingStore.ts @@ -1,5 +1,6 @@ import path from "node:path"; -import { readFile, rename, writeFile } from "node:fs/promises"; +import { readFile } from "node:fs/promises"; +import { writeJsonAtomic } from "../lib/jsonFile-server"; import type { Paths } from "../lib/paths"; // Leaf module for the maybe-missing record: the set of videos we have on disk @@ -50,10 +51,9 @@ export async function writeMaybeMissing( slug: string, record: MaybeMissingRecord, ): Promise<void> { - const file = maybeMissingPath(paths, slug); - const tmp = `${file}.tmp-${process.pid}`; - await writeFile(tmp, JSON.stringify(record, null, 2) + "\n"); - await rename(tmp, file); + // No mkdir: the channel dir exists. Chained with the roster/scan writes of + // every other actor in this process (lib/jsonFile-server.ts). + await writeJsonAtomic(maybeMissingPath(paths, slug), record); } // Diff a fresh listing's canonical ids against the videos we already have on diff --git a/common/controller/metadataScanStore.ts b/common/controller/metadataScanStore.ts @@ -24,10 +24,11 @@ // // Modelled on ./rosterStore.ts: versioned, normalize-on-read (unknown fields // dropped), additive merge that returns the same object when nothing changed, -// tmp+rename write. +// tmp+rename write (lib/jsonFile-server.ts). import path from "node:path"; -import { readFile, rename, writeFile } from "node:fs/promises"; +import { readFile } from "node:fs/promises"; +import { writeJsonAtomic } from "../lib/jsonFile-server"; import type { Paths } from "../lib/paths"; import type { ChannelConfig } from "../lib/channelConfig"; import { @@ -185,27 +186,18 @@ 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); - // 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); + // UNIQUE PER WRITE, not per process, and chained per path: two overlapping + // writers in the SAME process once shared a pid-named tmp and the second + // rename threw ENOENT, failing the job. The scan serializes its own flushes, + // but the download path writes here too. The shared writer's temp name and + // globalThis chain replace this module's own counter, which was per module + // COPY. + await writeJsonAtomic(metadataScanPath(paths, slug), scan); } export type MetadataScanUpsert = { diff --git a/common/controller/normalizeLiveChat.ts b/common/controller/normalizeLiveChat.ts @@ -4,7 +4,8 @@ // matches NormalizedTranscript with source: "live_chat". import path from "node:path"; -import { readFile, rename, stat, writeFile } from "node:fs/promises"; +import { readFile, stat } from "node:fs/promises"; +import { writeJsonAtomic } from "../lib/jsonFile-server"; import { parseLiveChat } from "../lib/liveChat"; import type { Cue } from "../lib/vtt"; import { summarize, type RawMetadata } from "../lib/transcripts-server"; @@ -95,9 +96,8 @@ export async function normalizeLiveChat( cues, }; - const tmp = `${cuesPath}.tmp-${process.pid}`; - await writeFile(tmp, JSON.stringify(out)); - await rename(tmp, cuesPath); + // Compact, no trailing newline: live_chat.cues.json's historical bytes. + await writeJsonAtomic(cuesPath, out, { indent: 0, newline: false }); opts.log?.( `Normalized live chat ${opts.channelSlug}/${path.basename(opts.videoDir)} (${cues.length} cues)`, ); diff --git a/common/controller/normalizeTranscript.ts b/common/controller/normalizeTranscript.ts @@ -4,7 +4,8 @@ // shape that buildIndex would emit, plus a `source` marker. import path from "node:path"; -import { readdir, readFile, rename, stat, writeFile } from "node:fs/promises"; +import { readdir, readFile, stat } from "node:fs/promises"; +import { writeJsonAtomic } from "../lib/jsonFile-server"; import { parseVtt, type Cue } from "../lib/vtt"; import { detectTranscriptFormat, @@ -136,9 +137,8 @@ export async function normalizeTranscript( cues, }; - const tmp = `${cuesPath}.tmp-${process.pid}`; - await writeFile(tmp, JSON.stringify(out)); - await rename(tmp, cuesPath); + // Compact, no trailing newline: transcript.cues.json's historical bytes. + await writeJsonAtomic(cuesPath, out, { indent: 0, newline: false }); opts.log?.( `Normalized ${opts.channelSlug}/${path.basename(opts.videoDir)} (${transcriptFormat}, ${cues.length} cues)`, ); diff --git a/common/controller/relocateDir.ts b/common/controller/relocateDir.ts @@ -4,11 +4,10 @@ import { readdir, readFile, readlink, - rename, rm, stat, - writeFile, } from "node:fs/promises"; +import { writeJsonAtomic } from "../lib/jsonFile-server"; import { execa } from "execa"; import { formatRsyncProgressDetail, @@ -136,9 +135,7 @@ export async function writeDirMarker( file: string, marker: DirRelocationMarker, ): Promise<void> { - const tmp = `${file}.tmp-${process.pid}`; - await writeFile(tmp, JSON.stringify(marker, null, 2) + "\n"); - await rename(tmp, file); + await writeJsonAtomic(file, marker); } export async function clearDirMarker(file: string): Promise<void> { diff --git a/common/controller/rosterStore.ts b/common/controller/rosterStore.ts @@ -1,6 +1,7 @@ import path from "node:path"; -import { mkdir, open, readdir, readFile, rename, writeFile } from "node:fs/promises"; +import { open, readdir, readFile } from "node:fs/promises"; import pLimit from "p-limit"; +import { writeJsonAtomic } from "../lib/jsonFile-server"; import type { Paths } from "../lib/paths"; import { extractVideoId } from "../lib/videoId"; @@ -224,17 +225,15 @@ export async function loadRoster(paths: Paths, slug: string): Promise<Roster> { // Atomic tmp+rename, like maybe-missing.json: a crashed write must never leave a // truncated roster behind, because a truncated roster is exactly the data loss -// this file exists to prevent. +// this file exists to prevent. Several actors write it (the sync, the quick +// availability check, the pipeline's server action — a different bundle), so +// it goes through the one per-path chain on globalThis. export async function writeRoster( paths: Paths, slug: string, roster: Roster, ): Promise<void> { - const file = rosterPath(paths, slug); - await mkdir(path.dirname(file), { recursive: true }); - const tmp = `${file}.tmp-${process.pid}`; - await writeFile(tmp, JSON.stringify(roster, null, 2) + "\n"); - await rename(tmp, file); + await writeJsonAtomic(rosterPath(paths, slug), roster, { mkdir: true }); } // Load, merge, write in one step — the shape every additive writer wants. Skips diff --git a/common/controller/scanCorruptMedia.ts b/common/controller/scanCorruptMedia.ts @@ -29,7 +29,8 @@ // first, delete second — not a dry-run flag on a deleting command. import path from "node:path"; -import { mkdir, readdir, readFile, rename, stat, writeFile } from "node:fs/promises"; +import { readdir, readFile, stat } from "node:fs/promises"; +import { writeJsonAtomic } from "../lib/jsonFile-server"; import { execa } from "execa"; import type { Paths } from "../lib/paths"; import { @@ -389,10 +390,8 @@ async function mergeAndWriteReport( ), }; const outPath = path.join(paths.transcriptsDir, MEDIA_SCAN_FILENAME); - await mkdir(path.dirname(outPath), { recursive: true }); - const tmp = `${outPath}.tmp-${process.pid}`; - await writeFile(tmp, JSON.stringify(merged)); - await rename(tmp, outPath); + // Compact, no trailing newline: the report's historical bytes. + await writeJsonAtomic(outPath, merged, { indent: 0, newline: false, mkdir: true }); return merged; } @@ -450,9 +449,6 @@ export async function updateMediaScanOverride( reviewed, }; const file = mediaScanOverridesPath(paths); - await mkdir(path.dirname(file), { recursive: true }); - const tmp = `${file}.tmp-${process.pid}`; - await writeFile(tmp, JSON.stringify(out, null, 2) + "\n"); - await rename(tmp, file); + await writeJsonAtomic(file, out, { mkdir: true }); return out; } diff --git a/common/controller/shard.ts b/common/controller/shard.ts @@ -1,5 +1,6 @@ import path from "node:path"; -import { readFile, rename, unlink, writeFile } from "node:fs/promises"; +import { readFile, unlink } from "node:fs/promises"; +import { writeJsonAtomic } from "../lib/jsonFile-server"; import type { Paths } from "../lib/paths"; export const SHARD_OPS = [ @@ -52,10 +53,7 @@ export async function saveShardConfig( op: ShardOp, cfg: ShardConfig, ): Promise<void> { - const file = shardFile(paths, slug, op); - const tmp = `${file}.tmp-${process.pid}`; - await writeFile(tmp, JSON.stringify(cfg, null, 2) + "\n"); - await rename(tmp, file); + await writeJsonAtomic(shardFile(paths, slug, op), cfg); } export async function clearShardConfig( diff --git a/common/controller/storageWatch.ts b/common/controller/storageWatch.ts @@ -318,6 +318,10 @@ export async function runStorageWatchPass( // of the process. export const STORAGE_WATCH_INTERVAL_MS = 5 * 60_000; +// A per-module-copy singleton, deliberately left so: it is not a temp-file +// name (slice W folded every tmp + rename onto lib/jsonFile-server.ts, whose +// state is on globalThis), and its one caller is editor/instrumentation.ts, +// so only one copy ever arms it. let timer: ReturnType<typeof setInterval> | null = null; // ARMED ONCE PER PROCESS. `unref()` so it never holds the event loop open — a diff --git a/common/jobs/autoQueueState.test.ts b/common/jobs/autoQueueState.test.ts @@ -1,6 +1,6 @@ import { test } from "node:test"; import assert from "node:assert/strict"; -import { mkdtemp, rm, writeFile } from "node:fs/promises"; +import { mkdtemp, readdir, rm, writeFile } from "node:fs/promises"; import { tmpdir } from "node:os"; import path from "node:path"; import type { Paths } from "../lib/paths"; @@ -126,3 +126,29 @@ test("a lane missing from an in-memory state object still writes", async () => { ]); }); }); + +// --- write safety ------------------------------------------------------------ + +test("overlapping writes do not collide on the tmp file", async () => { + await withPaths(async (paths) => { + // persist() is fire-and-forget and two lanes' runners can persist within + // milliseconds of each other. With one pid-named tmp the second write + // truncated the first's temp and one rename found it gone (ENOENT). + await Promise.all( + Array.from({ length: 12 }, (_, i) => { + const state = emptyAutoQueueState(); + state.download.platformBackoff = { youtube: { until: i, fails: i } }; + return writeAutoQueueState(paths, state); + }), + ); + // Every write landed and nothing threw; the last ISSUED is on disk (the + // writes are chained per path) and no temp is left behind. + const back = await readAutoQueueState(paths); + assert.deepEqual(back.download.platformBackoff, { + youtube: { until: 11, fails: 11 }, + }); + assert.deepEqual(await readdir(path.dirname(paths.autoQueueStateFile)), [ + "state.json", + ]); + }); +}); diff --git a/common/jobs/autoQueueState.ts b/common/jobs/autoQueueState.ts @@ -1,5 +1,5 @@ -import { mkdir, readFile, rename, writeFile } from "node:fs/promises"; -import path from "node:path"; +import { readFile } from "node:fs/promises"; +import { writeJsonAtomic } from "../lib/jsonFile-server"; import type { Paths } from "../lib/paths"; import { type AutoQueueRuntime, @@ -136,9 +136,9 @@ export async function readAutoQueueState(paths: Paths): Promise<AutoQueueState> // tmp name per process means the second write truncates the file the first is // still writing and then renames it over the real one — leaving a state.json // that reads as empty (tolerated) and loses a persisted platform cooldown -// across a restart (not). One counter is enough: this is one process, and the -// pid still separates two of them. -let writeSeq = 0; +// across a restart (not). `writeJsonAtomic` gives every write its own temp +// name and chains writes to one path on globalThis — which this module's own +// counter, one per module COPY, did not. export async function writeAutoQueueState( paths: Paths, @@ -152,10 +152,7 @@ export async function writeAutoQueueState( const out = Object.fromEntries( LANES.map((lane) => [lane, trim(state[lane] ?? emptyAutoQueueKindState())]), ) as AutoQueueState; - await mkdir(path.dirname(paths.autoQueueStateFile), { recursive: true }); - const tmp = `${paths.autoQueueStateFile}.tmp-${process.pid}-${++writeSeq}`; - await writeFile(tmp, JSON.stringify(out, null, 2) + "\n"); - await rename(tmp, paths.autoQueueStateFile); + await writeJsonAtomic(paths.autoQueueStateFile, out, { mkdir: true }); } export function recordPick(state: AutoQueueKindState, pick: AutoQueuePick): void { diff --git a/common/jobs/syncSchedulerState.ts b/common/jobs/syncSchedulerState.ts @@ -1,5 +1,5 @@ -import { mkdir, readFile, rename, writeFile } from "node:fs/promises"; -import path from "node:path"; +import { readFile } from "node:fs/promises"; +import { writeJsonAtomic } from "../lib/jsonFile-server"; import type { Paths } from "../lib/paths"; // Persistent state for the cron-driven sync scheduler. Unlike the in-memory job @@ -118,10 +118,7 @@ export async function writeSchedulerState( lastSavedVideoBackupAt: state.lastSavedVideoBackupAt ?? null, lastSyncAllAt: state.lastSyncAllAt ?? null, }; - await mkdir(path.dirname(paths.schedulerStateFile), { recursive: true }); - const tmp = `${paths.schedulerStateFile}.tmp-${process.pid}`; - await writeFile(tmp, JSON.stringify(out, null, 2) + "\n"); - await rename(tmp, paths.schedulerStateFile); + await writeJsonAtomic(paths.schedulerStateFile, out, { mkdir: true }); } // Get (or lazily create) the mutable state entry for a channel. diff --git a/common/jobs/workerDefaults.ts b/common/jobs/workerDefaults.ts @@ -1,5 +1,5 @@ import fs from "node:fs"; -import path from "node:path"; +import { writeJsonAtomic } from "../lib/jsonFile-server"; import type { Paths } from "../lib/paths"; // Persisted "default" worker arrangement. Unlike the in-memory worker pool @@ -55,10 +55,5 @@ export async function writeWorkerDefaults( if (typeof id === "string" && id) seen.add(id); } const out: WorkerDefaults = { enabledWorkerIds: Array.from(seen) }; - await fs.promises.mkdir(path.dirname(paths.workerDefaultsFile), { - recursive: true, - }); - const tmp = `${paths.workerDefaultsFile}.tmp-${process.pid}`; - await fs.promises.writeFile(tmp, JSON.stringify(out, null, 2) + "\n"); - await fs.promises.rename(tmp, paths.workerDefaultsFile); + await writeJsonAtomic(paths.workerDefaultsFile, out, { mkdir: true }); } diff --git a/common/lib/homepage.ts b/common/lib/homepage.ts @@ -1,4 +1,5 @@ import fs from "node:fs"; +import { writeJsonAtomic } from "./jsonFile-server"; import { getPaths, type Paths } from "./paths"; import { PROJECT_NAME, PROJECT_TAGLINE } from "./project"; import { @@ -124,9 +125,6 @@ export async function writeHomepageConfig( ? { cloudflareProject: config.cloudflareProject.trim() } : {}), }; - await fs.promises.mkdir(paths.homepageDir, { recursive: true }); - const file = paths.homepageConfigFile; - const tmp = `${file}.tmp-${process.pid}`; - await fs.promises.writeFile(tmp, JSON.stringify(merged, null, 2) + "\n"); - await fs.promises.rename(tmp, file); + // mkdir: the parent of homepageConfigFile IS homepageDir (lib/paths.ts). + await writeJsonAtomic(paths.homepageConfigFile, merged, { mkdir: true }); } diff --git a/common/lib/jsonFile-server.test.ts b/common/lib/jsonFile-server.test.ts @@ -2,9 +2,11 @@ import { test } from "node:test"; import assert from "node:assert/strict"; import fs from "node:fs"; import { mkdtemp, readFile, readdir, writeFile } from "node:fs/promises"; +import { createRequire, syncBuiltinESMExports } from "node:module"; import os from "node:os"; import path from "node:path"; import { + copyFileAtomic, jsonFileState, jsonText, pendingJsonWrites, @@ -12,9 +14,12 @@ import { readJsonFileSync, tmpPathFor, writeJsonAtomic, + writeFileAtomic, writeJsonAtomicSync, } from "./jsonFile-server"; +const require = createRequire(import.meta.url); + async function scratch(): Promise<string> { return mkdtemp(path.join(os.tmpdir(), "jsonfile-")); } @@ -123,3 +128,79 @@ test("writeJsonAtomicSync: same bytes, mkdir, no temp left", async () => { assert.equal(fs.readFileSync(file, "utf8"), '{\n "k": 1\n}'); assert.deepEqual(await readdir(path.dirname(file)), ["b.json"]); }); + +test("writeFileAtomic: exact bytes for a string and a Buffer, mkdir, no temp left", async () => { + const dir = await scratch(); + const txt = path.join(dir, "d", "playlist"); + await writeFileAtomic(txt, "a\nb\n", { mkdir: true }); + assert.equal(await readFile(txt, "utf8"), "a\nb\n"); + const bin = path.join(dir, "d", "blob"); + const bytes = Buffer.from([0, 255, 10, 13, 128]); + await writeFileAtomic(bin, bytes); + assert.deepEqual(await readFile(bin), bytes); + // An empty string is a real write (clearFailedTranscriptions relies on it). + await writeFileAtomic(txt, ""); + assert.equal(await readFile(txt, "utf8"), ""); + assert.deepEqual((await readdir(path.join(dir, "d"))).sort(), ["blob", "playlist"]); +}); + +test("writeFileAtomic: mode is on the file from creation (0o600 cookie jar)", async () => { + const dir = await scratch(); + const file = path.join(dir, "cookies.txt"); + await writeFileAtomic(file, "secret\n", { mode: 0o600 }); + assert.equal(fs.statSync(file).mode & 0o777, 0o600); + // The temp itself carries the mode BEFORE the rename: intercept the rename + // (the module calls fs/promises' `rename`, so patch it and sync the builtin + // ESM exports) and stat its source — the temp — at that moment. + const fsp = require("node:fs/promises") as typeof import("node:fs/promises"); + const realRename = fsp.rename; + const seen: number[] = []; + fsp.rename = (async (from: fs.PathLike, to: fs.PathLike) => { + seen.push(fs.statSync(from).mode & 0o777); + return realRename(from, to); + }) as typeof fsp.rename; + syncBuiltinESMExports(); + try { + await writeFileAtomic(file, "second\n", { mode: 0o600 }); + } finally { + fsp.rename = realRename; + syncBuiltinESMExports(); + } + assert.deepEqual(seen, [0o600], "the rename saw exactly one temp, already 0o600"); + assert.equal(await readFile(file, "utf8"), "second\n"); +}); + +test("writeFileAtomic and writeJsonAtomic share one chain on one path", async () => { + const dir = await scratch(); + const file = path.join(dir, "shared.json"); + const writes: Promise<void>[] = []; + for (let i = 0; i < 40; i++) { + writes.push( + i % 2 === 0 + ? writeJsonAtomic(file, { i }) + : writeFileAtomic(file, `{"i":${i}}`), + ); + } + assert.equal(pendingJsonWrites(), 1, "one chain entry for the path, not two"); + await Promise.all(writes); + assert.deepEqual(JSON.parse(await readFile(file, "utf8")), { i: 39 }); + assert.deepEqual(await readdir(dir), ["shared.json"]); +}); + +test("copyFileAtomic: dest gets src's bytes, src stays, no temp left", async () => { + const dir = await scratch(); + const src = path.join(dir, "src.bin"); + await writeFile(src, Buffer.from([1, 2, 3])); + const dest = path.join(dir, "out", "dest.bin"); + await copyFileAtomic(src, dest, { mkdir: true }); + assert.deepEqual(await readFile(dest), Buffer.from([1, 2, 3])); + assert.deepEqual(await readdir(path.join(dir, "out")), ["dest.bin"]); + await assert.rejects(copyFileAtomic(path.join(dir, "nope"), dest)); + assert.deepEqual(await readdir(path.join(dir, "out")), ["dest.bin"]); +}); + +test("the remark-empty transcript is the literal it replaced (videoActions)", () => { + // editor/app/channels/[slug]/videos/[id]/videoActions.ts wrote this literal + // by hand before slice W; it now writes the object with { indent: 0 }. + assert.equal(jsonText({ transcription: [] }, { indent: 0 }), '{"transcription":[]}\n'); +}); diff --git a/common/lib/jsonFile-server.ts b/common/lib/jsonFile-server.ts @@ -1,6 +1,8 @@ -// ONE JSON READER AND ONE ATOMIC JSON WRITER for the config files and sidecars. -// (Not yet every JSON writer: the slice 4b record lists the ones still on the -// per-pid temp name.) +// ONE JSON READER AND ONE ATOMIC WRITER — for the config files and sidecars, +// and (since release 4 slice W) for every tmp + rename write in common/ and +// editor/: JSON through `writeJsonAtomic`, text and binary through +// `writeFileAtomic`, copies through `copyFileAtomic`. What is left on its own +// temp name is listed in plans/one-core-phase-3.md, "Slice W, as shipped". // // one-core phase 3 slice 4b. Before this module the repo had seven private // copies of `writeJsonAtomic` and some twenty inline `tmp + rename` writes, all @@ -43,7 +45,7 @@ import { randomBytes } from "node:crypto"; import fs from "node:fs"; -import { mkdir, readFile, rename, rm, writeFile } from "node:fs/promises"; +import { copyFile, mkdir, readFile, rename, rm, writeFile } from "node:fs/promises"; import path from "node:path"; export type ReadJsonResult = @@ -134,15 +136,42 @@ export function tmpPathFor(file: string): string { return `${file}.tmp-${process.pid}-${seq}-${randomBytes(4).toString("hex")}`; } -async function writeNow( +export type WriteFileOptions = { + // Create the parent directory first (`mkdir -p`). + mkdir?: boolean; + // Permission bits for the NEW file, applied as the temp file is CREATED + // (`writeFile(tmp, data, { mode })`, under the umask), so the bytes are never + // readable more widely than `mode` — not even for the instant before the + // rename. The X cookie jar writes 0o600 through this. + mode?: number; +}; + +// Run `op` for `file` after every earlier chained operation on the same +// absolute path this process issued has settled. A failed earlier op does not +// block a later one; each caller sees only its own op's error. +function chained(key: string, op: () => Promise<void>): Promise<void> { + const { chains } = jsonFileState(); + const prev = chains.get(key) ?? Promise.resolve(); + const next = prev.then(op, op); + chains.set(key, next); + const release = () => { + if (chains.get(key) === next) chains.delete(key); + }; + next.then(release, release); + return next; +} + +// Put the temp file in place with `fill`, then rename it over `file`; on any +// failure remove the temp and rethrow. +async function replaceNow( file: string, - text: string, makeDir: boolean, + fill: (tmp: string) => Promise<void>, ): Promise<void> { if (makeDir) await mkdir(path.dirname(file), { recursive: true }); const tmp = tmpPathFor(file); try { - await writeFile(tmp, text); + await fill(tmp); await rename(tmp, file); } catch (err) { await rm(tmp, { force: true }).catch(() => {}); @@ -150,30 +179,47 @@ async function writeNow( } } -// Write `value` as JSON to `file` atomically (tmp + rename), after every write -// to the same absolute path this process issued earlier has settled. The value -// is serialised NOW, at the call, so a caller that goes on mutating its object -// cannot change what lands. A failed earlier write does not block a later one; -// each caller sees only its own write's error. +// THE ONE ATOMIC WRITE (tmp + rename) for text and binary files: the playlist, +// failed-transcriptions, the cookie jar, a CHANGELOG, a VTT — and every JSON +// file, through `writeJsonAtomic` below. Same chain, same unique temp name. +// `data` is not copied: a Buffer must not change before the promise settles. +export function writeFileAtomic( + file: string, + data: string | Buffer, + opts: WriteFileOptions = {}, +): Promise<void> { + const key = path.resolve(file); + const { mode } = opts; + return chained(key, () => + replaceNow(key, opts.mkdir === true, (tmp) => + mode === undefined ? writeFile(tmp, data) : writeFile(tmp, data, { mode }), + ), + ); +} + +// Copy `src` to `dest` atomically: copy into a temp beside `dest`, then rename, +// so a crash mid-copy never leaves a partial file under the final name. On the +// same chain as the writers above (keyed on `dest`). +export function copyFileAtomic( + src: string, + dest: string, + opts: { mkdir?: boolean } = {}, +): Promise<void> { + const key = path.resolve(dest); + return chained(key, () => + replaceNow(key, opts.mkdir === true, (tmp) => copyFile(src, tmp)), + ); +} + +// Write `value` as JSON to `file` atomically: `jsonText` → `writeFileAtomic`. +// The value is serialised NOW, at the call, so a caller that goes on mutating +// its object cannot change what lands. export function writeJsonAtomic( file: string, value: unknown, opts: WriteJsonOptions = {}, ): Promise<void> { - const key = path.resolve(file); - const text = jsonText(value, opts); - const { chains } = jsonFileState(); - const prev = chains.get(key) ?? Promise.resolve(); - const next = prev.then( - () => writeNow(key, text, opts.mkdir === true), - () => writeNow(key, text, opts.mkdir === true), - ); - chains.set(key, next); - const release = () => { - if (chains.get(key) === next) chains.delete(key); - }; - next.then(release, release); - return next; + return writeFileAtomic(file, jsonText(value, opts), { mkdir: opts.mkdir }); } // The synchronous twin, for the three stores whose API is synchronous @@ -215,7 +261,8 @@ export function withJsonFileLock<T>(file: string, fn: () => Promise<T>): Promise return next; } -// For tests: how many paths currently have a write in flight. +// For tests: how many paths currently have a write (of any of the three +// kinds) in flight. export function pendingJsonWrites(): number { return jsonFileState().chains.size; } diff --git a/common/lib/savedVideo-server.ts b/common/lib/savedVideo-server.ts @@ -1,7 +1,6 @@ -import { writeJsonAtomic } from "./jsonFile-server"; +import { copyFileAtomic, writeJsonAtomic } from "./jsonFile-server"; import path from "node:path"; import { - copyFile, mkdir, readFile, rename, @@ -76,9 +75,7 @@ async function moveFileCrossDevice(src: string, dest: string): Promise<void> { } catch (err) { if ((err as NodeJS.ErrnoException).code !== "EXDEV") throw err; } - const tmp = `${dest}.tmp-${process.pid}`; - await copyFile(src, tmp); - await rename(tmp, dest); + await copyFileAtomic(src, dest); await rm(src, { force: true }); } diff --git a/common/lib/widgetPresets.ts b/common/lib/widgetPresets.ts @@ -1,5 +1,5 @@ import fs from "node:fs"; -import path from "node:path"; +import { writeJsonAtomic } from "./jsonFile-server"; import type { Paths } from "./paths"; // Persisted monitor-widget presets: named arrangements the /widget/builder board @@ -77,12 +77,7 @@ async function writePresets( presets: WidgetPreset[], ): Promise<void> { const out: PresetsFile = { v: 1, presets }; - await fs.promises.mkdir(path.dirname(paths.widgetPresetsFile), { - recursive: true, - }); - const tmp = `${paths.widgetPresetsFile}.tmp-${process.pid}`; - await fs.promises.writeFile(tmp, JSON.stringify(out, null, 2) + "\n"); - await fs.promises.rename(tmp, paths.widgetPresetsFile); + await writeJsonAtomic(paths.widgetPresetsFile, out, { mkdir: true }); } // Save a preset under `name`. Saving over an existing name OVERWRITES its query diff --git a/common/social/xSessionBroker.ts b/common/social/xSessionBroker.ts @@ -20,7 +20,8 @@ // build-time one. import path from "node:path"; -import { mkdir, readFile, rename, writeFile, rm, stat } from "node:fs/promises"; +import { mkdir, readFile, rm, stat } from "node:fs/promises"; +import { writeFileAtomic } from "../lib/jsonFile-server"; import type { Paths } from "../lib/paths"; import { importPlaywright } from "./playwrightRuntime"; @@ -212,11 +213,12 @@ async function writeCookieJar( log: (line: string) => void, ): Promise<void> { const file = xCookieFile(paths); - await mkdir(path.dirname(file), { recursive: true }); - const body = toNetscapeCookieFile(cookies); - const tmp = `${file}.tmp-${process.pid}`; - await writeFile(tmp, body, { mode: 0o600 }); - await rename(tmp, file); + // mode is applied as the temp is CREATED, before the rename: the jar is + // never world-readable, not even for an instant. + await writeFileAtomic(file, toNetscapeCookieFile(cookies), { + mkdir: true, + mode: 0o600, + }); const authed = cookies.some((c) => c.name === AUTH_COOKIE); log( `Exported ${cookies.length} X cookie(s) to ${file}` + diff --git a/common/ytdlp/runYtdlp.ts b/common/ytdlp/runYtdlp.ts @@ -1,5 +1,6 @@ import path from "node:path"; -import { mkdir, readdir, readFile, rename, writeFile } from "node:fs/promises"; +import { mkdir, readdir, readFile } from "node:fs/promises"; +import { writeFileAtomic } from "../lib/jsonFile-server"; import { execa } from "execa"; import pLimit from "p-limit"; import { readArchive } from "../lib/archive"; @@ -457,9 +458,7 @@ async function writePlaylistFile( onLog: (s: string) => void, ): Promise<void> { const playlistPath = path.join(root, "playlist"); - const tmpPath = `${playlistPath}.tmp-${process.pid}`; - await writeFile(tmpPath, urls.join("\n") + (urls.length ? "\n" : "")); - await rename(tmpPath, playlistPath); + await writeFileAtomic(playlistPath, urls.join("\n") + (urls.length ? "\n" : "")); onLog(`Wrote ${urls.length} URLs to ${playlistPath}\n`); } diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -3,6 +3,7 @@ ## [Unreleased] - **Channel rows no longer scroll over a group's controls on `/channels`.** Scrolled down and to the right, the pinned Slug column of every row painted over the pinned group header and its five station buttons (Sync, Download, Transcribe, Digest and the speaker lane), and took the clicks. The pinned Slug cell and the group header sat at the same stacking level, and the later rows won. The rack now has one named layer order, kept in one file: the Advanced panel, then the column header, then the group header, then the pinned checkbox and Slug cells. Nothing ties any more. The screenshot audit found four more problems, fixed as well. A group header's name and buttons now stay on screen however far the columns scroll across (they used to scroll off to the left). An Advanced panel opened near the bottom or the right edge scrolls itself into view instead of being cut off. The rule above a pinned group header moves with it instead of leaving a gap the rows showed through. On a phone, the column header no longer paints over the selection bar pinned to the bottom of the screen. - **A group's Transcribe works for YouTube channels, and it counts what it queues.** The station used to be disabled for every `youtube`-handling channel with the message "a youtube-handling channel never runs whisper". That was wrong. A YouTube video that came down with no captions is transcription work like any other, and the automatic runner already treats it that way. Transcribe now counts two kinds of video, after the usual members-only, deleted and private exclusions: downloaded videos with no transcript at all, and downloaded videos whose only transcript is YouTube's auto-captions. Pressing it queues exactly those videos, by id, as the channel page does: up to two jobs per channel on the transcription queue. A video downloaded before it went private, members-only or deleted is no longer transcribed by the group button, because it was never in the figure. Pressing it again while either job runs says *already running*. The wording names no method ("…has downloaded audio to transcribe", "…each takes minutes"). **This figure can now be higher than the Transcription band in the same rack on channels with many auto-caption-only videos.** The band counts videos with no transcript at all, while the station counts everything its button would queue. That is intended. +- **The editor's atomic JSON, text and binary writes now go one way, and a failed write no longer leaves a temp file behind.** Nineteen JSON write sites and seven text and binary ones each wrote `<file>.tmp-<pid>` and renamed it over the original — the channel roster, maybe-missing and metadata-scan records, the scheduler and auto-queue state, worker defaults, widget presets, the homepage config, relocation markers, shard configs, the duplicate and media-scan reports and their review decisions, the saved-video backup manifest, both cue normalizers, the playlist, the failed-transcriptions list, the X cookie jar, a site's CHANGELOG cut, a saved video copied into its store across drives, and the video page's VTT promote and remark. They now all go through one writer (`common/lib/jsonFile-server.ts`), which gives every write its own temp name and queues writes to the same file one behind another, so two jobs touching one channel's roster at once cannot trip over each other's temp file. A write that fails now removes its temp: the live `.auto-queue/` holds 175 `state.json.tmp-…` files (173 of them empty) from the day `/home` filled up (2026-09-11), each one a failed write the old code left behind; nothing deletes those old ones for you — `find transcripts -name '*.tmp-*'` lists them. Four temp names stay, on purpose: the export build's two page writers `buildIndex.ts` (a streaming page writer) and `buildStats.ts` (a hand-joined array) — folding them is a restructuring, not a swap — and `transcode.ts` / `transcribeOne.ts` name the output file ffmpeg or the transcription app writes, which is not our write to fold. No file's contents change — every writer puts the same bytes on disk it did before, measured over the live corpus. The cookie jar is still created readable only by you. - **Every channel table and every job-in-flight line is now drawn one way.** The /channels rack, the dashboard's Channels table and the work tables on the operation pages and /cleanup are one table with a column set per page, over one channel row built on the server (which no longer ships a channel's config to the browser); the dashboard's "Needs work" seed is computed by the same code the widget endpoint serves. On the jobs side, /jobs rows, the "Active jobs" cards on channel/video/operation pages, the monitor widget's Active jobs strip and the operations board's "In flight" list are one job row in three sizes, with one rule for which buttons (Retry / Reorder / Drain / Cancel / Force-release) a job gets. **What you might notice:** a work table's report column reads "stale"/"missing" like the rack's instead of a date; the dashboard's Sync button is the rack's; a lane line on /jobs offers Force-release while its runner is running; widget job lines show who asked for the job; an in-flight download on the operations board links to its job page. Nothing a count says moved. - **`site.json`, each channel's `config.json` and the per-video sidecars now have one schema each, and the two config files have generated key tables.** **`SITE.md`** and **`CHANNEL.md`** (new, repo root) list every key with its default and meaning, generated by `common/bin/file-schemas-docs.ts` and checked by a test. Nothing an operator has configured reads or saves differently: every live `site.json` and `config.json`, and a 1,763-file sample of sidecars, read and write back byte-for-byte as before. **Fixed:** a social-channel fetch no longer undoes Configure-form edits made while it was running (it used to write back the whole config it read when it started). Every change to a channel's config now re-reads the file at the moment it saves and changes only its own fields, so a sync stamping its time and a form save made at the same moment both land. Two writes to the same file from the editor no longer share one temporary file. - **`settings.json` has one schema and one writer, and its key table is generated.** Every key, its default, its clamp and its documentation is now one zod schema (`common/lib/settingsSchema.ts`); `getSettings`/`writeSettings` both parse through it, and every settings form saves through one helper (`editor/app/settings/saveSettings.ts`) that merges only what the form changed. **`SETTINGS.md`** (new, repo root) lists every key with its default and what it does, and `settings.json.example` is now the full default object — both generated by `common/bin/settings-example.ts` and checked by a test, so neither can drift. Nothing an operator has configured reads differently. **Fixed:** adding or editing a storage location on `/storage` no longer erases the record of which location the saved-video store is on (`storage.savedVideosLocationId`). diff --git a/editor/app/channels/[slug]/videos/[id]/videoActions.ts b/editor/app/channels/[slug]/videos/[id]/videoActions.ts @@ -1,7 +1,7 @@ "use server"; import path from "node:path"; -import { readdir, readFile, rename, rm, stat, writeFile } from "node:fs/promises"; +import { readdir, readFile, rm, stat } from "node:fs/promises"; import { revalidatePath } from "next/cache"; import { redirect } from "next/navigation"; import type { @@ -71,6 +71,7 @@ import { import { requestChannelSnapshot } from "yt-dlp-transcript-common/jobs/snapshotScheduler"; import { makeTaskTracker } from "yt-dlp-transcript-common/jobs/taskHooks"; import { fixIncompleteTranscriptOne } from "../../lib/fixIncompleteTranscript"; +import { writeFileAtomic, writeJsonAtomic } from "yt-dlp-transcript-common/lib/jsonFile-server"; function videoQueueKey(config: ChannelConfig, override: string | undefined): string { return resolveQueueKey(downloadQueueKey(config), override); @@ -495,10 +496,7 @@ export async function setPrimaryTranscriptAction( return { ok: false, error: `File not found: ${filename}` }; } if (filename !== VTT_FILENAME) { - const dest = path.join(videoDir, VTT_FILENAME); - const tmp = path.join(videoDir, `${VTT_FILENAME}.tmp-${process.pid}`); - await writeFile(tmp, raw); - await rename(tmp, dest); + await writeFileAtomic(path.join(videoDir, VTT_FILENAME), raw); } // The video page reads the dir directly, so it reflects the new primary right // away. The channel list + diagnostics bucket read the cached snapshot and @@ -607,9 +605,8 @@ export async function markVideoUntranscribableAction( } catch { // good — file does not exist } - const tmp = path.join(videoDir, `transcript.tmp-${process.pid}.json`); - await writeFile(tmp, '{"transcription":[]}\n'); - await rename(tmp, transcriptPath); + // Compact + "\n": exactly the literal `{"transcription":[]}\n` this wrote. + await writeJsonAtomic(transcriptPath, { transcription: [] }, { indent: 0 }); const failureListFile = path.join( getPaths().channelsDir, slug, diff --git a/editor/app/sites/lib/cutReleaseAction.ts b/editor/app/sites/lib/cutReleaseAction.ts @@ -9,6 +9,7 @@ import { CutReleaseError, } from "yt-dlp-transcript-common/lib/changelog"; import { commitPath, listDirtyPaths } from "yt-dlp-transcript-common/lib/git"; +import { writeFileAtomic } from "yt-dlp-transcript-common/lib/jsonFile-server"; export type CutReleaseState = | { ok: true; version: string; committed: boolean } @@ -93,9 +94,7 @@ export async function cutReleaseAction( } throw err; } - const tmp = `${filePath}.tmp-${process.pid}`; - await fs.promises.writeFile(tmp, next); - await fs.promises.rename(tmp, filePath); + await writeFileAtomic(filePath, next); if (shouldCommit) { const result = await commitPath( paths.monorepoRoot, diff --git a/plans/one-core-phase-3.md b/plans/one-core-phase-3.md @@ -835,6 +835,125 @@ gave **diff empty, 5,119 bytes each**. - An opened Advanced panel lengthens the region's scroll extent while it is open, because it is absolute inside the scroll box. Closing it restores the extent. +### Slice W, as shipped — one write idiom (2026-09-24) + +Branch `one-core/phase-3-w` off `main` `4130aca1`, five commits, unmerged. Slice 4b left +**16 JSON write sites in 13 files** on the per-pid temp name `${file}.tmp-${process.pid}`, +plus the text and binary tmp + rename writers, plus two modules (`metadataScanStore`, +`autoQueueState`) that had grown their own per-module-copy write counters to dodge the +collision. All of them now go through `common/lib/jsonFile-server.ts`, byte-for-byte. + +| sha | what | +|---|---| +| `db9e3f44` | `writeFileAtomic(file, data: string \| Buffer, {mkdir, mode})` under `writeJsonAtomic` (now `jsonText` → `writeFileAtomic`): the same per-path chain on `globalThis`, the same `${file}.tmp-${pid}-${seq}-${random}` name. `mode` goes to `writeFile(tmp, data, {mode})`, i.e. onto the temp at creation, before the rename — the ordering `xSessionBroker` had. `copyFileAtomic(src, dest, {mkdir})` on the same chain (decided: the saved-video move folds rather than stays). No new sync twin. Tests: string / Buffer / empty bytes, mode 0o600 on the file and on every temp seen, one chain shared by the two writers on one path (40 interleaved writes, last issued lands), copy | +| `9bfd15cd` | Race class: `rosterStore.writeRoster` (mkdir), `maybeMissingStore.writeMaybeMissing` (no mkdir), `metadataScanStore.writeMetadataScan` (no mkdir) on `writeJsonAtomic`. The counters `metadataScanStore.ts:188-191` and `autoQueueState.ts:141` deleted; `autoQueueState` on `writeJsonAtomic` (mkdir). New `autoQueueState.test.ts` "overlapping writes do not collide on the tmp file" (12 concurrent writes, last issued lands, no temp left); `metadataScanStore.test.ts:253` and `rosterStore.test.ts:184-186` unchanged and green | +| `ae6ed9d2` | Process-global + per-video JSON: `syncSchedulerState`, `workerDefaults`, `widgetPresets`, `homepage` (mkdir `homepageDir` = the file's parent), `migrate-channel-priority`, `relocateDir.writeDirMarker`, `shard.saveShardConfig`, `duplicateShorts` ×2, `scanCorruptMedia` ×2, `backupSavedVideos` manifest, `normalizeLiveChat`, `normalizeTranscript`, `videoActions.ts` remark (`'{"transcription":[]}\n'` → `writeJsonAtomic(file, {transcription: []}, {indent: 0})`). Compact without newline (`{indent: 0, newline: false}`) for the two reports and the two cue files. mkdir exactly where a site had one. New `controller/compactJsonWriters.test.ts` (the four compact writers' bytes = `JSON.stringify` of their parse, no newline; passes on the parent too) and the literal pinned in `jsonFile-server.test.ts` | +| `588fd7e9` | Text / binary: `failedTranscriptions` prune + clear, `runYtdlp.writePlaylistFile`, `xSessionBroker.writeCookieJar` (`{mkdir: true, mode: 0o600}`), `savedVideo-server.moveFileCrossDevice` (`copyFileAtomic`), `sites/lib/cutReleaseAction.ts`, `videoActions.ts` VTT promote. One comment on `storageWatch.ts`'s `let timer` (a per-copy singleton, not a temp name, one caller) | +| `be4769de` | `plans/tools/phase3-writers-numbers.ts`, this record, the changelog bullet | +| (review fixes) | record wording (188, writes-not-a-lock, `transcribeOne.ts:173`), changelog scope, the mode test made non-vacuous — see below | + +`videoActions.ts` changed at its two write sites and one added import line only — no export +renamed, no signature changed (slice 3b owns its import list). + +**After commit 4,** `git grep -n 'tmp-${process.pid}' -- common editor`: + +``` +common/controller/buildIndex.ts:971: tmpPath = `${outPath}.tmp-${process.pid}`; +common/controller/buildStats.ts:220: const tmp = `${outPath}.tmp-${process.pid}`; +common/controller/transcode.ts:31: `audio.tmp-${process.pid}.${opts.targetFormat}`, +common/controller/transcribeOne.ts:142: const tmpBase = `transcript.tmp-${process.pid}`; +common/lib/jsonFile-server.test.ts:65: assert.ok(a.startsWith(`/x/config.json.tmp-${process.pid}-`)); +common/lib/jsonFile-server.ts:9:// of them spelling the temp file `${file}.tmp-${process.pid}`. That name is the +common/lib/jsonFile-server.ts:136: return `${file}.tmp-${process.pid}-${seq}-${randomBytes(4).toString("hex")}`; +``` + +The two export page writers, as planned, **plus two the plan did not list**: `transcode.ts:31` +and `transcribeOne.ts:142` are not writers of ours — they name the output file an external +process (ffmpeg; whisper / chough) writes, which the code then renames. Nothing to fold; left +and named. The last three hits are the shared writer itself, its history comment and its test. +`git grep 'rename(tmp'` over `common editor` finds only `buildIndex`, `buildStats`, `transcode`. + +**Out of scope by name:** `buildIndex.ts:971` (the streaming `createWriteStream` page writer) +and `buildStats.ts:220` (a hand-joined array) — restructuring, not a fold; +`scripts/diarize.mjs`; everything under `umtool/`. And one the grep cannot see: +`transcribeOne.ts:173` writes the remote transcript straight to `transcript.json` +(`writeFile(transcriptPath, bytes)` — no temp, no rename), so a crash mid-write can leave it +truncated; it never had the tmp idiom. A candidate for `writeFileAtomic` in a later slice. The module-level `storageWatch` timer and +the TTL caches in `autoRunner.ts` / `recencyIndex.ts` are not temp names. + +**Behaviour changes (intended).** (1) Every folded write is chained per absolute path with every +other `writeFileAtomic`/`writeJsonAtomic`/`copyFileAtomic` in the process, across module +copies — the roster's writers in `runYtdlp`, `quickAvailabilityCheck` and the +`pipelineActions` server action now have their WRITES serialised. That is not a lock: a +`load → merge → writeRoster` from two actors can still lose one merge, exactly as before +(the read-modify-write lock is `withJsonFileLock`, which these callers do not take). (2) A failed write removes its temp; the old +code left it. (3) Temp names changed shape (`<file>.tmp-<pid>-<seq>-<hex>`); the remark's temp +was `transcript.tmp-<pid>.json` and is now `transcript.json.tmp-…`. Nothing reads temp names. + +**Found and left: 188 orphan temps in the live corpus.** `transcripts/.auto-queue/` holds +**175** `state.json.tmp-2514131-NNNN` files, all dated 2026-09-11 — the day `/home` hit 100 % — +**173 of them 0 bytes** (two are partial, 16–20 KB): each a failed `writeFile` (ENOSPC) the old +code never cleaned. Thirteen more elsewhere: seven `snapshot.json.tmp-<pid>`, three sidecar +temps (`availability.json`, `download-outcome.json`) — all from pre-4b writers — and three +`transcript.tmp-<pid>` whisper output bases (`transcribeOne`, an external process's file). The +new writer cannot leave more of the first kind; the existing ones are corpus files and were +**not touched** — the operator's to delete (`find transcripts -name '*.tmp-*'` lists all 188). + +**Numbers** (`plans/tools/phase3-writers-numbers.ts`, one process, never writes the corpus, +never boots a server, clock frozen, only names present on both sides). Inputs frozen once +(`FREEZE_TO`, 1,787 files) and both runs read the frozen tree: the parent `4130aca1` from a +detached scratch worktree, the branch from `588fd7e9` + the tool. Per writer and sample it +loads through the module's reader, writes through the module's writer into scratch, and prints +the md5 of the bytes plus whether they equal the OLD idiom's bytes (`JSON.stringify(v, null, +2) + "\n"`, or `JSON.stringify(v)` for the compact four). Coverage: 53 rosters, 52 +maybe-missing, 1 metadata-scan, scheduler, auto-queue (33,004 B), worker defaults, widget +presets, homepage, the duplicates report (empty scan), the media-scan report (a merge of the +live 27,756 B report), a relocation marker (synthetic — none live), the backup manifest, 100 +`transcript.cues.json` and 100 `live_chat.cues.json` re-normalized. 323 lines each side, +**315 `old-idiom=same`, 0 DIFF, 0 THREW, 0 leftover temps; `diff` before/after empty** +(`w-numbers-{before,after}.txt`). Not measurable without a live sample: shard configs (none +live), duplicate overrides (the live file has no clusters), media-scan overrides (no file), +the priority migration (a CLI whose only write is the default-options `writeJsonAtomic`, +pinned by `jsonText`'s tests) — and the remark literal, pinned by a unit test. + +**Gates** (worktree `one-core-phase-3-w`, ports 3301/3311/3310/3320): +- `pnpm -r --no-bail --workspace-concurrency=1 exec tsc --noEmit` — clean before every commit. +- common **1733/1733** (1723 + 4 `writeFileAtomic`/`copyFileAtomic` + 1 literal + 1 + `autoQueueState` overlap + 4 compact-writer bytes); editor unit (`tsx --test + "app/**/*.test.ts"`) **67/67**; `test:scripts` **156 + 1 skip** (a first run while this + slice's own e2e held the machine lock failed the queue-lock banner test — the lock was + taken; re-run with it free, green); mcp **219/219**. +- `pnpm --filter editor exec next build` exit 0, `ƒ /api/view/[name]` in the table; + `pnpm --filter export exec next build` exit 0; `ZodError|_zod` over both `.next/static`: 0 + files each. +- e2e, the prompt's 20 specs (all exist): **140 passed, 0 failed, 10.7 min**. + +**Review fixes** (review verdict: ship after fixes; four nits, none blocking). The heading's +orphan count is corrected from 173 to 188, the body's total. The record now says the roster +writers' writes are serialised but a load-merge-write is not locked, and lists +`transcribeOne.ts:173` as out of scope. The changelog no longer says "every file": it names +the four temp names that stay. The mode test (`jsonFile-server.test.ts`) could pass vacuously, +because a `fs.watch` could see no temp. It now intercepts `rename` (patched on +`node:fs/promises` + `syncBuiltinESMExports`) and stats the temp at that moment, asserting +exactly one temp, already 0o600. With `mode` removed from `writeFileAtomic` it fails. +Gates after: tsc clean, common **1733/1733**, editor unit **67/67**. No runtime code changed, +so the builds, e2e and numbers were not re-run. + +**Merged with `main` `77a32de2` (slice P) as `44dd843f`**. The merge conflicted only in +`editor/CHANGELOG.md` and this file, and both sides were kept, P then W. Gates on the merged +tip: +- tsc clean. +- common **1738/1738**: W's 1733 plus P's 5. +- editor unit **69/69**: 67 plus P's 2. +- `test:scripts` **156 + 1 skip**, run once the e2e lock was free. +- mcp **219/219**. +- Editor `next build` exit 0, with `ƒ /api/view/[name]` in the route table. Export + `next build` exit 0. No `ZodError|_zod` in either `.next/static`. +- Numbers tool: inputs re-frozen (1,787 files). The before run is `main` `77a32de2` (the old + writers), from a detached scratch worktree; the after run is `44dd843f`. 323 lines each, + 315 `old-idiom=same`, 0 DIFF/THREW/LEFT, **diff empty** (`w-m-numbers-{before,after}.txt`). +- e2e, the same 20 specs, on ports 3411/3410: **140 passed, 0 failed, 8.8 min**. + ## Next release — slice 3b and Phase 4 (inventory kept from 2026-09-23) Slices 3a and 4b shipped in the release above (2026-09-24); the slice 3b bullets and the diff --git a/plans/tools/phase3-writers-numbers.ts b/plans/tools/phase3-writers-numbers.ts @@ -0,0 +1,439 @@ +#!/usr/bin/env tsx +// The one-core Phase 3 release 4 slice W measurement: the bytes every folded +// tmp + rename JSON writer puts on disk, over live samples — printed +// deterministically so a run on the parent commit and a run on the branch can +// be diffed. Slice W moved these writers onto lib/jsonFile-server.ts and +// claims it changed no byte; an empty diff is that claim. +// +// Model: phase3-files-numbers.ts next door, and the same rules. +// +// NEVER WRITES THE CORPUS. Every input is COPIED into a scratch tree under +// os.tmpdir() first; every read-for-measurement and every write happens on the +// copy, and the scratch tree is deleted at the end. The live corpus is only +// ever read (readdir, stat, readFile, copyFile source). +// +// NEVER BOOTS A SERVER. The writers are called in-process; instrumentation.ts +// is not loaded, so nothing is armed. +// +// ONE PROCESS. `TRANSCRIPTS_DIR` is pointed at the scratch tree BEFORE any +// module that memoises `getPaths()` is imported. +// +// THE CLOCK IS FROZEN. Several writers stamp `new Date()` into what they write +// (a decision's `decidedAt`, a report's `generatedAt`); `Date` is replaced with +// a subclass whose "now" is fixed, so two runs write the same bytes. +// +// USES ONLY NAMES PRESENT ON BOTH SIDES of slice W (the modules' exported +// readers and writers, unchanged by the slice), so the SAME file runs on the +// parent and on the branch. +// +// Per writer and sample it prints the md5 of the bytes written and whether they +// equal the OLD idiom's bytes for the same value — `JSON.stringify(v, null, 2) +// + "\n"` or, for the compact four, `JSON.stringify(v)` — where v is the parse +// of what was written. The old-idiom column pins the FORMAT on each side; the +// before/after diff of the md5 column pins the CONTENT. +// +// Usage, from the repo root: +// LIVE_TRANSCRIPTS_DIR=/abs/transcripts node_modules/.bin/tsx plans/tools/phase3-writers-numbers.ts > out.txt +// Default LIVE_TRANSCRIPTS_DIR: the primary checkout's, a sibling of this repo. +// SAMPLE_N (default 100): video dirs sampled per cue file. +// +// FREEZING THE INPUTS. The live editor rewrites the scheduler and auto-queue +// state, rosters and maybe-missing records all day, so a parent run and a +// branch run minutes apart would differ for reasons that are not this code. +// `FREEZE_TO=/abs/dir` copies exactly what a measurement reads, in the same +// layout, into that directory and exits; both runs then take +// `LIVE_TRANSCRIPTS_DIR=/abs/dir`, and a frozen tree samples itself. + +import { createHash } from "node:crypto"; +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; +import { fileURLToPath } from "node:url"; + +const HERE = path.dirname(fileURLToPath(import.meta.url)); +const REPO = path.resolve(HERE, "..", ".."); +const LIVE = + process.env.LIVE_TRANSCRIPTS_DIR ?? + path.join(path.dirname(REPO), "yt-dlp-transcript-browser", "transcripts"); +const SAMPLE_N = Number(process.env.SAMPLE_N ?? 100); +// A live-chat file can be hundreds of MB; the sample takes small ones. +const LIVE_CHAT_MAX_BYTES = 4 * 1024 * 1024; + +const SCRATCH = fs.mkdtempSync(path.join(os.tmpdir(), "phase3-writers-")); +process.env.TRANSCRIPTS_DIR = SCRATCH; +process.env.TZ = "UTC"; + +const FIXED_NOW = Date.UTC(2026, 8, 24, 12, 0, 0); +const RealDate = Date; +class FixedDate extends RealDate { + constructor(...args: unknown[]) { + if (args.length === 0) super(FIXED_NOW); + else super(...(args as [number])); + } + static now(): number { + return FIXED_NOW; + } +} +globalThis.Date = FixedDate as DateConstructor; + +type Format = "pretty" | "compact"; + +function md5(b: string | Buffer): string { + return createHash("md5").update(b).digest("hex"); +} + +function oldIdiom(v: unknown, format: Format): string { + return format === "compact" ? JSON.stringify(v) : JSON.stringify(v, null, 2) + "\n"; +} + +function sortedDir(dir: string): string[] { + try { + return fs.readdirSync(dir).sort(); + } catch { + return []; + } +} + +// Copy a live file (path relative to the corpus root) into the scratch tree. +function stage(rel: string): boolean { + const src = path.join(LIVE, rel); + if (!fs.existsSync(src)) return false; + const dest = path.join(SCRATCH, rel); + fs.mkdirSync(path.dirname(dest), { recursive: true }); + fs.copyFileSync(src, dest); + return true; +} + +// Report the file a writer just wrote. +function report(label: string, file: string, format: Format): void { + if (!fs.existsSync(file)) { + console.log(`${label} ABSENT`); + return; + } + const bytes = fs.readFileSync(file); + let idiom: string; + try { + idiom = oldIdiom(JSON.parse(bytes.toString("utf8")), format) === bytes.toString("utf8") + ? "old-idiom=same" + : "old-idiom=DIFF"; + } catch { + idiom = "old-idiom=UNPARSEABLE"; + } + const leftovers = sortedDir(path.dirname(file)).filter((n) => + n.startsWith(`${path.basename(file)}.tmp-`), + ); + console.log( + `${label} ${bytes.length}B md5=${md5(bytes)} ${idiom}` + + (leftovers.length ? ` LEFT ${leftovers.length} temp(s)` : ""), + ); +} + +async function run(label: string, fn: () => Promise<void>): Promise<void> { + try { + await fn(); + } catch (e) { + console.log(`${label} THREW: ${(e as Error).message}`); + } +} + +async function channelStores(slugs: string[]): Promise<void> { + const { getPaths } = await import("../../common/lib/paths"); + const roster = await import("../../common/controller/rosterStore"); + const mm = await import("../../common/controller/maybeMissingStore"); + const ms = await import("../../common/controller/metadataScanStore"); + const shard = await import("../../common/controller/shard"); + const paths = getPaths(); + + console.log("# per-channel stores (every live file)"); + for (const slug of slugs) { + const ch = path.join("channels", slug); + if (stage(path.join(ch, roster.ROSTER_FILENAME))) { + await run(`roster ${slug}`, async () => { + await roster.writeRoster(paths, slug, await roster.loadRoster(paths, slug)); + report(`roster ${slug}`, roster.rosterPath(paths, slug), "pretty"); + }); + } + if (stage(path.join(ch, mm.MAYBE_MISSING_FILENAME))) { + await run(`maybe-missing ${slug}`, async () => { + const rec = await mm.loadMaybeMissing(paths, slug); + if (!rec) return void console.log(`maybe-missing ${slug} load=null`); + await mm.writeMaybeMissing(paths, slug, rec); + report( + `maybe-missing ${slug}`, + path.join(paths.channelsDir, slug, mm.MAYBE_MISSING_FILENAME), + "pretty", + ); + }); + } + if (stage(path.join(ch, ms.METADATA_SCAN_FILENAME))) { + await run(`metadata-scan ${slug}`, async () => { + // The writer is private; an upsert that changes one entry's title by a + // fixed suffix forces exactly one write of the whole store. + const scan = await ms.loadMetadataScan(paths, slug); + const id = Object.keys(scan.entries).sort()[0]; + const upsert = id + ? { entries: { [id]: { ...scan.entries[id], title: `${scan.entries[id].title}·` } } } + : {}; + await ms.upsertMetadataScan(paths, slug, upsert, "2026-09-24T12:00:00.000Z"); + report(`metadata-scan ${slug}`, ms.metadataScanPath(paths, slug), "pretty"); + }); + } + for (const op of shard.SHARD_OPS) { + const rel = path.relative(SCRATCH, shard.shardFile(paths, slug, op)); + if (!stage(rel)) continue; + await run(`shard-${op} ${slug}`, async () => { + const cfg = await shard.loadShardConfig(paths, slug, op); + if (!cfg) return void console.log(`shard-${op} ${slug} load=null`); + await shard.saveShardConfig(paths, slug, op, cfg); + report(`shard-${op} ${slug}`, shard.shardFile(paths, slug, op), "pretty"); + }); + } + } +} + +async function globalStores(): Promise<void> { + const { getPaths } = await import("../../common/lib/paths"); + const paths = getPaths(); + const rel = (abs: string) => path.relative(SCRATCH, abs); + console.log(""); + console.log("# process-global stores"); + + const sched = await import("../../common/jobs/syncSchedulerState"); + if (stage(rel(paths.schedulerStateFile))) { + await run("scheduler", async () => { + await sched.writeSchedulerState(paths, await sched.readSchedulerState(paths)); + report("scheduler", paths.schedulerStateFile, "pretty"); + }); + } + const aq = await import("../../common/jobs/autoQueueState"); + if (stage(rel(paths.autoQueueStateFile))) { + await run("auto-queue", async () => { + await aq.writeAutoQueueState(paths, await aq.readAutoQueueState(paths)); + report("auto-queue", paths.autoQueueStateFile, "pretty"); + }); + } + const wd = await import("../../common/jobs/workerDefaults"); + if (stage(rel(paths.workerDefaultsFile))) { + await run("worker-defaults", async () => { + const d = wd.readWorkerDefaults(paths); + if (!d) return void console.log("worker-defaults load=null"); + await wd.writeWorkerDefaults(paths, d.enabledWorkerIds); + report("worker-defaults", paths.workerDefaultsFile, "pretty"); + }); + } + const wp = await import("../../common/lib/widgetPresets"); + if (stage(rel(paths.widgetPresetsFile))) { + await run("widget-presets", async () => { + const list = await wp.readWidgetPresets(paths); + if (list.length === 0) return void console.log("widget-presets empty"); + // Saving over an existing name rewrites the whole list, unchanged. + await wp.addWidgetPreset(paths, list[0].name, list[0].query); + report("widget-presets", paths.widgetPresetsFile, "pretty"); + }); + } + const hp = await import("../../common/lib/homepage"); + if (stage(rel(paths.homepageConfigFile))) { + await run("homepage", async () => { + await hp.writeHomepageConfig(hp.getHomepageConfig(paths), paths); + report("homepage", paths.homepageConfigFile, "pretty"); + }); + } + + const dup = await import("../../common/controller/duplicateShorts"); + if (stage(rel(dup.duplicateOverridesPath(paths)))) { + await run("duplicate-overrides", async () => { + const cur = await dup.readDuplicateOverrides(paths); + const id = Object.keys(cur.clusters).sort()[0]; + if (!id) return void console.log("duplicate-overrides empty"); + const c = cur.clusters[id] as Record<string, unknown>; + await dup.updateDuplicateOverride(paths, id, { + ...(typeof c.canonicalSlug === "string" ? { canonicalSlug: c.canonicalSlug } : {}), + ...(typeof c.notDuplicate === "boolean" ? { notDuplicate: c.notDuplicate } : {}), + ...(typeof c.confirmed === "boolean" ? { confirmed: c.confirmed } : {}), + ...(typeof c.note === "string" ? { note: c.note } : {}), + }); + report("duplicate-overrides", dup.duplicateOverridesPath(paths), "pretty"); + }); + } + // The report writer runs at the end of a detection pass; over the scratch + // tree (configs and stores, no data/) that pass sees no videos. + await run("duplicates-report", async () => { + await dup.detectDuplicateShorts({ paths, onLog: () => {} }); + report("duplicates-report (no videos)", path.join(paths.transcriptsDir, "duplicates.json"), "compact"); + }); + + const scm = await import("../../common/controller/scanCorruptMedia"); + if (stage(rel(scm.mediaScanOverridesPath(paths)))) { + await run("media-scan-overrides", async () => { + const cur = await scm.readMediaScanOverrides(paths); + const key = Object.keys(cur.reviewed).sort()[0]; + if (!key) return void console.log("media-scan-overrides empty"); + const note = (cur.reviewed[key] as { note?: string }).note; + await scm.updateMediaScanOverride(paths, key, { reviewed: true, ...(note ? { note } : {}) }); + report("media-scan-overrides", scm.mediaScanOverridesPath(paths), "pretty"); + }); + } + // A scan of a channel that does not exist merges the live report's every + // finding back unchanged, through the report writer. + if (stage("media-scan.json")) { + await run("media-scan-report", async () => { + await scm.scanCorruptMedia({ paths, channels: ["__phase3-writers-none__"], onLog: () => {} }); + report("media-scan-report (merge of live)", path.join(paths.transcriptsDir, "media-scan.json"), "compact"); + }); + } + + const rd = await import("../../common/controller/relocateDir"); + // No live marker exists outside a move; a fixed one exercises the writer. + await run("relocation-marker", async () => { + const file = path.join(SCRATCH, "channels", ".phase3-writers-marker.json"); + fs.mkdirSync(path.dirname(file), { recursive: true }); + await rd.writeDirMarker(file, { + target: "/mnt/x/slug/data", + direction: "out", + startedAt: "2026-09-24T12:00:00.000Z", + phase: "copy", + }); + report("relocation-marker (synthetic)", file, "pretty"); + }); + + const bk = await import("../../common/controller/backupSavedVideos"); + await run("backup-manifest", async () => { + const dest = path.join(SCRATCH, ".phase3-writers-backup"); + const r = await bk.backupSavedVideos({ paths, dest, onLog: () => {} }); + report(`backup-manifest (${r.entries} saved)`, r.manifestPath, "pretty"); + }); +} + +// The first SAMPLE_N video dirs (sorted slugs × sorted ids) holding `filename`. +function sampleDirs(slugs: string[], filename: string, maxBytes?: number): string[] { + const out: string[] = []; + for (const slug of slugs) { + const data = path.join(LIVE, "channels", slug, "data"); + for (const id of sortedDir(data)) { + if (out.length >= SAMPLE_N) return out; + const f = path.join(data, id, filename); + try { + const st = fs.statSync(f); + if (maxBytes !== undefined && st.size > maxBytes) continue; + out.push(path.join(data, id)); + } catch { + // not in this dir + } + } + } + return out; +} + +const MEDIA = /\.(opus|m4a|mp3|webm|mp4|mkv|wav|ogg|aac|flac|part|ytdl|jpg|jpeg|png|webp)$/i; + +// Copy a video dir's non-media files (metadata, raw transcripts, the prior cue +// file whose format tag a re-normalize reads) into the scratch tree. +function stageVideoDir(live: string, n: number): string { + const dir = path.join(SCRATCH, "videos", String(n)); + fs.mkdirSync(dir, { recursive: true }); + for (const name of sortedDir(live)) { + if (MEDIA.test(name)) continue; + const src = path.join(live, name); + if (!fs.statSync(src).isFile()) continue; + fs.copyFileSync(src, path.join(dir, name)); + } + return dir; +} + +async function cueWriters(slugs: string[]): Promise<void> { + const nt = await import("../../common/controller/normalizeTranscript"); + const nl = await import("../../common/controller/normalizeLiveChat"); + const vs = await import("../../common/lib/videoStatus"); + let n = 0; + console.log(""); + console.log(`# ${vs.CUES_JSON_FILENAME} (first ${SAMPLE_N} dirs holding one, force re-normalize)`); + for (const live of sampleDirs(slugs, vs.CUES_JSON_FILENAME)) { + const rel = path.relative(path.join(LIVE, "channels"), live); + const dir = stageVideoDir(live, n++); + await run(rel, async () => { + const out = await nt.normalizeTranscript({ videoDir: dir, channelSlug: "x", force: true } as never); + if ((out as { status: string }).status !== "wrote") { + return void console.log(`${rel} ${(out as { status: string }).status}`); + } + report(rel, path.join(dir, vs.CUES_JSON_FILENAME), "compact"); + }); + } + console.log(""); + console.log(`# ${vs.LIVE_CHAT_CUES_FILENAME} (first ${SAMPLE_N} dirs with a live chat ≤ 4 MB)`); + for (const live of sampleDirs(slugs, vs.LIVE_CHAT_FILENAME, LIVE_CHAT_MAX_BYTES)) { + const rel = path.relative(path.join(LIVE, "channels"), live); + const dir = stageVideoDir(live, n++); + await run(rel, async () => { + const out = await nl.normalizeLiveChat({ videoDir: dir, channelSlug: "x", force: true }); + if (out.status !== "wrote") return void console.log(`${rel} ${out.status}`); + report(rel, path.join(dir, vs.LIVE_CHAT_CUES_FILENAME), "compact"); + }); + } +} + +// The files the measurement reads, relative to the corpus root. +async function inputs(slugs: string[]): Promise<string[]> { + const { getPaths } = await import("../../common/lib/paths"); + const shard = await import("../../common/controller/shard"); + const dup = await import("../../common/controller/duplicateShorts"); + const scm = await import("../../common/controller/scanCorruptMedia"); + const vs = await import("../../common/lib/videoStatus"); + const paths = getPaths(); + const rel = (abs: string) => path.relative(SCRATCH, abs); + const out: string[] = []; + for (const s of slugs) { + const ch = path.join("channels", s); + out.push( + path.join(ch, "config.json"), + path.join(ch, "roster.json"), + path.join(ch, "maybe-missing.json"), + path.join(ch, "metadata-scan.json"), + ...shard.SHARD_OPS.map((op) => rel(shard.shardFile(paths, s, op))), + ); + } + out.push( + rel(paths.schedulerStateFile), + rel(paths.autoQueueStateFile), + rel(paths.workerDefaultsFile), + rel(paths.widgetPresetsFile), + rel(paths.homepageConfigFile), + rel(dup.duplicateOverridesPath(paths)), + rel(scm.mediaScanOverridesPath(paths)), + "media-scan.json", + ); + const dirs = [ + ...sampleDirs(slugs, vs.CUES_JSON_FILENAME), + ...sampleDirs(slugs, vs.LIVE_CHAT_FILENAME, LIVE_CHAT_MAX_BYTES), + ]; + for (const d of dirs) { + for (const name of sortedDir(d)) { + if (MEDIA.test(name) || !fs.statSync(path.join(d, name)).isFile()) continue; + out.push(path.relative(LIVE, path.join(d, name))); + } + } + return [...new Set(out)].filter((r) => fs.existsSync(path.join(LIVE, r))); +} + +try { + const slugs = sortedDir(path.join(LIVE, "channels")).filter((s) => + fs.existsSync(path.join(LIVE, "channels", s, "config.json")), + ); + if (process.env.FREEZE_TO) { + const to = path.resolve(process.env.FREEZE_TO); + const files = await inputs(slugs); + for (const r of files) { + fs.mkdirSync(path.dirname(path.join(to, r)), { recursive: true }); + fs.copyFileSync(path.join(LIVE, r), path.join(to, r)); + } + console.error(`froze ${files.length} files into ${to}`); + fs.rmSync(SCRATCH, { recursive: true, force: true }); + process.exit(0); + } + // Configs first: several writers read the channel list. + for (const s of slugs) stage(path.join("channels", s, "config.json")); + await channelStores(slugs); + await globalStores(); + await cueWriters(slugs); +} finally { + fs.rmSync(SCRATCH, { recursive: true, force: true }); +}