[omp major-2] Overflow fan-out not restart-safe despite doc claim #77

Closed
opened 2026-09-05 00:15:04 +00:00 by crueber · 5 comments
Owner

[omp major-2] Overflow fan-out not restart-safe despite doc claim

06 §335 claims "the activity payload is the only queue that survives a restart", but the only driver is Run (wake channel + 1-min webhook sweep + daily retention); sweepWebhooks (internal/notify/tasks.go:177-203) starts only webhooks tasks. Nothing on startup or in any sweep re-drains activity events whose notify-fanout task was in flight at process death. Overflow emissions (>100 recipients) and shortfall repairs pending at restart are lost.

Fix

Either add a startup/sweep pass that re-drains undrained fan-out events from the activity log, or correct the doc claim if restart-safety is deliberately out of scope (law 12 — but fix preferred: the claim is load-bearing for correctness). Regression test (seed undrained events, restart/sweep, assert delivery). Coverage gate holds.

Acceptance criteria

  • Undrained fan-out events are redelivered after restart (or doc honestly corrected with rationale).
# [omp major-2] Overflow fan-out not restart-safe despite doc claim 06 §335 claims "the activity payload is the only queue that survives a restart", but the only driver is `Run` (wake channel + 1-min webhook sweep + daily retention); `sweepWebhooks` (`internal/notify/tasks.go:177-203`) starts only `webhooks` tasks. Nothing on startup or in any sweep re-drains activity events whose `notify-fanout` task was in flight at process death. Overflow emissions (>100 recipients) and shortfall repairs pending at restart are lost. ## Fix Either add a startup/sweep pass that re-drains undrained fan-out events from the activity log, or correct the doc claim if restart-safety is deliberately out of scope (law 12 — but fix preferred: the claim is load-bearing for correctness). Regression test (seed undrained events, restart/sweep, assert delivery). Coverage gate holds. ## Acceptance criteria - [ ] Undrained fan-out events are redelivered after restart (or doc honestly corrected with rationale).
Author
Owner

Fixed by PR #86 (#86): pending-marked activity events now redrain via a startup + minutely sweep that reuses the webhooks repo enumeration — no new global LIST. Regression tests + doc Decisions entries included.

Fixed by PR #86 (https://git.packden.us/crueber/walhub/pulls/86): pending-marked activity events now redrain via a startup + minutely sweep that reuses the webhooks repo enumeration — no new global LIST. Regression tests + doc Decisions entries included.
Author
Owner

Review: PR #86 (fix/issue-77 @ ff23bc0)

Reviewed in scratch worktree /tmp/pr86. One real bug found and fixed (see below); everything else checks out.

BUG (fixed in fix/issue-77-sweep-retry @ b35725c — needs folding into this PR branch)

  • Transient probe failure orphans a pending drain — internal/notify/tasks.go:477-481 (was), internal/notify/activity.go:87-98 (was). readActivity returns nil for both honest gaps AND transient store errors, and the sweep loop continued on nil while the high-water advanced to end unconditionally. A single flaked GET on a pending seq during the restart sweep permanently skipped it (high-water never revisits) — the exact loss this PR fixes. Fix: new readActivityErr distinguishes error from absent/corrupt; on error the window stops (end = seq-1, next pass retries), gaps/corrupt still skip. docs/features/06_notifications.md:311-312 updated in the same commit (law 12). New TestSweepFanoutTransientProbeRetries fails on the unfixed code (high-water = 1 after a failed probe, want 0) and passes with the fix.

