Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 6 additions & 17 deletions cmd/ingestor/mqtt_watchdog_1749_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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 {
Expand Down
6 changes: 2 additions & 4 deletions cmd/ingestor/mqtt_watchdog_1810_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand Down
136 changes: 136 additions & 0 deletions cmd/ingestor/mqtt_watchdog_stop_join_test.go
Original file line number Diff line number Diff line change
@@ -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")
}
}
70 changes: 70 additions & 0 deletions cmd/ingestor/mqtt_watchdog_testhelpers_test.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package main

import (
"sync"
"testing"
"time"
)
Expand Down Expand Up @@ -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")
}
}
Loading