commit 9fdd996f96979e768ba3761405c68e9e6045585e
parent af39d6ccf573b389477c73660efef5c3401fa268
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Mon, 7 Sep 2026 22:00:44 -0400
common: a refreshed run keeps the lane's whole ledger, and two lanes never share a tmp file
Slice 1.2 review fix-up. Three things, one of them a hole in the fix that
carried the accumulators across a context refresh.
**The refresh copied two fields of the ledger, not the object.** Its
wait-for-idle check runs BEFORE the await that opens the new run, so it cannot
close the window it looks like it closes: `limit()` keeps reading the old run
while the new one is being built, and `next()` can dispatch one more unit
against it. That unit then settles afterwards — releasing its llm or unit
lease, and reporting its own `costUsd` — against the run it was dispatched
with. With only `costUsd` and `diskFloorHit` copied forward, the release and
the spend land on a run `laneLimit` no longer reads: a lease under-count that
never recovers, and a unit's spend dropped from the metered cap. Reachable on a
`maxWorkers: 1` digest lane whenever a unit finishes with the TTL lapsed.
So the whole `live` object is carried across. Every field on it is per-RUN and
a runner's run is its lifetime; there was never a reason for the copy to be
field-by-field, only an assumption that the two counters would be zero at the
moment of the swap. A unit test pins it: dispatch on the old run, settle after
the swap, and assert the limit reading the NEW run sees both the released lease
and the spend — the spend-cap arm firing is the visible half.
**Two lanes could rename over each other's tmp file.** `writeAutoQueueState`
built `${file}.tmp-${pid}` — one name per PROCESS — and `persist()` is
fire-and-forget, so two runners persisting within a few milliseconds could have
the second truncate the file the first was still writing and then rename it
over the real one. The result reads as empty state, which is tolerated, and
loses a persisted platform cooldown across a restart, which is not.
Pre-existing; two fast lanes made it reachable. A per-write counter is enough —
this is one process, and the pid still separates two of them.
**`resetChannelSnapshotMemo` is wired in rather than left exported and
uncalled.** `/api/test/invalidate-cache` clears it beside the job registry, the
worker pool, the scheduler and the auto-runner. The memo is self-invalidating
on (mtime, size) in production, but `resetData()` rewrites the SAME fixture
paths from the same source tree, and two specs writing an identical snapshot
inside one filesystem timestamp tick would hand the second one the first's
parsed object. Cleared with the other process-wide caches rather than left to a
coincidence.
The as-shipped note records this commit, and two things review asked to have
named: `backfillLaneLimit` returns 0 for an empty operations list where the old
`Math.min(...)` returned `Infinity`, and `runOperationBatch`'s digest path uses
`run.operations[0]`, correct while the digest lane carries one operation.
common 897 tests (896 + the shared-ledger case), tsc clean in common and
editor, numbers diff against the after-1.2 dump empty. No full e2e — 1.3's run
covers it.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Diffstat:
5 files changed, 135 insertions(+), 12 deletions(-)
diff --git a/common/controller/autoRunner.ts b/common/controller/autoRunner.ts
@@ -885,14 +885,21 @@ async function runLoop(
...(clusterPlan ? { clusterPlan } : {}),
});
clusterPlan = opened.digest?.clusterPlan ?? clusterPlan;
- // THE PER-RUN ACCUMULATORS SURVIVE THE REFRESH, and they have to: the
- // metered spend cap is a per-RUN ceiling, and a runner's run is its whole
- // lifetime — resetting the total every sixty seconds would turn a $5 cap
- // into $5 a minute. The disk-floor latch is the same shape.
- if (laneRun.run) {
- opened.live.costUsd = laneRun.run.live.costUsd;
- opened.live.diskFloorHit = laneRun.run.live.diskFloorHit;
- }
+ // THE WHOLE LEDGER SURVIVES THE REFRESH — the OBJECT, not a copy of two
+ // of its fields.
+ //
+ // Every field on it is per-RUN, and a runner's run is its whole lifetime:
+ // the metered spend cap is a per-run ceiling (resetting it every sixty
+ // seconds would turn a $5 cap into $5 a minute) and the disk-floor latch
+ // is the same shape. The two COUNTERS are why this has to be the object.
+ // The idle check above cannot close the window: it runs before the await
+ // below, so `limit()` keeps reading the old run while this one is being
+ // opened and `next()` may dispatch one more unit against it. That unit's
+ // `llmActive--` / `unitActive--` and its own `costUsd` land on the run it
+ // was dispatched with — and if that run's ledger is not this one's, the
+ // lane under-counts its leases forever and drops that unit's spend from
+ // the cap.
+ if (laneRun.run) opened.live = laneRun.run.live;
// The engine fail-fast, asked ONCE per context rather than per video: an
// unreachable ollama would otherwise produce one failure per candidate
// over 55,956 of them. A refusal is an IDLE, not a stop — it is a
diff --git a/common/controller/operationBatch.test.ts b/common/controller/operationBatch.test.ts
@@ -1,6 +1,12 @@
import { test } from "node:test";
import assert from "node:assert/strict";
-import { backfillLimit, candidateAction, laneLimit } from "./operationBatch";
+import {
+ backfillLimit,
+ candidateAction,
+ laneLimit,
+ operationLaneLive,
+ type OperationRun,
+} from "./operationBatch";
import {
defaultBackfill,
defaultDigest,
@@ -437,3 +443,59 @@ function fakeOperation(id: string, contendsFor: "cpu" | "gpu" | "network") {
lane: { queueKey: "backfill", contendsFor },
} as unknown as Operation;
}
+
+// ---------------------------------------------------------------------------
+// The ledger a refreshed run keeps.
+//
+// A long-lived runner re-opens its OperationRun when the operation settings
+// change or its TTL lapses, and the idle check that guards that swap cannot
+// close the window: it runs before the await that opens the new run, so one
+// more unit can be dispatched against the OLD one meanwhile. That unit settles
+// afterwards — releasing its lease and reporting its cost — against the run it
+// was dispatched with. So the refresh must carry the whole ledger OBJECT, not
+// copies of the fields that happened to look like state.
+
+function fakeDigestRun(live: OperationRun["live"]): OperationRun {
+ return {
+ lane: "digest",
+ live,
+ digest: {
+ laneChoice: "local",
+ engine: {
+ app: { lane: "remote-api", metered: true, fields: { baseUrl: true } },
+ },
+ },
+ } as unknown as OperationRun;
+}
+
+test("a refreshed run keeps the lane's whole ledger, leases and spend alike", () => {
+ const ledger = { llmActive: 0, unitActive: 0, costUsd: 0, diskFloorHit: false };
+ const before = fakeDigestRun(ledger);
+ // A unit is dispatched on `before` and takes a remote lease.
+ before.live.llmActive += 1;
+
+ // The refresh installs a NEW run carrying the SAME ledger — what
+ // autoRunner's refreshLaneRun does with `opened.live = laneRun.run.live`.
+ const after = fakeDigestRun(before.live);
+ assert.equal(after.live, before.live, "the ledger is shared, not copied");
+
+ // The in-flight unit now settles against the run it was dispatched with.
+ before.live.llmActive -= 1;
+ before.live.costUsd += 6;
+
+ // And the limit, which reads the NEW run, sees both halves. Without the
+ // shared object it would see llmActive stuck at 0 with a lease outstanding
+ // (an under-count that never recovers) and costUsd at 0 (a spend cap the
+ // metered lane silently stops enforcing).
+ const liveNow = operationLaneLive(after, 1);
+ assert.equal(liveNow.lane, "digest");
+ assert.equal(liveNow.lane === "digest" ? liveNow.llmActive : -1, 0);
+ assert.equal(liveNow.lane === "digest" ? liveNow.costUsd : -1, 6);
+
+ const verdict = laneLimit(
+ settingsWith({ digest: { ...defaultDigest(), spendCapUsd: 5 } }),
+ liveNow,
+ );
+ assert.equal(verdict.limit, 0);
+ assert.equal(verdict.hold?.reason, "spend-cap");
+});
diff --git a/common/jobs/autoQueueState.ts b/common/jobs/autoQueueState.ts
@@ -131,6 +131,15 @@ export async function readAutoQueueState(paths: Paths): Promise<AutoQueueState>
// Write the state atomically (tmp + rename), creating the .auto-queue dir on
// first use. Each kind's pick log is trimmed to the limit on the way out.
+// UNIQUE PER WRITE, not per process. `persist()` is fire-and-forget and two
+// lanes' runners can persist within a few milliseconds of each other, so one
+// 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;
+
export async function writeAutoQueueState(
paths: Paths,
state: AutoQueueState,
@@ -144,7 +153,7 @@ export async function writeAutoQueueState(
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}`;
+ const tmp = `${paths.autoQueueStateFile}.tmp-${process.pid}-${++writeSeq}`;
await writeFile(tmp, JSON.stringify(out, null, 2) + "\n");
await rename(tmp, paths.autoQueueStateFile);
}
diff --git a/editor/app/api/test/invalidate-cache/route.ts b/editor/app/api/test/invalidate-cache/route.ts
@@ -1,6 +1,7 @@
import { NextResponse } from "next/server";
import { revalidatePath } from "next/cache";
import { resetSnapshotScheduler } from "yt-dlp-transcript-common/jobs/snapshotScheduler";
+import { resetChannelSnapshotMemo } from "yt-dlp-transcript-common/controller/channels";
export const dynamic = "force-dynamic";
@@ -80,6 +81,13 @@ function invalidate() {
// next iteration; clearing this lets the next spec start fresh runners.
// eslint-disable-next-line @typescript-eslint/no-explicit-any
(globalThis as any).__yttAutoRunner__ = undefined;
+ // Drop the shared snapshot parse memo. It is keyed on (mtime, size), so it is
+ // self-invalidating in production — but resetData() rewrites the SAME fixture
+ // paths from the same source tree, and two specs writing an identical
+ // snapshot inside one filesystem timestamp tick would hand the second one the
+ // first's parsed object. Cleared here with the other process-wide caches
+ // rather than left to a coincidence.
+ resetChannelSnapshotMemo();
revalidatePath("/", "layout");
return NextResponse.json({ ok: true });
}
diff --git a/plans/one-core-phase-1.md b/plans/one-core-phase-1.md
@@ -333,8 +333,9 @@ three-second poll, so widen the projection in `run()`'s loop rather than in
`83954f1` (one executor) → `596ddb6` (the runner runs operations) → `8c5a49d`
(the console per lane) → `ddbdb30` (the e2e spec) → `1f531fb` (dead half-rules
removed) → `24a2069` (three fixes to the long-lived half) → `9a16d65` (two
-runners, one state object — the bug the e2e spec found), on `one-core/phase-1`
-off `3ced7dd`. **Numbers: the before/after diff of `plans/tools/phase1-numbers.ts`
+runners, one state object — the bug the e2e spec found) → the commit after
+`3bdf5ee` (the review fix-up, which carries this paragraph and so cannot name
+its own sha), on `one-core/phase-1` off `3ced7dd`. **Numbers: the before/after diff of `plans/tools/phase1-numbers.ts`
over the live corpus is EMPTY.** It stays empty because the live `settings.json`
carries no `autoQueue.digest` or `.backfill` block, and the script's lane
sections come from the KEYS PRESENT IN THE FILE — which is exactly why it was
@@ -452,6 +453,32 @@ version that shipped compares the two timestamped pick logs AFTER the work, and
it failed the same way in isolation, which is what turned "flaky test" into
"empty pick log for a lane that is demonstrably working".
+### The review fix-up, the commit after `3bdf5ee`
+
+Three things review found after the e2e run, one of them a hole in the fix
+above. **`refreshLaneRun` copied two fields of the ledger, not the object.** Its
+idle check runs BEFORE the await that opens the new run, so `limit()` keeps
+reading the old run while the new one is being built and `next()` can dispatch
+one more unit against it — reachable on a `maxWorkers: 1` digest lane whenever
+a unit finishes with the TTL lapsed. That unit's `llmActive--` and its own
+`costUsd` then land on a run `laneLimit` no longer reads: a lease under-count
+that never recovers, and a unit's spend dropped from the cap. The whole `live`
+object is carried across now (`opened.live = laneRun.run.live`) — every field on
+it is per-run, and a runner's run is its lifetime — with a unit test that
+dispatches on the old run, settles after the swap, and asserts the limit reading
+the new run sees both the released lease and the spend.
+
+Also: `writeAutoQueueState`'s tmp path was `${file}.tmp-${pid}`, one per
+PROCESS, and `persist()` is fire-and-forget — so two lanes persisting within a
+few milliseconds could rename over each other's half-written file, leaving a
+`state.json` that reads as empty (tolerated) and loses a persisted platform
+cooldown across a restart (not). Pre-existing; two fast lanes made it reachable.
+It carries a per-write counter now. And `resetChannelSnapshotMemo` is wired into
+`/api/test/invalidate-cache` beside the other process-wide caches rather than
+left exported and uncalled: the memo is self-invalidating on (mtime, size) in
+production, but `resetData()` rewrites the same fixture paths and two specs
+writing an identical snapshot inside one filesystem tick would share a parse.
+
### Divergences from the 1.2 bullets, and why
- **One batch function, not two.** The plan kept `runDigestBatch` /
@@ -515,6 +542,16 @@ it failed the same way in isolation, which is what turned "flaky test" into
verbs must work per channel.
- **The concurrency spec measures overlapping pick-log INTERVALS**, not
simultaneous in-flight counts. See the bug note above for why.
+- **`backfillLaneLimit` returns 0 for an empty operations list** where the old
+ `Math.min(...kinds.map(...))` returned `Infinity`. Unreachable before (the
+ batch returned early on no kinds) and unreachable now (the batch still does),
+ but the limit is a shared function called by a runner too, and an unbounded
+ limit for a lane with nothing to dispatch is not a defensible answer to have
+ sitting there. Named because it is a silent behaviour fix, not a rewrite.
+- **`runOperationBatch`'s digest path uses `run.operations[0]`.** Correct while
+ the digest lane carries exactly one operation, which its queue keys guarantee
+ today; a second digest-queue operation would need the candidate list to carry
+ its operation the way the backfill path's `flatMap` does.
- **`backfill.spec.ts`'s concurrency test is kept as well as re-asserted.** The
plan says the `backfill.spec.ts:615` invariant "moves, does not vanish". At
the runner the mechanism is different — two loops on queueKey `""` dispatching