Verified OK (no change needed)

  • No new global LIST — sweepFanout reuses the shared eachRepo (repos/ + repos/<owner>/ prefix LISTs, tasks.go:237-258), the same enumeration sweepWebhooks uses. Fan-out probes are GET-by-key only (collab_state, collab-events/<seq>, collab-fanout/<seq>); no LIST of collab-events/ or collab-fanout/ in prod code. All backends implement ListPrefixes (memory, filesystem, S3, Prefixed).
  • Quiet-repo cost 1 GET/pass in steady state (tasks.go:462-472: collab_state GET, NextSeq <= seen → return, no fan-out). Note: early passes after a restart probe up to 32 activity GETs/repo with history — bounded by the cap, converges, no hidden fan-out (sweep only enqueues in-memory; the worker writes).
  • 32-seq cap can't starve — high-water advances only to the probed end; deeper backlogs converge across minute passes; restart rebuild (empty map → from seq 1) is what makes the restart pass complete. With the fix above, no path permanently skips a pending seq.
  • Completion records bounded, deleted with events — written only for drained pending seqs (tasks.go:331-335; gaps never marked since fanoutOne==false skips the mark); retention deletes event + marker together (tasks.go:722-727). Orphan markers are inert (probe reads the event first).
  • Sync path byte-identical, zero added store ops — sync branch passes pending=false (emit.go:230); omitempty keeps payload bytes unchanged; no enqueue/mark/CAS added on the hot path.
  • At-least-once collapse safe — duplicate drains derive deterministic NotificationIDs: same object key (412 collapses), indexAdd dedups by ID (emit.go:505-548). Durable double-delivery is impossible. (Transient duplicate SSE publish under a sweep/worker race is possible — tasks.go:398 publishes even on the 412 path — but that is the pre-existing at-least-once pattern shared with createOne, ephemeral only.)
  • Sweep/worker race — no coordination needed; completion Create arbitrates (second writer 412s into success). fanoutMu guards only the map, never held across store calls (law 3 ✓). Lock order table→entry only in endIfQuiescent ✓.
  • Coverage 97.1% (≥95% gate ✓); gofmt/go vet clean; stdlib-only imports (added sync in notify.go, errors/context in test only).
  • Docs accurate — 06 §8.1 + Decisions + 14 Seam 5 amendment describe the shipped shape (per-seq records vs cursor rationale cites law 6 correctly).

Test results (with fix)

  • go test -race -count=1 ./internal/notify/... — ok
  • redrain subset (TestSweepRedrains|TestOverflowDrain|TestSweepFanout|TestRetentionDeletesFanout) -race -count=5 — ok
  • coverage 97.1% of statements

MERGE RECOMMENDATION

Blocked: fold fix/issue-77-sweep-retry (b35725c) into fix/issue-77 first (cherry-pick or merge — 4 files, +96/-6). With that in, ready to merge.

