commit 75ed0f0ec2260aba9c586216ce1eae68a1948c51
parent d9a69726a8532bdb47151d10e2ece13de239a8d5
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Thu, 24 Sep 2026 18:53:07 -0400
common: roster, maybe-missing and metadata-scan on writeJsonAtomic; two counters gone
The three per-channel stores several actors write — the sync, the quick
availability check, the pipeline's server action in another bundle — now
share the one per-path chain on globalThis. metadataScanStore's and
autoQueueState's own write counters (one per module copy) are deleted:
the shared temp name makes them redundant. Bytes unchanged (null, 2 +
"\n"); roster and auto-queue keep their mkdir, maybe-missing and the scan
keep none. New: the overlapping-writes test for autoQueueState.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Diffstat:
5 files changed, 54 insertions(+), 40 deletions(-)
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/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/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 {