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_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 348c2f2d1..670c4e1f9 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,72 @@ 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 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 func() { + once.Do(func() { close(done) }) + <-exited + } +} + +// TestStartWatchdogTestLoop_StopIsIdempotent turns the "safe to call more +// than once" claim in startWatchdogTestLoop's doc comment above into an +// 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{}) + 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") + } +}