Files
xarmian c73584088f fix(watchevents): detect a half-open Redis connection with a bus heartbeat (BUG-2769) (#1199)
* fix(watchevents): detect a half-open Redis connection with a bus heartbeat (BUG-2769)

internal/watchevents had the same defect as internal/events did, by the same
mechanism: ChannelWithSubscriptions on a connection whose go-redis health check
only writes. PubSub.Ping calls writeCmd and returns without reading a reply
(v9.22.0), so a route that stops carrying traffic without closing is invisible —
the instance blocks on a read forever while its replay buffer goes on looking
complete.

Named as a class sweep in BUG-2738's filing and deferred there. It became
load-bearing when that unit shipped: docs/deployment.md told operators the gap
was "closed on the activity stream and still open on the watch stream". This
diff falsifies that, which is why the prose sweep is part of it.

THE PORT IS SMALLER THAN THE ORIGINAL BY DESIGN. This bus holds ONE
process-wide subscription created in its constructor, off any request path, so
none of BUG-2747's establishment machinery exists to interact with: no
per-workspace map, no establishment record, no single-establisher wall, no
concurrency cap, no bounded-parallel recovery, and no per-workspace cycle
scoping. Cost is flat too — one frame per instance per interval regardless of
workspace count.

NO COMPANION COUNTER, and that was CHECKED rather than inherited.
internal/events needs pad_event_subscription_cycled_total because its
dropWorkspaceCoverage returns early when a workspace has no buffer, so the reset
reason under-reports the early-wedge case. dropCoverage here has no such branch:
it replaces the buffer and reports unconditionally, so idle_timeout is a
complete count on its own and a second metric would be a number needing to be
explained against its neighbour for no signal.

THE RECEIVE LOOP NOW OWNS ITS SUBSCRIPTION AND CONTEXT. A cycle replaces the
subscription under a running bus, and the loop reading the old one must tell "I
was replaced" from "the client died" — the second logs an ERROR and moves a
counter documented to mean the instance has gone deaf. The cycle cancels that
loop's own context before closing its PubSub, so it leaves by the quiet door.
Its own test.

I PORTED A FLAW ALONG WITH THE STRUCTURE, and the wiring test caught it: both
maintenance halves shared one kick channel, so whichever goroutine was waiting
consumed it and the other stayed on the stale cadence. internal/events' mutation
matrix found exactly that (M11c) and fixed it; the fix did not come across. That
is the contamination hazard this port's grounding warned about, in its literal
form, caught by the CONVE-19 test rather than by review.

Two more found by mutation, both missing tests rather than missing code: nothing
asserted that ordinary traffic keeps the instance alive (removing the per-frame
stamp survived, because every other test drives idleness through the clock), and
nothing asked for a SECOND detection (a replacement inheriting stale stamps
gives a detector that works exactly once, which is worse than one that never
runs because it looks like it works). The second needed a direct assertion on
the install stamps, because the behavioural route re-stamps the field it was
meant to be testing.

Trio in one commit as required: reason enumeration, the
pad_watchevents_sequence_resets_total Help string, and docs/deployment.md — plus
the two BUG-2738 sentences this falsifies and a new section explaining how the
watch bus differs from the activity one.

Claude-Session: https://claude.ai/code/session_01JVDBKbgn3Xt7ndW1YoYd8X

* fix(watchevents): fence stragglers and re-validate before the drop (codex r1)

Three findings, and two of them are BUG-2738 fixes I again failed to bring
across with the structure. That is now three times in one port: the shared kick
channel, the stale idle decision, and the missing generation. The mechanism is
the same each time — I ported what the code DOES and not what its review
history taught it, and each was caught by a test or a reviewer rather than by
me reading the source I was copying.

STALE IDLE DECISION. cycleIfIdle decided under one lock and tore down under
another; a heartbeat or notification arriving between them left a demonstrably
alive subscription being dropped and every client on the instance resynced for
nothing. BUG-2738 fixed exactly this at its round 11. Re-validated immediately
before the drop, with a positional seam so a test can land the recovery inside
the window rather than racing it.

NO GENERATION FENCE. Cancelling a receive loop and closing its PubSub does not
JOIN the goroutine, and go-redis's channel is buffered, so a frame from a
replaced subscription could still stamp the replacement's liveness, append to
its buffer, or drop its coverage. On a wedged route that is the worst
direction: the dead connection's buffered tail suppressing the detector for its
successor. One check at the top of the frame handler covers all three, because
the three must agree about whether a frame belongs to the live subscription.
The probe stamp is fenced separately, since a slow publish can outlive the
subscription it was sent for.

A COPIED COST PARAGRAPH THAT CONTRADICTED ITS OWN SECTION. The activity bus's
"each workspace has its own subscription, N frames per interval" text sat below
the new watch-specific section saying the opposite. Retitled and moved above it.

FOUR INSTRUMENT DEFECTS ON THE WAY, all found by mutation:

- Nothing asserted ordinary traffic keeps the instance alive — every other test
  drives idleness through the clock, so removing the per-frame stamp survived.
- Nothing asked for a SECOND detection, so a replacement inheriting stale stamps
  gave a detector that works exactly once — worse than one that never runs,
  because it looks like it works. Needed a direct assertion on the install
  stamps, since the behavioural route re-stamps the field under test.
- The generation tests asserted the PREDICATE, not that the loop calls it.
- And that wiring test could not discriminate on a frozen clock, where a stamp
  writes the value already there. It advances the clock first now.

Claude-Session: https://claude.ai/code/session_01JVDBKbgn3Xt7ndW1YoYd8X

* fix(watchevents): make the generation fence atomic with what it guards (r2)

Two P1s, both mine, both the same shape: a check in one lock acquisition and the
mutation it guards in another.

THE FENCE WAS NOT ATOMIC WITH ITS MUTATIONS. One check at the top of the frame
handler read well and guarded nothing reliably — a replacement between that
check and stampLastSeen / fanOutFromRedis / dropCoverage let a straggler through
to any of them. The generation now travels TO each mutation and is re-checked
under the same lock that mutates. A stale notification entering the
replacement's buffer is the worst of the three: it makes the instance vouch for
a span it never received, which is the false coverage claim this whole family
exists to remove.

THE OLD GENERATION STAYED CURRENT ACROSS THE REPLACEMENT. subGen was
incremented only after the new subscription was confirmed, leaving the cancel,
the close, the dial and a round trip during which the OLD generation still
passed every fence. Retired at teardown now, so during resubscribe NO generation
is current and a late frame is ignored everywhere. That also makes the failure
path honest: the "no notifications until restarted" log was false — no
generation is current, so the next idle tick tries again.

Revalidation and the drop are now ONE critical section rather than two, for the
same reason at one level down: a frame arriving between them was silently
discarded by a drop already decided on.

Also: phase 1 no longer starts the maintenance goroutines, and the watch bus's
phase is logged at startup — an operator cannot read an absence of idle_timeout
without knowing whether the detector was running, and the two flags are
independent.

DOCS still described the workspace model in the section that claims to cover
both buses: one heartbeat "per subscribed workspace", a phase table naming only
PAD_EVENTS_HEARTBEAT, and coverage described as a workspace's. Generalised.

Two more instrument gaps, both found by mutation: nothing asserted a straggler
cannot enter the replacement's BUFFER (only the stamp was covered), and the
phase-1 goroutine gate is untested by design — removing it changes no behaviour,
only goroutine count, and the only assertion is a census that would be flaky
here. Said out loud rather than left to look like coverage.

Claude-Session: https://claude.ai/code/session_01JVDBKbgn3Xt7ndW1YoYd8X

* test(watchevents): prove each generation fence on its own

Round 3's fix put a generation check in each of the four places a frame
from a replaced subscription can mutate shared state, rather than one
check at the top of the receive path — a check in one lock acquisition
and a write in another is a TOCTOU, which is what codex blocked.

Four checks means four mutations, and the matrix found the first pass of
tests could not tell them apart: removing the append's check, or the
coverage drop's, left every test green. Not because the guards were
redundant — because no test drove those paths with a stale generation.
The straggler tests all enter through fanOutFromRedis, whose own guard
returns early and hides the one below it, and nothing at all drove
dropCoverageForGen with a straggler.

So the fences are asserted one at a time, each through the entry point
that actually reaches it:

  epoch bookkeeping   fanOutFromRedis with a foreign epoch — the loudest
                      of the four, since an accepted straggler would
                      rewrite the id space and resync every client on the
                      instance
  buffer append       fanOutLocally directly, under the guard above it
  coverage drop       dropCoverageForGen, previously undriven
  liveness stamp      stampLastSeen, which would otherwise let a dead
                      socket's traffic hold detection open

Each fails against removal of the single check it names (M16/M17/M19 and
the existing stamp mutation), and the four together still pass the
end-to-end straggler tests unchanged.

Refs BUG-2769

* test(metrics): prove the two new watch signals reach the registry

Both were wired and neither was asserted at the metrics layer, which is
where docs/deployment.md's claims about them actually live. A reason or
a callback that never reaches the registry is a runbook pointing at a
series that does not exist, and nothing in internal/watchevents can
catch that — its observer is an interface, satisfied by a test double.

  pad_watchevents_heartbeat_publish_failures_total  incremented six
  times, a count no other assertion in that test uses, so a callback
  wired to the wrong counter cannot land on the right number by
  coincidence. Fails when the increment is pointed at a neighbour.

  sequence_resets_total{reason="idle_timeout"}  asserted with the
  literal label, alongside the four spellings already pinned there and
  for the same reason BUG-2739's rename left that test behind. Fails
  when the constant drifts.

Also corrects the shared "what happens if you run them out of order"
paragraph, which moved under a heading covering both buses while still
describing only one: it said the frame travels on "the workspace's event
channel" and that an un-upgraded instance resyncs "for every workspace",
neither of which is the watch bus, where there is one channel and one
buffer per instance. The blast radius differs in scale between the two
and the paragraph now says so.

Refs BUG-2769

* docs(watchevents): correct three counted claims that stopped being true

All three said "three" where the code now has four, and each was
accurate when written — the fourth fence (the epoch bookkeeping in
fanOutFromRedis) was identified after them, in the pass that found the
matrix could not tell the guards apart.

That is the whole failure mode: a count is a claim, and a claim written
before the last change is wrong afterwards with nothing to notice it.
Two of the three sat inside a comment ABOUT how carefully the guards
were enumerated, and one names them now instead of counting them, so
the next site added has to appear in the list or contradict it visibly.

Found by sweeping the branch diff for counted prose rather than by
rereading, which is what had already missed them twice.

Refs BUG-2769

* test(config): close the other half of the two-flag independence claim

The flag tests asserted PAD_WATCH_HEARTBEAT does not move
EventsHeartbeat and stopped there, while the comment above them and the
deployment doc both claim the two buses roll INDEPENDENTLY. That is a
biconditional and one leg does not establish it: a Load() that pointed
PAD_EVENTS_HEARTBEAT at both fields passed everything. Now both
directions are asserted, and the events leg checks its own premise
first, so a fixture that stopped setting the flag fails as a fixture
rather than as a pass.

Also pins env-over-file precedence for the watch flag, in the direction
that actually matters: PAD_WATCH_HEARTBEAT=false over
watch_heartbeat=true in config.toml. That is the rollback for a bad
phase-2 flip, and an operator reaching for it mid-incident cannot be
editing a file on every host.

Mutation matrix, each detected: the env var wired to the neighbouring
field, the env var never read at all, and the toml tag dropped.

Refs BUG-2769

* test(watchevents): fix five tests that passed for the wrong reason

Codex round 4 went at test honesty rather than correctness and found no
BLOCK, but it found five assertions that hold whether or not the thing
they name works. Each is now driven through the path it claims, and each
was mutation-checked against the specific defect it exists to catch.

  the malformed-frame contract  only ever called isWatchHeartbeat. The
  predicate can be perfect while the receive loop routes every "hb|…"
  payload to the ignore arm without asking it, which is the defect, and
  the test's name promises coverage ends — a claim about the loop. Now
  published on the real channel, with a well-formed frame as the control
  so the assertion cannot be satisfied by a loop that finds everything
  undecodable.

  the receive-loop wiring test  published, slept 300ms, and asserted
  nothing had changed. A loop that stalled or never started satisfies
  that perfectly. There is no natural signal to wait on instead, because
  a frame the fence refuses is by design invisible — hence a seam that
  fires after the loop handles a frame whichever arm it took. Bounded,
  so a stalled loop fails with a message rather than a package timeout,
  and followed by a control that the same loop still accepts a frame
  whose generation matches.

  the quiet-exit test  asserted only that no loud exit was reported,
  which a replaced goroutine that never exits at all also satisfies —
  a leak, and the worse outcome. Now joins the loop first via a
  process-wide live-loop count, then checks the counter, so it is a
  statement about a goroutine that has finished.

  the maintenance-loop wiring test  claimed both halves and observed a
  heartbeat, which a loop that started only the publisher passes. The
  idle half cannot be proved there at all: against a live miniredis this
  bus's own heartbeats come back and refresh liveness every cadence, so
  wedging it with the loop running is a race against the publisher —
  which is what my first fix for this turned out to be, flaky at 2 in 3.
  Renamed to what it proves, pointing at the blackhole end-to-end test,
  which drives the scanner for real and detects both mutations.

  the straggler test  never delivered a straggler. It incremented subGen
  by hand, called isCurrentGen, and compared an unchanged timestamp
  without touching a mutation path — green with every fence removed.
  Deleted rather than repaired: the four-way per-fence test added
  earlier covers it properly, and isCurrentGen went with it.

Plus two ordering changes in Close/resubscribe that ARE NOT fixes for an
observed race, and say so in the test. Making b.pubsub reassignable made
Close's unlocked read of it look wrong, and resubscribe's wg.Add outside
the lock look like it could land after Close reached Wait. Both windows
turn out to be shut already by resubscribe's b.closed check, which sits
under the same acquisition as the count — reverting either fix leaves
the new Close-during-cycle test green. Kept as defence because the
invariant they lean on is three functions away, and documented so
nobody later reads them as evidence of a bug that existed.

Also corrects the metric help and two comments that said an idle cycle
"replaced the connection" when it attempts a replacement that can fail;
the deployment doc already said attempted. And the deployment doc's
rollback, frame-validation, what-to-watch and startup-log paragraphs,
all of which moved under a heading covering both buses while still
describing only the activity one.

Refs BUG-2769

* refactor(watchevents): drop an always-empty return and the branch reading it

dropCoverageIfStillIdle returned (string, bool) where the string was
never anything but empty — the reset it reports goes out through the
pending/flush path inside the lock, so the caller's `if report != ""`
was unreachable. A second reporting path that exists in the signature
and never fires is a thing a later change wires up by accident.

Refs BUG-2769

* fix(watchevents): a failed re-dial retries without re-dropping coverage

Codex round 5, on behaviour across a full Redis outage. No BLOCK; this
was its one P2 and it is real.

The probe-failure suspension does not cover this case, and the reason is
worth stating because the suspension looks like it should. Suspension
asks "did our last probe get through", and that can be YES with the
route already gone: the last successful publish stamps lastProbeOK,
Redis dies before that frame comes back, and lastSeen stays behind it.
From there both timestamps are frozen — the probe fails so nothing
stamps lastProbeOK, nothing arrives so nothing stamps lastSeen — and the
cycle's precondition stays true for the whole outage. Every pass then
dropped coverage, announced to every subscriber, and re-dialled.

Only the re-dial is owed. The second drop empties an already-empty
buffer and re-announces a hole every subscriber has been told about,
and it moves pad_watchevents_sequence_resets_total{reason="idle_timeout"}
once per cadence — so a five-minute outage read as ten incidents on the
series operators are told to alert on.

cycleIfIdle now has a retry-only arm ahead of the decision, entered when
there is no subscription at all, and the teardown clears b.pubsub /
b.subCancel so that state is representable. Clearing them also stops
Close closing an already-closed PubSub a second time.

Two tests, discriminating in OPPOSITE directions, because the obvious
fix for the noise is to suspend the pass and that would trade a noisy
outage for one the instance never returns from — retrying the dial IS
the recovery path:

  three passes with Redis away        one reset, not three
  Redis returns after a failed pass   the subscription is re-established
                                      and the counter does not move again

Matrix: removing the retry arm, making it return without retrying, and
leaving the torn-down subscription in place are each detected, the
middle one only by the recovery test.

internal/events has no equivalent defect. Its teardown deletes the
workspace's subscription entry, so its next scan finds nothing live and
abandons; recovery there runs off the request path.

Refs BUG-2769

* fix(watchevents): only one caller may install a replacement subscription

Codex round 6, verifying round 5's fix. No BLOCK; this was its P2.

Both the cycle and its new retry arm dial with the lock RELEASED, which
is deliberate — a Redis round trip under the bus's hot mutex would stall
every fan-out on the instance — so two passes can each find no
subscription and each dial one. Installing both is wrong twice over: two
receive loops would run on the SAME generation, so both accept every
frame and each notification is processed twice, and the loser's PubSub
would be untracked, closed by nothing including Close.

The install is what needs serialising, not the dial, so the loser
discards its own connection under the lock rather than the two racing to
overwrite b.pubsub.

Only the idle scanner calls this today, so this guards an invariant
rather than fixing an observed fault. Written down because the invariant
lives in a different file from the code relying on it, and because the
failure is silent duplication rather than a crash.

The test races two resubscribes through the install seam. Two details it
needed, both found by running it rather than reading it:

  the loop count is incremented INSIDE the goroutine, so sampling it
  right after the constructor returns reads zero — the first version
  did, and measured every later count against that wrong baseline. It
  waits for the loop now.

  the seam release is deferred, because without it the guard's mutation
  parks both callers in the callback, Close waits on receive loops that
  cannot start, and the detection arrives as a package-wide hang with no
  message. That is how the mutation first appeared to pass.

Also completes the idle_timeout reason in three comment/help sites that
still enumerated four reasons and said "the last two" — the same stale
count corrected in the observer contract earlier on this branch, missed
in its neighbours because I fixed the one the reviewer named instead of
grepping for the claim.

Refs BUG-2769

* test(watchevents): count installs instead of waiting for one that never comes

Codex round 7 returned no BLOCK and no P2 on the production code, and
two NITs on what round 6 added. Both are real.

The concurrency test synchronised on a WaitGroup expecting BOTH callers
to reach the install seam. Only the winner does — that is the property
under test — so in the passing case the goroutine waiting on it blocks
forever. A leak inside a test written to prove a leak does not happen is
not a shape to leave standing. An atomic the abandoning caller never
touches carries the same information and blocks nobody, and it removes
the release channel and its deferred close along with it.

The final assertion also moved off liveReceiveLoops and onto that
count. A loop starts AFTER its install, so reading the loop count can
catch a second caller's goroutine before it has begun and see the
passing value on a failing run. Both callers have returned by the time
the install count is read, so it is final. Detection over ten runs with
the guard removed: 10/10, where the loop-count version was a race
against a goroutine's first instruction.

Also softens the retry arm's log line. It said the instance receives no
notifications until an attempt succeeds, which is true for today's
single scanner and stale the moment there are two: one caller's dial can
fail while another has already installed. It now claims only what the
failing call knows.

Refs BUG-2769

* test(watchevents): hold both callers at the window, and say what that misses

Codex round 8's P2, on the test the previous commit rewrote. Starting
two goroutines from a start gate makes overlap likely and guarantees
nothing: one can finish resubscribe before the other begins, so the
window the install guard closes need never have been open.

A seam at the dial/install boundary — connection dialled, lock not yet
taken — lets both callers announce their arrival and wait for each
other. Now the window is open by construction rather than by luck, and
the test fails as a fixture if only one caller ever reaches it, instead
of passing on evidence it never gathered.

AND IT STILL DOES NOT DETECT EVERYTHING, which the test now says in
place of leaving it implied. Measured:

  guard removed entirely                        10 runs, 10 detected
  guard checked in its own acquisition, then    10 runs,  0 detected
  the lock retaken to install

The second is the regression round 8 asked about, and catching it would
mean landing the second caller inside a check-to-install gap that exists
only in the mutant — there is nothing to yield on there, and no seam can
be placed in code that is not written. So this test covers "a guard
exists", not "the guard is in the right critical section". The latter is
held by the comment at the guard and by review, and a test comment
claiming otherwise would be worth less than the honest note.

Refs BUG-2769

* fix(watchevents): make the frame seam and the cycle log tell the truth

Codex round 9 was asked whether this should merge and said hold for a
cleanup pass. Five findings, no correctness blocker, and every one of
them a claim that had stopped matching the code.

  the frame seam did not fire for every arm, though its comment said so.
  The arms that decline to act — a heartbeat, an undecodable payload, an
  unsubscribe confirmation — were `continue` statements, which skipped
  everything after the switch. A test waiting on the seam for one of
  those frames would have HUNG rather than failed, which is the worst
  way to find this out. The switch is now its own method so every arm
  ends the frame by returning, and a test drives one frame per
  publisher-reachable arm and counts three. Detected against restoring
  the skip.

  the idle-cycle warning was emitted before the revalidation that can
  abandon the cycle, so it could announce coverage ending and resumes
  answering sync_required for a subscription that was then left alone —
  a log line with no counter behind it, and an on-call hunting a bug
  that is not there. internal/events learned this at its own round 6;
  the reason did not come across with the port. Moved after the decision
  is final, still saying "attempting" to replace because the resubscribe
  can fail.

  the quiet-exit test sampled liveReceiveLoops instead of waiting for
  it, so its "the replaced loop left" assertion could be satisfied by a
  loop that never ran. Same defect fixed in the sibling concurrency test
  a commit earlier and missed here, because I looked at the test the
  reviewer named rather than at the pattern. Latent rather than
  observed: sampling survives 10 runs, so this removes a possibility.

  the probe-failure log and metric help called an errored Publish a
  failure to publish. A returned error can also mean the reply was lost
  after Redis accepted the frame, so the honest claim is that the probe
  is UNCONFIRMED. It changes no behaviour — an unconfirmed probe is not
  evidence about the receive path either, so detection suspends the same
  way — but an operator reading the counter should not be told more than
  the instance knows.

  the deployment doc said the watch stream differs in "three things" and
  listed four, the fourth being the bullet I added last round. Third
  instance of that species on this branch; the count is gone rather than
  corrected.

Refs BUG-2769

* docs(watchevents): stop one unconfirmed probe standing in for a broken path

Codex round 10 confirmed four of round 9's five fixes and held the fifth
as partial. It was right on all three residual sites.

Renaming the condition to "could not confirm" did not fix the sentences
downstream of it. The log still said silence cannot be read as a finding
"when we could not ask" — but we may well have asked, and lost only the
answer. And both the metric help and the observer contract said an
instance in this state "is also failing to deliver its own notifications
to every other instance", which is a conclusion about the outbound path
drawn from a single call that did not come back.

The inference is sound at a SUSTAINED rate and worthless at one
increment, so both now say which is which. That distinction is the whole
value of the counter to an on-call: a blip is a lost reply, a rate is a
broken path, and the same wording for both makes the first look like the
second.

No behaviour change. An unconfirmed probe suspends detection exactly as
a definite failure does, because it is not evidence about the receive
path either way.

Refs BUG-2769

* docs: sweep the BUG-2738 prose this change makes false

BUG-2738 shipped documentation that describes the watch stream as still
carrying the half-open defect. Merging this makes those sentences wrong,
and I flagged the sweep as owed twice during the groundwork and then did
not do it — the lead caught that the package said nothing about it.

Five sites, each re-read after editing rather than grepped for, because
grepping for a phrasing I chose is how I have twice verified a sweep
that had not landed:

  the residual enumeration opened "One gap remains everywhere, and a
  second remains on the watch stream only", then described one gap and
  said it was open on both. The second WAS the half-open case. Now
  states one gap, on both streams, and says where the second went.

  the half-open paragraph already said "closed on both streams" — the
  one site I had fixed — but omitted that each half is behind its own
  phase-2 flag, so a reader takes it as closed on their deployment when
  it is closed only once they turn it on.

  "A third residual" counted the item it followed. With the second gone
  the ordinal was wrong; it does not need one.

  "these two gaps" in the closing sentence, same arithmetic.

  the pad_event_subscription_cycled_total row told an operator to read
  heartbeat_phase off the startup log. There are now two such fields on
  two lines under two flags, and only one bears on that counter. It
  names the line.

No code change; suite 28/28 and lint 0 re-run because the branch is
under review and a docs commit that skips them is a commit nobody
checked.

Refs BUG-2769
2026-08-25 11:04:18 -04:00
..