## Review: PR #86 (fix/issue-77 @ ff23bc0) Reviewed in scratch worktree /tmp/pr86. One real bug found and fixed (see below); everything else checks out. ### BUG (fixed in `fix/issue-77-sweep-retry` @ b35725c — needs folding into this PR branch) - **Transient probe failure orphans a pending drain — `internal/notify/tasks.go:477-481` (was), `internal/notify/activity.go:87-98` (was).** `readActivity` returns nil for both honest gaps AND transient store errors, and the sweep loop `continue`d on nil while the high-water advanced to `end` unconditionally. A single flaked GET on a pending seq during the restart sweep permanently skipped it (high-water never revisits) — the exact loss this PR fixes. Fix: new `readActivityErr` distinguishes error from absent/corrupt; on error the window stops (`end = seq-1`, next pass retries), gaps/corrupt still skip. `docs/features/06_notifications.md:311-312` updated in the same commit (law 12). New `TestSweepFanoutTransientProbeRetries` **fails on the unfixed code** (`high-water = 1 after a failed probe, want 0`) and passes with the fix. ### Verified OK (no change needed) - **No new global LIST** — `sweepFanout` reuses the shared `eachRepo` (`repos/` + `repos/<owner>/` prefix LISTs, `tasks.go:237-258`), the same enumeration `sweepWebhooks` uses. Fan-out probes are GET-by-key only (`collab_state`, `collab-events/<seq>`, `collab-fanout/<seq>`); no LIST of `collab-events/` or `collab-fanout/` in prod code. All backends implement `ListPrefixes` (memory, filesystem, S3, Prefixed). - **Quiet-repo cost 1 GET/pass in steady state** (`tasks.go:462-472`: `collab_state` GET, `NextSeq <= seen` → return, no fan-out). Note: early passes after a restart probe up to 32 activity GETs/repo with history — bounded by the cap, converges, no hidden fan-out (sweep only enqueues in-memory; the worker writes). - **32-seq cap can't starve** — high-water advances only to the probed `end`; deeper backlogs converge across minute passes; restart rebuild (empty map → from seq 1) is what makes the restart pass complete. With the fix above, no path permanently skips a pending seq. - **Completion records bounded, deleted with events** — written only for drained pending seqs (`tasks.go:331-335`; gaps never marked since `fanoutOne==false` skips the mark); retention deletes event + marker together (`tasks.go:722-727`). Orphan markers are inert (probe reads the event first). - **Sync path byte-identical, zero added store ops** — sync branch passes `pending=false` (`emit.go:230`); `omitempty` keeps payload bytes unchanged; no enqueue/mark/CAS added on the hot path. - **At-least-once collapse safe** — duplicate drains derive deterministic `NotificationID`s: same object key (412 collapses), `indexAdd` dedups by ID (`emit.go:505-548`). Durable double-delivery is impossible. (Transient duplicate SSE `publish` under a sweep/worker race is possible — `tasks.go:398` publishes even on the 412 path — but that is the pre-existing at-least-once pattern shared with `createOne`, ephemeral only.) - **Sweep/worker race** — no coordination needed; completion Create arbitrates (second writer 412s into success). `fanoutMu` guards only the map, never held across store calls (law 3 ✓). Lock order table→entry only in `endIfQuiescent` ✓. - **Coverage 97.1%** (≥95% gate ✓); `gofmt`/`go vet` clean; stdlib-only imports (added `sync` in `notify.go`, `errors`/`context` in test only). - **Docs accurate** — 06 §8.1 + Decisions + 14 Seam 5 amendment describe the shipped shape (per-seq records vs cursor rationale cites law 6 correctly). ### Test results (with fix) - `go test -race -count=1 ./internal/notify/...` — ok - redrain subset (`TestSweepRedrains|TestOverflowDrain|TestSweepFanout|TestRetentionDeletesFanout`) `-race -count=5` — ok - coverage 97.1% of statements ### MERGE RECOMMENDATION **Blocked: fold `fix/issue-77-sweep-retry` (b35725c) into `fix/issue-77` first** (cherry-pick or merge — 4 files, +96/-6). With that in, ready to merge.
Author
Owner

Follow-up to PR #86: reviewer commit b35725c transplanted as PR #87 (branch fix/issue-77-retry). Cherry-pick applied cleanly onto origin/main (cf3383c); notify tests + vet/gofmt green. Not merging per instructions.

Follow-up to PR #86: reviewer commit b35725c transplanted as PR #87 (branch fix/issue-77-retry). Cherry-pick applied cleanly onto origin/main (cf3383c); notify tests + vet/gofmt green. Not merging per instructions.
Author
Owner

PR #87 review (branch fix/issue-77-retry, head 950d64b) — transplant of b35725c, follow-up for #77

FIDELITY: PASS. git diff b35725c origin/fix/issue-77-retry is empty — identical modulo hash (re-cherry-pick). 4 files: internal/notify/tasks.go, internal/notify/activity.go, internal/notify/redrain_test.go, docs/features/06_notifications.md.

