From b2ad17a78903395813150b2102ebc20aed64c51a Mon Sep 17 00:00:00 2001 From: efiten Date: Fri, 11 Sep 2026 10:58:56 +0200 Subject: [PATCH 1/4] test(ingestor): join the watchdog loop goroutine instead of only asking it to stop (#2003) Verified before merging: on upstream/master `go test ./cmd/ingestor -run TestMQTTStallWatchdog -count=20` fails; on this branch the same command passes. The flake blocked CI on #2000. (cherry picked from commit 02feb2a88ef05284445f022eb87aca78273fbb4c) Co-Authored-By: Claude Opus 5 --- cmd/ingestor/mqtt_watchdog_1749_test.go | 23 ++++----------- cmd/ingestor/mqtt_watchdog_1810_test.go | 6 ++-- .../mqtt_watchdog_testhelpers_test.go | 29 +++++++++++++++++++ 3 files changed, 37 insertions(+), 21 deletions(-) diff --git a/cmd/ingestor/mqtt_watchdog_1749_test.go b/cmd/ingestor/mqtt_watchdog_1749_test.go index 105c18ad2..43e407829 100644 --- a/cmd/ingestor/mqtt_watchdog_1749_test.go +++ b/cmd/ingestor/mqtt_watchdog_1749_test.go @@ -58,15 +58,8 @@ func TestMQTTStallWatchdog_EscalateOnPersistentDisconnect_1749(t *testing.T) { t.Fatalf("setup: %v", err) } - tick := make(chan time.Time) - done := make(chan struct{}) - defer close(done) - - exited := make(chan struct{}) - go func() { - runLivenessWatchdogLoop(tick, done, threshold, func(args ...any) {}) - close(exited) - }() + tick, stopLoop := startWatchdogTestLoop(t, threshold, func(args ...any) {}) + defer stopLoop() // Feed ticks spanning > (multiplier × threshold) of wall clock so // the escalation path fires. We control the `now` parameter by @@ -114,10 +107,8 @@ func TestMQTTStallWatchdog_DisconnectedEscalationThrottled_1749(t *testing.T) { t.Fatalf("setup: %v", err) } - tick := make(chan time.Time) - done := make(chan struct{}) - defer close(done) - go runLivenessWatchdogLoop(tick, done, threshold, func(args ...any) {}) + tick, stopLoop := startWatchdogTestLoop(t, threshold, func(args ...any) {}) + defer stopLoop() base := time.Now() // Pre-stamp DisconnectedSinceUnix so that the first tick is @@ -259,10 +250,8 @@ func TestMQTTStallWatchdog_LastTickUnixExposed_1749(t *testing.T) { watchdogLastTickUnix.Store(before) }) - tick := make(chan time.Time) - done := make(chan struct{}) - defer close(done) - go runLivenessWatchdogLoop(tick, done, time.Minute, func(args ...any) {}) + tick, stopLoop := startWatchdogTestLoop(t, time.Minute, func(args ...any) {}) + defer stopLoop() stamp := time.Now().Add(48 * time.Hour) // guaranteed > before select { diff --git a/cmd/ingestor/mqtt_watchdog_1810_test.go b/cmd/ingestor/mqtt_watchdog_1810_test.go index aeefd4a0f..b68fa75d3 100644 --- a/cmd/ingestor/mqtt_watchdog_1810_test.go +++ b/cmd/ingestor/mqtt_watchdog_1810_test.go @@ -212,10 +212,8 @@ func TestWatchdog_EscalationWarnThrottled_1810(t *testing.T) { } } - tick := make(chan time.Time) - done := make(chan struct{}) - defer close(done) - go runLivenessWatchdogLoop(tick, done, threshold, emit) + tick, stopLoop := startWatchdogTestLoop(t, threshold, emit) + defer stopLoop() base := time.Now() // First tick: stamps DisconnectedSinceUnix (no escalation yet). diff --git a/cmd/ingestor/mqtt_watchdog_testhelpers_test.go b/cmd/ingestor/mqtt_watchdog_testhelpers_test.go index 348c2f2d1..1db9b99a5 100644 --- a/cmd/ingestor/mqtt_watchdog_testhelpers_test.go +++ b/cmd/ingestor/mqtt_watchdog_testhelpers_test.go @@ -1,6 +1,7 @@ package main import ( + "sync" "testing" "time" ) @@ -45,3 +46,31 @@ func sendTickOrFail(t *testing.T, tick chan<- time.Time, stamp time.Time, timeou t.Fatalf("%s: tick blocked after %s — loop dead?", label, timeout) } } + +// startWatchdogTestLoop is setupWatchdogTestLoop for the common case: a test +// that wants the loop running for its own duration and nothing more. The +// returned stop closes done AND waits for the loop goroutine to return. It is +// safe to call more than once. +// +// Waiting is the part that must not be skipped. Closing done only asks the +// loop to stop; the goroutine can still be inside a scan, and that scan walks +// the package-level livenessRegistry. The next test registers its own source +// there, so a loop that has not returned yet will process that source on a +// tick carrying the previous test's fabricated clock. +// +// That is measured, not theoretical. With four call sites closing done and +// walking away, TestMQTTStallWatchdog_DisconnectedEscalationThrottled_1749 +// failed 2 to 3 times per 20 runs of the watchdog tests and never once in 50 +// runs on its own. Debug tracing showed two loops reaching maybeForceReconnect +// for the same source with tick clocks 420s apart, both reading +// LastForceReconnectUnix as 0 before either wrote it, so the throttle the test +// asserts on let both through. +func startWatchdogTestLoop(t *testing.T, threshold time.Duration, emit func(...any)) (tick chan time.Time, stop func()) { + t.Helper() + tick, done, exited := setupWatchdogTestLoop(t, threshold, emit) + var once sync.Once + return tick, func() { + once.Do(func() { close(done) }) + <-exited + } +} From 3729e12a88cbda8b89476e757142f54fe5fe3400 Mon Sep 17 00:00:00 2001 From: Dennis Jakobsen Date: Fri, 18 Sep 2026 05:43:45 +0200 Subject: [PATCH 2/4] test(ingestor): prove the watchdog test stop joins the loop Add TestStartWatchdogTestLoop_StopJoinsLoop. It holds the watchdog loop inside a tick with a blocking emit callback and checks that stop does not return until the callback is released and the loop has exited. Ordering uses channels only. The test waits for done to be closed, so "stop has not returned" cannot just mean the stop goroutine was not scheduled yet. Every goroutine is joined on every path, including failures. With a close-only stop the test fails; with the join it passes. Factor the stop closure into joiningWatchdogStop so the test can hold done and exited. startWatchdogTestLoop behaves exactly as before. Correct the helper comment. The previous commit message and comment said two loops both read LastForceReconnectUnix as 0 before either wrote it. The failure traced locally was different: the loop from EscalateOnPersistentDisconnect_1749, whose clock was 420s ahead, was still running when the throttle test started. It read the non-zero stamp that the throttle test's own loop had just written, measured more than forceReconnectThrottle against its own clock, and forced a second reconnect. Test-only. No production code changes. Co-Authored-By: Claude Opus 5 --- cmd/ingestor/mqtt_watchdog_stop_join_test.go | 136 ++++++++++++++++++ .../mqtt_watchdog_testhelpers_test.go | 23 ++- 2 files changed, 153 insertions(+), 6 deletions(-) create mode 100644 cmd/ingestor/mqtt_watchdog_stop_join_test.go diff --git a/cmd/ingestor/mqtt_watchdog_stop_join_test.go b/cmd/ingestor/mqtt_watchdog_stop_join_test.go new file mode 100644 index 000000000..2fa759aea --- /dev/null +++ b/cmd/ingestor/mqtt_watchdog_stop_join_test.go @@ -0,0 +1,136 @@ +package main + +import ( + "sync" + "sync/atomic" + "testing" + "time" +) + +// TestStartWatchdogTestLoop_StopJoinsLoop guards the test harness, not +// production code. The stop returned by startWatchdogTestLoop must not return +// while the loop goroutine is still inside a tick, because the caller restores +// livenessRegistry right after it and the next test registers its own sources +// there. A stop that only closes done lets the old loop finish that tick, and +// possibly scan the next test's registry, with its own fabricated clock. That +// is how TestMQTTStallWatchdog_DisconnectedEscalationThrottled_1749 saw two +// forced reconnects. +// +// The loop is held inside a tick by an emit callback that blocks until the +// test releases it. Everything is ordered by channels: +// +// - emitEntered proves the loop is inside the tick before stop is called. +// - done being closed proves the stop goroutine has actually run stop and +// got past close(done). Without this, "stop has not returned" could only +// mean the goroutine had not been scheduled yet. +// - stopReturned is then checked while emit is still blocked. A close-only +// stop returns right after close(done) and is caught here. A joining stop +// cannot return at all until release, so nothing timing-based can make +// the correct helper fail. notReturnedWindow only bounds how long a +// close-only stop gets to show itself. +// - As a second check, the stop goroutine records whether exited was +// already closed at the moment stop returned. A joining stop always sees +// it closed. +// +// The time.After bounds are safety nets that turn a hang into a failure. +// They do not order anything. +func TestStartWatchdogTestLoop_StopJoinsLoop(t *testing.T) { + defer snapshotAndResetRegistry(t)() + + const ( + safetyTimeout = 5 * time.Second + notReturnedWindow = 100 * time.Millisecond + ) + + // Connected but silent for 10m against a 1m threshold: the first tick + // classifies it as LivenessStalled and calls emit. + s := &SourceLivenessState{ + Tag: "stop-join-harness", + Broker: "tcp://x:1883", + IsConnectedFn: func() bool { return true }, + } + atomic.StoreInt64(&s.LastMessageUnix, time.Now().Add(-10*time.Minute).Unix()) + if err := registerLivenessState(s); err != nil { + t.Fatalf("setup: %v", err) + } + + emitEntered := make(chan struct{}) + release := make(chan struct{}) + var enteredOnce, releaseOnce sync.Once + releaseEmit := func() { releaseOnce.Do(func() { close(release) }) } + emit := func(...any) { + enteredOnce.Do(func() { close(emitEntered) }) + <-release + } + + tick, done, exited := setupWatchdogTestLoop(t, time.Minute, emit) + stop := joiningWatchdogStop(done, exited) + // Deferred calls run last-in first-out. On every path, including + // t.Fatal, emit is released first, then the stop goroutine (if started) + // is joined, then stop joins the loop, then the registry is restored. + defer stop() + + sendTickOrFail(t, tick, time.Now(), safetyTimeout, "tick into stalled source") + + var stopStarted bool + stopReturned := make(chan struct{}) + var exitedWhenStopReturned atomic.Bool + defer func() { + if !stopStarted { + return + } + select { + case <-stopReturned: + case <-time.After(safetyTimeout): + t.Errorf("stop goroutine still running %s after release", safetyTimeout) + } + }() + defer releaseEmit() + + select { + case <-emitEntered: + case <-time.After(safetyTimeout): + t.Fatalf("loop never called emit within %s", safetyTimeout) + } + + stopStarted = true + go func() { + defer close(stopReturned) + stop() + select { + case <-exited: + exitedWhenStopReturned.Store(true) + default: + } + }() + + select { + case <-done: + case <-time.After(safetyTimeout): + t.Fatalf("stop did not close done within %s", safetyTimeout) + } + + select { + case <-stopReturned: + t.Fatal("stop returned while the loop was still blocked inside a tick; it must wait for the loop to exit") + case <-exited: + t.Fatal("loop exited while emit was still blocked") + case <-time.After(notReturnedWindow): + } + + releaseEmit() + + select { + case <-exited: + case <-time.After(safetyTimeout): + t.Fatalf("loop did not exit within %s after emit was released", safetyTimeout) + } + select { + case <-stopReturned: + case <-time.After(safetyTimeout): + t.Fatalf("stop did not return within %s after the loop exited", safetyTimeout) + } + if !exitedWhenStopReturned.Load() { + t.Fatal("stop returned before the loop goroutine had exited") + } +} diff --git a/cmd/ingestor/mqtt_watchdog_testhelpers_test.go b/cmd/ingestor/mqtt_watchdog_testhelpers_test.go index 1db9b99a5..3e398654c 100644 --- a/cmd/ingestor/mqtt_watchdog_testhelpers_test.go +++ b/cmd/ingestor/mqtt_watchdog_testhelpers_test.go @@ -60,16 +60,27 @@ func sendTickOrFail(t *testing.T, tick chan<- time.Time, stamp time.Time, timeou // // That is measured, not theoretical. With four call sites closing done and // walking away, TestMQTTStallWatchdog_DisconnectedEscalationThrottled_1749 -// failed 2 to 3 times per 20 runs of the watchdog tests and never once in 50 -// runs on its own. Debug tracing showed two loops reaching maybeForceReconnect -// for the same source with tick clocks 420s apart, both reading -// LastForceReconnectUnix as 0 before either wrote it, so the throttle the test -// asserts on let both through. +// failed intermittently when run with the other watchdog tests and not when +// run on its own. A trace of one failure showed the loop from +// TestMQTTStallWatchdog_EscalateOnPersistentDisconnect_1749, whose last tick +// carried a clock 420s ahead, still running after its test returned. The +// throttle test's own loop had already forced a reconnect and stamped +// LastForceReconnectUnix. The old loop then read that non-zero stamp, measured +// it against its own clock, saw more than forceReconnectThrottle elapse and +// forced a second reconnect. The throttle itself was correct; it was given two +// clocks for one source. func startWatchdogTestLoop(t *testing.T, threshold time.Duration, emit func(...any)) (tick chan time.Time, stop func()) { t.Helper() tick, done, exited := setupWatchdogTestLoop(t, threshold, emit) + return tick, joiningWatchdogStop(done, exited) +} + +// joiningWatchdogStop returns the stop function used by startWatchdogTestLoop: +// close done once, then wait for exited. It is separate only so +// TestStartWatchdogTestLoop_StopJoinsLoop can hold done and exited itself. +func joiningWatchdogStop(done, exited chan struct{}) func() { var once sync.Once - return tick, func() { + return func() { once.Do(func() { close(done) }) <-exited } From 29b3573e17a2a98161027c36a872a59a1e12a191 Mon Sep 17 00:00:00 2001 From: Openclaw Date: Fri, 18 Sep 2026 05:58:43 +0200 Subject: [PATCH 3/4] test(ingestor): assert startWatchdogTestLoop's stop() is idempotent (#2003) startWatchdogTestLoop's doc comment claims "It is safe to call more than once". joiningWatchdogStop (factored out by the StopJoinsLoop commit) guards close(done) with sync.Once, so a second call skips that close and falls straight to <-exited, which returns immediately once exited is closed. That claim was never exercised. Call stop() twice and assert the second call returns within 1s, matching this package's existing timeout/select idiom. This is a different guarantee than TestStartWatchdogTestLoop_ StopJoinsLoop (previous commit): that test proves a single stop() call waits for the loop to actually exit, not just for done to close. This test proves a second, concurrent-or-sequential stop() call does not hang. Independently authored, rebuilt on top of the StopJoinsLoop commit after that commit landed and changed the tail of this file; not cherry-picked, no unrelated changes carried over. Verified: both new tests together, watchdog/liveness/asyncemit group (count=5, 180/180 pass, 0 fail), same group under -race (count=5, 180/180 pass, 0 data races), gofmt clean, git diff --check clean, no production code touched. Full cmd/ingestor package has two pre-existing issues, independently reproduced on the clean base (3729e12a) with this change stashed out: TestPruneOldNeighborMetrics fails 5/5 (dated-fixture time-bomb) and TestBackfillTxLastSeen_ ResolvesFromMaxObservationTimestamp is flaky, 2/5 fails with no watchdog code involved (async-migration timing flake). Neither is touched or caused by this commit. Co-Authored-By: Claude Opus 5 Co-Authored-By: Claude Sonnet 5 --- .../mqtt_watchdog_testhelpers_test.go | 24 +++++++++++++++++++ 1 file changed, 24 insertions(+) diff --git a/cmd/ingestor/mqtt_watchdog_testhelpers_test.go b/cmd/ingestor/mqtt_watchdog_testhelpers_test.go index 3e398654c..abe03871e 100644 --- a/cmd/ingestor/mqtt_watchdog_testhelpers_test.go +++ b/cmd/ingestor/mqtt_watchdog_testhelpers_test.go @@ -85,3 +85,27 @@ func joiningWatchdogStop(done, exited chan struct{}) func() { <-exited } } + +// TestStartWatchdogTestLoop_StopIsIdempotent turns the "safe to call more +// than once" claim in startWatchdogTestLoop's doc comment above into an +// executable assertion. joiningWatchdogStop's returned func uses sync.Once +// to guard close(done), so a second call skips that close and falls +// straight to <-exited — which returns immediately once exited is closed, +// on every subsequent read. If either the once-guard or that closed-channel +// behaviour ever regressed, a second stop() call would hang instead of +// returning immediately, and this test would time out. +func TestStartWatchdogTestLoop_StopIsIdempotent(t *testing.T) { + _, stop := startWatchdogTestLoop(t, time.Minute, func(args ...any) {}) + stopped := make(chan struct{}) + go func() { + defer close(stopped) + stop() + stop() + }() + + select { + case <-stopped: + case <-time.After(time.Second): + t.Fatal("calling stop twice blocked; stop must remain idempotent") + } +} From adf429e628e410da962ad40ac71e630a6977e42d Mon Sep 17 00:00:00 2001 From: Dennis Jakobsen Date: Fri, 18 Sep 2026 06:39:43 +0200 Subject: [PATCH 4/4] test(ingestor): fix TestStartWatchdogTestLoop_StopIsIdempotent comment The comment claimed that losing the sync.Once guard would make a second stop() call hang, so the test would catch it via the 1s timeout. Verified by mutation: removing sync.Once instead makes the second close(done) call panic with "close of closed channel", which fails the test immediately, not by timing out. The 1s timeout only catches a regression in <-exited's closed-channel read. Also states plainly that the two stop() calls are sequential, from the same goroutine, not concurrent. Comment-only. No change to the test's assertions, to StopJoinsLoop, or to production code. Co-Authored-By: Claude Sonnet 5 --- cmd/ingestor/mqtt_watchdog_testhelpers_test.go | 18 ++++++++++++------ 1 file changed, 12 insertions(+), 6 deletions(-) diff --git a/cmd/ingestor/mqtt_watchdog_testhelpers_test.go b/cmd/ingestor/mqtt_watchdog_testhelpers_test.go index abe03871e..670c4e1f9 100644 --- a/cmd/ingestor/mqtt_watchdog_testhelpers_test.go +++ b/cmd/ingestor/mqtt_watchdog_testhelpers_test.go @@ -88,12 +88,18 @@ func joiningWatchdogStop(done, exited chan struct{}) func() { // TestStartWatchdogTestLoop_StopIsIdempotent turns the "safe to call more // than once" claim in startWatchdogTestLoop's doc comment above into an -// executable assertion. joiningWatchdogStop's returned func uses sync.Once -// to guard close(done), so a second call skips that close and falls -// straight to <-exited — which returns immediately once exited is closed, -// on every subsequent read. If either the once-guard or that closed-channel -// behaviour ever regressed, a second stop() call would hang instead of -// returning immediately, and this test would time out. +// executable assertion. It calls stop() twice, sequentially, from the same +// goroutine. joiningWatchdogStop's returned func uses sync.Once to guard +// close(done), so the second call skips that close and falls straight to +// <-exited — which returns immediately once exited is closed, on every +// subsequent read. +// +// If the once-guard regressed, a second close(done) call panics with +// "close of closed channel" (verified: this is what removing sync.Once +// actually produces) and fails the test that way rather than by timing +// out. If instead the closed-channel behaviour of <-exited regressed so a +// second call blocked, that second call would hang and the select below +// would catch it via the 1s timeout. func TestStartWatchdogTestLoop_StopIsIdempotent(t *testing.T) { _, stop := startWatchdogTestLoop(t, time.Minute, func(args ...any) {}) stopped := make(chan struct{})