commit 4e50a1b902e81695f6765e5bd6c570fb0cdc8fb7
parent 35b5c5bc1c08ab90e9dab93257eeb70f9bcb5fe9
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Tue, 8 Sep 2026 16:18:05 -0400
plans: root-cause plan for the attribution batch hang
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Diffstat:
1 file changed, 168 insertions(+), 0 deletions(-)
diff --git a/plans/attribution-batch-hang.md b/plans/attribution-batch-hang.md
@@ -0,0 +1,168 @@
+# Root-cause the attribution batch hang
+
+## Context
+
+The de-flake job (2026-09-08, `a32612a`→`fea2996` on `one-core/phase-1`) fixed three of four
+e2e flakes. The fourth reproduces: `editor/e2e/attribution.spec.ts:349` "Run from the video
+page runs that video and no other" fails 1-2 in 10 under `--repeat-each 10`, waiting the full
+90 s for the batch's closing "1 done" line. It runs through `runOperationBatch`, the runner
+Phase 1.2 introduced, so a batch that never finishes after its unit's write is a candidate
+production bug and could surface on operator gate A. The operator chose to root-cause it
+before merge. Land on `one-core/phase-1` (same branch as the de-flake commits; cherry-picks
+to `main` cleanly).
+
+## What is known, and why it does not add up
+
+Recorded in `plans/STATE.md` ("The fourth REPRODUCED") from the earlier experiment:
+
+1. Server log on every failure ends on three lines: the "Backfill attribution-channel …"
+ header, "Attribute attrvid0002: 120 cues → 1 chunk(s), text-only.", and the ollama
+ timing line. The unit's closing "Attribute attrvid0002 (text-only): …" line never appears.
+2. The freshness pill reads "current", so `attribution.json` was written.
+3. `/api/pulse` answered with an unchanging `rev` for 90 s, so the job record never moved.
+
+The code says these three cannot all be true. `common/controller/attributeOne.ts:306-314`:
+`await writeAttribution(...)` then `log(...)` are adjacent statements; `log` is
+`opts.onLog` = `task.onLog` (`operationBatch.ts:1768`) → `taskHooks.ts:87-104` forwards first
+→ the job's `onLog` (`common/jobs/streamCommand.ts:311-315`) which is a synchronous
+`fileStream.write` + `safeEnqueue`. `writeAttribution` (`common/lib/attribution-server.ts:40-48`)
+is `writeFile(tmp)` + `rename`, no lock, no hook. So if (2) then the closing line reached the
+log FILE, even if the SSE stream was already closed (`safeStreamController.ts:28-41` drops
+silently after `cancel()` → `markClosed()`, `streamCommand.ts:289-295`). The two panel-side
+"reconcile against the file" attempts were reverted with "the file has no more than the
+panel" — that reading was taken when the stream ended, which under hypothesis 1 below is
+BEFORE the batch continues. At least one of (1)–(3) is a measurement artifact; the plan's
+first job is to find which.
+
+## Code path (verified 2026-09-08 at `fea2996`)
+
+1. Button: `editor/app/channels/[slug]/videos/[id]/components/SpeakerBodies.tsx:246-252` →
+ `runOperationForVideoAction` (`operationActions.ts:26-56`: `kindIds:[op]`, `ids:[videoId]`).
+2. Client: `common/components/StreamActionLog.tsx:127-163` reads the stream until `done`,
+ then `setRunning(false)` and `router.refresh()` (`:163`). The 1 s status poll (`:101-126`)
+ runs only while `running`.
+3. Job: `common/controller/operationJobs.ts:78-152`, `runManagedFunction` on
+ `BACKFILL_QUEUE`; the "N done" summary is `:128-146`, after `runOperationBatch` returns.
+4. Stream/registry: `common/jobs/streamCommand.ts:264-370`; finalize `:338-362` →
+ `registry.finalize` (`common/jobs/registry.ts:202-218`).
+5. Batch: `common/controller/operationBatch.ts:1544`; `next()` `:1682-1710`; `runOne`
+ `:1751-1785`; `runPool` call `:1799-1826` with `limit()` `:1801-1824`, `finite: true`.
+6. Pool: `common/jobs/concurrentRunner.ts:82-133`. If `limit()` reads 0 when the last unit
+ settles, `:92-93` waits forever and the finite-exhaustion break (`:104-108`) is
+ unreachable. Every 0 from `backfillLaneLimit` (`operationBatch.ts:313-410`) carries a
+ `hold` whose message is logged edge-triggered at `:1806-1812` — a line the failures
+ do NOT show, which argues against this path unless that log was also lost.
+7. Unit: `runOperationUnit` `:947-1058` → `executeUnit` `:1066-1110` (`acquireLlmSlot`
+ `:1086`, `op.run` `:1090`, release in `finally` `:1106-1110`) → `attributeOne.ts:124`;
+ engine call `common/lib/digestApps.ts:168-270` (`await res.json()` at `:220` precedes the
+ timing log at `:255`, so the body is drained; stub `editor/e2e/fixtures/ollama-stub.mjs:204-221`).
+8. Pulse: `editor/app/api/pulse/route.ts:65-118` folds `id:status:startedAt:endedAt:
+ draining:progress:tasks.length` per job off `globalThis.__yttJobRegistry__`.
+
+## Hypotheses, ranked, each with the observation that decides it
+
+| # | Hypothesis | Blocking point | Decided by |
+|---|---|---|---|
+| H1 | The SSE stream is torn down early (RSC/action teardown from a pulse-driven `router.refresh()`); the batch finishes normally but the panel froze, so the test's 90 s wait on the panel text never sees "1 done". Pill reads "current" because `router.refresh()` at `:163` re-rendered. (3) was misread. | none in the server | server-side: closing line + summary present in the job log FILE with timestamps after the client's `cancel()`; `record.status` becomes `done` |
+| H2 | `limit()` returned 0 at the wake after the unit settled; `runPool` idle-waits forever. | `concurrentRunner.ts:93` | per-wait trace shows `target=0, inFlight=0, drain=false`; hold line at `operationBatch.ts:1806` in the file |
+| H3 | The unit promise never settles after `op.run` returns (LLM slot release, `task.end()`, `reportProgress`). | `operationBatch.ts:1106-1110` or `:1779-1783` | `runOne`'s `finally` timestamp missing while `attributeOne.ts:314` timestamp present |
+| H4 | Finalize skipped or registry global replaced, so pulse folds a stale/empty list while the batch is actually done. | `streamCommand.ts:338-362`, `invalidate-cache` route | `/api/jobs/<id>/log` status after the failure: `done` vs `running` |
+| H5 | A late abort (previous test's `cancelLiveJobs`) fired `opts.signal`, `runPool` broke without the summary. | `runPool` exit | `signal.aborted === true` at pool exit |
+
+## Step 1 — Instrumented reproduction (worktree only, never committed)
+
+Worktree of `fea2996` per `plans/FACTS.md` (composed fixture site copied into
+`export/public`; add-by-path rule irrelevant since nothing is committed from it). Set
+`stdout: "pipe"`/`stderr: "pipe"` on the editor webServer in the worktree's
+`editor/playwright.config.ts`. All instrumentation goes to `console.error` with an ISO
+timestamp and the job id, never through the job's `onLog` (which is part of what is under
+test):
+
+- `common/controller/attributeOne.ts:306` and `:315` — before `writeAttribution`, after the
+ closing `log` call.
+- `common/jobs/streamCommand.ts:311` — per `onLog`: line counter, `safe.isClosed()`,
+ first 60 chars. `:289` — the `cancel()` with timestamp. `:338/:351/:362` — which settle
+ branch ran and `record.status`.
+- `common/jobs/concurrentRunner.ts:88-93` — per wait iteration: `target`, `inFlight.size`,
+ `drainSignal.aborted`, `signal.aborted`, `finite`.
+- `common/controller/operationBatch.ts:1783` (unit `finally`) and `:1826` (after `runPool`).
+- `editor/e2e/attribution.spec.ts:349` — on failure only (wrap the "1 done" expect in
+ try/catch, rethrow): fetch `/api/jobs/<id>/log` (job id from the panel's `data-*` or the
+ status poll URL, see `StreamActionLog.tsx:101-126`) and `/api/pulse`, print status, line
+ count and the last five lines to the test output; then rethrow.
+
+Run: `node ../scripts/queue-lock.mjs --ports PORT:3011,EXPORT_PORT:3010,OLLAMA_STUB_PORT:11435
+-- pnpm exec playwright test attribution.spec.ts --repeat-each 15` from the worktree's
+`editor/`, detached, watched with Monitor. Expected 1-3 failures. Collect, for each failure,
+the ordered stderr trace for that job id and the test's failure dump, and decide H1–H5 by
+the table. If 15/15 pass, run once more with `--repeat-each 25`; if still clean, the
+instrumentation changed the timing — drop the `concurrentRunner` trace (the noisiest) and
+re-run 15. Stop after three attempts and record "unreproduced under instrumentation" with
+the run log paths.
+
+## Step 2 — Fix at the cause (one commit, with a test that fails before it)
+
+Exactly one of these, chosen by Step 1; do not ship a fix for a hypothesis the trace killed.
+
+- **H1 (test/UI-side):** the product is correct (the batch finished, disk is right), the
+ panel is not: after the stream ends `StreamActionLog` must show the job's terminal state.
+ Fix in `common/components/StreamActionLog.tsx:141-163`: when `reader.read()` returns
+ `done` and the job is not terminal per the status poll, keep polling `/api/jobs/<id>/log`
+ until it is, and render its tail (the file has every line, `streamCommand.ts:311`). Test:
+ a unit test for the reconcile (or an e2e route mock that cancels the stream early) that
+ fails at `fea2996`. Also fix the trigger if it is ours: if the teardown comes from the
+ pulse-driven refresh re-rendering the panel, gate that refresh while a run is streaming.
+ (The de-flake implementer tried two reconcile variants and measured no change — but its
+ measurement read the file at stream end, before the batch continued. A reconcile that
+ keeps polling until terminal is a different fix; verify against the trace, not the memo.)
+- **H2:** find why `backfillLaneLimit` (or `isGateHeld`) reads 0 after the unit — the trace
+ says which term. Fix that term (e.g. a settings read racing the test's `writeSettings`, or
+ `transcriptionActivity()` seeing a stale worker). Do NOT change `runPool`'s hold
+ semantics ("a zero limit is a hold, never a stop", `operationBatch.ts:22-25`). Add a
+ `concurrentRunner.test.ts` case only if the runner itself is at fault.
+- **H3:** the pending await in `executeUnit`'s `finally` or `runOne`; fix and add a unit test
+ in `common/controller/` that runs a unit whose engine call resolves and asserts the pool
+ exits.
+- **H4:** make finalize unconditional (`streamCommand.ts:338-362`, a throwing `onLog` in
+ `.catch` must not skip `registry.finalize`) or stop `invalidate-cache` from replacing the
+ global registry under a live job. Unit test in `common/jobs/`.
+- **H5:** `resetData` ordering / `cancelLiveJobs` aborting a record the next test owns;
+ fix in `editor/app/api/test/invalidate-cache/route.ts` or `editor/e2e/helpers.ts`.
+
+Remove every instrumentation line before committing; `git diff fea2996 --stat` in the
+primary checkout must list only the fix, its test, and the plan records.
+
+## Step 3 — Records
+
+- `plans/STATE.md`: replace "The fourth REPRODUCED, is NOT fixed" with what it was, the
+ discriminating trace excerpt (≤ 8 lines), and the fix sha. If the earlier claim "the log
+ file has no more than the panel" was the artifact, say so, so nobody trusts it again.
+- `plans/deflake-e2e.md`: a "### 4, as shipped" section with the same.
+- `plans/FACTS.md` one-core Phase 1 section: one paragraph on the runner/stream fact
+ learned (whichever hypothesis held), file:line pinned to the fix sha.
+
+## Not in scope
+
+- Any change to lane hold semantics, `retries`, `test.slow`, or serial markers.
+- The homepage theme test's early read (`homepage/e2e/theme.spec.ts:47`), noted in the
+ de-flake review as a weakened assertion, not a flake.
+- Phase 2 of one-core.
+
+## Verification
+
+- `tsc --noEmit` in common, export, editor; `pnpm --filter yt-dlp-transcript-common test`
+ (the new test fails at `fea2996`, passes at the fix — show both runs).
+- `attribution.spec.ts --repeat-each 20` from a worktree at the fix sha, no instrumentation:
+ 20/20 (baseline was 1-2 failures per 10, so 20 clean runs is the bar, not 10).
+- Full editor e2e from the worktree at the final sha, detached, Monitor-watched: 515/515.
+- Commit shape: (1) the fix + its test, message carrying the trace excerpt and the 20/20;
+ (2) `plans:` records. Trailers as on the de-flake commits. Never `git add -A`
+ (`editor/content` symlink).
+
+## Execution
+
+Per the standing cadence: one Opus general-purpose implementer with the fish/`commit -F`,
+no-boot-against-`transcripts/`, detached-e2e and DO-NOT-POLL (Monitor) instructions;
+report contract = commit shas, exact gate outputs, the decisive trace lines, which
+hypothesis held and which were killed, every divergence. A second Opus agent reviews the
+diff before the final full e2e.