CORRECTNESS:

  • tasks.go:479-486 — transient probe error sets end=seq-1 and breaks; high-water write (tasks.go:501-504, if end > seen) can never advance past the unprobed seq. First-seq failure leaves seen==0, next pass retries. Correct.
  • activity.go:99-112 — readActivityErr splits transient store failure (nil, err) from absence/corruption (nil, nil). Corrupt reads as absent = pre-existing contract (old readActivity returned nil for both too), so no behavior change there; gaps still skip. Callers of readActivity (20+ sites, incl. tasks.go:356,723, webhooks.go:365) untouched via the thin wrapper — no seam change.
  • Doc lines (06_notifications.md:311-313) match the code exactly: stop-the-window, high-water never passes unprobed seq, next pass retries.
  • No new deps (stdlib context/errors in test only); no lock changes (probes outside fanoutMu, as before); Law 12 satisfied — behavior paragraph updated in the same change.

TEST GENUINENESS: TestSweepFanoutTransientProbeRetries fails without the fix by construction — old code conflated error->nil->continue, loop ran to end=1, fanoutSeen advanced to 1, and the seen!=0 assertion fires (plus the healed pass would never drain the orphaned seq). No base-check run needed; the assertion directly distinguishes.

RESULTS (scratch worktree /tmp/pr87 @950d64b, removed after):

  • go vet ./internal/notify/... clean; gofmt clean.
  • go test -race -count=1 ./internal/notify/... ok (1.7s), incl. TestSweepFanoutTransientProbeRetries PASS.
  • coverage 97.1% statements, gate (>=95%) holds.

MERGE RECOMMENDATION: ready to merge (not merging per instructions).

PR #87 review (branch fix/issue-77-retry, head 950d64b) — transplant of b35725c, follow-up for #77 FIDELITY: PASS. git diff b35725c origin/fix/issue-77-retry is empty — identical modulo hash (re-cherry-pick). 4 files: internal/notify/tasks.go, internal/notify/activity.go, internal/notify/redrain_test.go, docs/features/06_notifications.md. CORRECTNESS: - tasks.go:479-486 — transient probe error sets end=seq-1 and breaks; high-water write (tasks.go:501-504, if end > seen) can never advance past the unprobed seq. First-seq failure leaves seen==0, next pass retries. Correct. - activity.go:99-112 — readActivityErr splits transient store failure (nil, err) from absence/corruption (nil, nil). Corrupt reads as absent = pre-existing contract (old readActivity returned nil for both too), so no behavior change there; gaps still skip. Callers of readActivity (20+ sites, incl. tasks.go:356,723, webhooks.go:365) untouched via the thin wrapper — no seam change. - Doc lines (06_notifications.md:311-313) match the code exactly: stop-the-window, high-water never passes unprobed seq, next pass retries. - No new deps (stdlib context/errors in test only); no lock changes (probes outside fanoutMu, as before); Law 12 satisfied — behavior paragraph updated in the same change. TEST GENUINENESS: TestSweepFanoutTransientProbeRetries fails without the fix by construction — old code conflated error->nil->continue, loop ran to end=1, fanoutSeen advanced to 1, and the seen!=0 assertion fires (plus the healed pass would never drain the orphaned seq). No base-check run needed; the assertion directly distinguishes. RESULTS (scratch worktree /tmp/pr87 @950d64b, removed after): - go vet ./internal/notify/... clean; gofmt clean. - go test -race -count=1 ./internal/notify/... ok (1.7s), incl. TestSweepFanoutTransientProbeRetries PASS. - coverage 97.1% statements, gate (>=95%) holds. MERGE RECOMMENDATION: ready to merge (not merging per instructions).
Author
Owner

Follow-up PR #87 (transplant of the reviewed sweep-retry fix) merged. #77 complete. Closing.

Follow-up PR #87 (transplant of the reviewed sweep-retry fix) merged. #77 complete. Closing.
crueber added this to the v1 milestone 2026-09-10 22:27:20 +00:00
Sign in to join this conversation.
No milestone
No project
No assignees
1 participant
Notifications
Due date
The due date is invalid or out of range. Please use the format "yyyy-mm-dd".

No due date set.

Dependencies

No dependencies set

Reference
crueber/walhub#77
No description provided.