commit 66f21c97b8e5f542173d185c99d5b2a6320956ae
parent ed9c10c00e135db5579540babdeaee60fcb9626e
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Mon, 7 Sep 2026 21:21:08 -0400
common: two runners, one state file — so one state object
`writeAutoQueueState` serializes the WHOLE file, all four lanes, and every
runner held its own copy read at startup. So every persist wrote back that
copy's idea of the other lanes: whichever runner dispatched most recently
erased the others' pick log and fairness memory.
It was nearly invisible with one slow lane — auto-transcribe persists minutes
apart, auto-download rarely faster — and stopped being invisible the moment two
fast lanes ran together. The digest lane would dispatch a unit, write its pick,
and have it overwritten by the backfill runner's next persist a few hundred
milliseconds later; `/api/auto-queue/status` reads the FILE, so the console
showed a lane with zero picks and `no-pending` while its sidecars were being
written on disk in front of you.
The runners share one object now, on the auto-runner singleton — whose lifetime
is already the right one, since the e2e harness clears it between specs, so a
reset corpus cannot inherit a previous spec's picks. `computeLeafPending` and
the download lane's cooldown merge still read the file, deliberately: the first
wants a value it can clone without touching live fairness, the second is there
precisely to pick up what a manual sync wrote from outside the runner.
**Found by the new lane-runner concurrency spec, which is the answer to "what
was that test for".** It asked for the two lanes' pick logs to overlap in time
and got an empty log for the digest lane; the instrumented form said
`picks=0 running=true idle=no-pending` for a lane whose ai-digest.json files
were already on disk. That spec is also rewritten here: it compares the two
timestamped pick logs AFTER the work, rather than sampling for "both in flight
at one instant" — the same property measured by a coin toss, which a loaded
machine loses to sampling rather than to a regression.
lane-runner.spec 4 passed / 0 failed in 39.0 s (the concurrency case 4.7 s,
against a 120 s timeout it used to exhaust). tsc clean in six packages, common
896, mcp 205, numbers diff empty.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Diffstat:
2 files changed, 94 insertions(+), 20 deletions(-)
diff --git a/common/controller/autoRunner.ts b/common/controller/autoRunner.ts
@@ -54,6 +54,7 @@ export type RecencyKeyView = RecencyKey;
import { LANES } from "../lib/autoQueueTypes";
import {
type AutoQueueKind,
+ type AutoQueueState,
readAutoQueueState,
recordPick,
writeAutoQueueState,
@@ -198,7 +199,13 @@ type RunnerLive = {
idleReason: AutoRunnerIdleReason | null;
};
-type AutoRunnerSingleton = { runners: Map<AutoQueueKind, RunnerLive> };
+type AutoRunnerSingleton = {
+ runners: Map<AutoQueueKind, RunnerLive>;
+ // ONE persisted state object, shared by every lane's runner in this process.
+ // See sharedAutoQueueState.
+ state: AutoQueueState | null;
+ stateFile: string | null;
+};
declare global {
// eslint-disable-next-line no-var
@@ -207,11 +214,42 @@ declare global {
function getSingleton(): AutoRunnerSingleton {
if (!globalThis.__yttAutoRunner__) {
- globalThis.__yttAutoRunner__ = { runners: new Map() };
+ globalThis.__yttAutoRunner__ = {
+ runners: new Map(),
+ state: null,
+ stateFile: null,
+ };
}
return globalThis.__yttAutoRunner__;
}
+// THE PERSISTED STATE IS ONE OBJECT FOR EVERY LANE'S RUNNER, and it has to be.
+//
+// `writeAutoQueueState` serializes the WHOLE file — all four lanes — so two
+// runners each holding their own copy means every persist clobbers the other
+// lane's pick log and fairness memory with whatever that copy was read with.
+// With one slow lane that is nearly invisible (a transcription is minutes
+// apart), and it became obvious the moment two fast lanes ran together: the
+// digest lane's Recent picks panel emptied itself on every backfill dispatch,
+// and `/api/auto-queue/status` — which reads the FILE — showed zero picks for a
+// lane that was demonstrably working. `lane-runner.spec.ts`'s concurrency test
+// is what found it.
+//
+// One object mutated by both is what makes a write of it true for both. Its
+// lifetime is the auto-runner singleton's, which the e2e harness already clears
+// between specs (api/test/invalidate-cache) — so a reset corpus does not
+// inherit a previous spec's picks.
+async function sharedAutoQueueState(paths: Paths): Promise<AutoQueueState> {
+ const singleton = getSingleton();
+ if (singleton.state && singleton.stateFile === paths.autoQueueStateFile) {
+ return singleton.state;
+ }
+ const state = await readAutoQueueState(paths);
+ singleton.state = state;
+ singleton.stateFile = paths.autoQueueStateFile;
+ return state;
+}
+
export type AutoRunnerStatus = {
kind: AutoQueueKind;
running: boolean;
@@ -652,7 +690,10 @@ async function runLoop(
ctx: JobRunContext,
): Promise<void> {
const tracker = makeTaskTracker(ctx, onLog);
- const state = await readAutoQueueState(paths);
+ // SHARED with every other lane's runner in this process — see
+ // sharedAutoQueueState. A per-runner copy makes each persist erase the other
+ // lanes' pick logs.
+ const state = await sharedAutoQueueState(paths);
const kindState = state[kind];
const runtime = kindState.runtime;
// Drop long-lapsed platform cooldowns on boot; entries still in (or recently
diff --git a/editor/e2e/lane-runner.spec.ts b/editor/e2e/lane-runner.spec.ts
@@ -36,8 +36,13 @@ import {
// 5. THE TWO LANES OVERLAP. This is the `backfill.spec.ts` "a digest runs
// concurrently with a backfill, not behind it" invariant, at the runner:
// there it is two per-channel JOBS holding distinct queue keys, here it is
-// two runner loops dispatching in-process units at the same moment. The
-// mechanism differs; the property must not.
+// two runner loops, both on queueKey "", dispatching in-process units that
+// hold no key at all. The mechanism differs; the property must not. It is
+// measured from the two timestamped PICK LOGS afterwards — do their
+// intervals overlap — rather than by catching both lanes in flight at one
+// sampled instant. The second is the same property measured by a coin
+// toss, and a loaded machine loses it to sampling rather than to a
+// regression.
//
// Engines: the digest and attribution lanes use the ollama HTTP stub
// (e2e/fixtures/ollama-stub.mjs, wired via OLLAMA_URL), diarization uses
@@ -137,9 +142,17 @@ function laneSettings(over: Record<string, unknown> = {}) {
}
type RunnerStatus = {
- runner: { running: boolean; jobId: string | null; idleReason: string | null; inFlight: unknown[] };
+ runner: {
+ running: boolean;
+ jobId: string | null;
+ idleReason: string | null;
+ inFlight: unknown[];
+ };
policy: { enabled?: boolean; order?: string };
pendingByLeaf: Record<string, number>;
+ // The persisted pick log — one entry per unit the runner handed out, newest
+ // first, each stamped with the moment it was dispatched.
+ picks: { at: number; leafId: string; videoId: string }[];
};
async function status(
@@ -347,9 +360,13 @@ test("the digest and backfill lanes dispatch at the same time", async ({
request,
}) => {
test.setTimeout(SLOW);
- // ENOUGH WORK that both lanes are provably still going when the other is
- // looked at. One video each finishes faster than a poll.
- const ids = Array.from({ length: 8 }, (_, i) => `lanevid01${i}0`);
+ // ENOUGH WORK that each lane's dispatching spans seconds rather than one
+ // tick: 12 digests on one side, 12 diarizations plus 12 attributions on the
+ // other. Both stay inside the pick log's 50-entry cap.
+ const ids = Array.from(
+ { length: 12 },
+ (_, i) => `lanevid1${String(i).padStart(3, "0")}`,
+ );
await seedDigestChannel(page, ids);
for (const id of ids) {
await writeFile(resolvePath(dataRel(id, "audio.mp3")), "fake audio\n");
@@ -367,20 +384,36 @@ test("the digest and backfill lanes dispatch at the same time", async ({
await control(request, "digest", "start");
await control(request, "backfill", "start");
- // Both dispatching at ONE moment. A lane sitting at zero while the other
- // works is exactly the failure a shared resource gate would produce, so it is
- // this simultaneity that is polled rather than "both eventually did work".
+ // Let both lanes do their work, then read the two PICK LOGS — which are
+ // timestamped and persisted, so this measures what happened rather than what
+ // a poll happened to catch. An earlier version of this test sampled for "both
+ // in flight at one instant" and lost the race on a loaded machine: same
+ // property, measured by a coin toss.
await expect
.poll(
async () => {
- const res = await request.get(`${baseUrl}/api/auto-queue/status`);
- const body = (await res.json()) as Record<string, RunnerStatus>;
- return (
- body.digest.runner.inFlight.length > 0 &&
- body.backfill.runner.inFlight.length > 0
- );
+ const digest = await status(request, "digest");
+ const backfill = await status(request, "backfill");
+ return Math.min(digest.picks.length, backfill.picks.length);
},
- { timeout: 90_000, intervals: [100] },
+ { timeout: 120_000, intervals: [500] },
)
- .toBe(true);
+ .toBeGreaterThanOrEqual(4);
+
+ const digest = await status(request, "digest");
+ const backfill = await status(request, "backfill");
+ const span = (s: RunnerStatus): [number, number] => {
+ const at = s.picks.map((p) => p.at);
+ return [Math.min(...at), Math.max(...at)];
+ };
+ const [dFrom, dTo] = span(digest);
+ const [bFrom, bTo] = span(backfill);
+ // OVERLAPPING INTERVALS IS THE WHOLE ASSERTION. If either lane were waiting on
+ // the other, every one of its picks would fall after the other's last — two
+ // disjoint intervals, which is exactly what a shared queue key produces for
+ // the per-channel jobs backfill.spec pins.
+ expect(
+ dFrom <= bTo && bFrom <= dTo,
+ `digest picks ${dFrom}-${dTo} vs backfill picks ${bFrom}-${bTo} do not overlap`,
+ ).toBe(true);
});