From fc0c068f5a93a17de6e6107094c75b5f55fa63ee Mon Sep 17 00:00:00 2001 From: dborup Date: Sun, 4 Oct 2026 10:37:58 +0200 Subject: [PATCH 1/4] test(ingestor): watchdog retry noise and force-reconnect at shutdown (#102, #103) Red on master: - #102: a real paho client in its initial ConnectRetry loop (status connecting). Five force-reconnects log "Connect() failed" five times, though paho keeps retrying and connects once the broker is up. - #103: newAsyncEmit's emit after stop panics with "send on closed channel", also when racing stop (and -race reports a data race). The issue's scenario: a force-reconnect whose ForceReconnectFn blocks, the watchdog stops, then the reconnect returns and its goroutine emits "reconnect attempt issued": the test binary panics. Green on master, pinning what must stay an error: Connect()'s status error while paho is disconnecting with no reconnect pending, and any other Connect() error. fakeClient.IsConnected now returns a field. Relates to #102, #103 Co-Authored-By: Claude Opus 5.5 --- .../mqtt_force_reconnect_race_test.go | 8 +- cmd/ingestor/mqtt_watchdog_102_103_test.go | 221 ++++++++++++++++++ 2 files changed, 226 insertions(+), 3 deletions(-) create mode 100644 cmd/ingestor/mqtt_watchdog_102_103_test.go diff --git a/cmd/ingestor/mqtt_force_reconnect_race_test.go b/cmd/ingestor/mqtt_force_reconnect_race_test.go index 8697cd3ea..76df03e21 100644 --- a/cmd/ingestor/mqtt_force_reconnect_race_test.go +++ b/cmd/ingestor/mqtt_force_reconnect_race_test.go @@ -25,13 +25,15 @@ func (t *fakeToken) WaitTimeout(time.Duration) bool { return true } func (t *fakeToken) Done() <-chan struct{} { ch := make(chan struct{}); close(ch); return ch } func (t *fakeToken) Error() error { return t.err } -// fakeClient implements mqtt.Client, recording calls to the three methods +// fakeClient implements mqtt.Client, recording calls to the methods // buildForceReconnectFn actually uses (IsConnectionOpen, Disconnect, -// Connect). Every other method panics — buildForceReconnectFn must never +// Connect; IsConnected after a Connect() error, #102). Every other method +// panics — buildForceReconnectFn must never // touch subscriptions, publishes, or options, so a call there indicates the // fix drifted from its intended scope. type fakeClient struct { isConnectionOpen bool + isConnected bool connectErr error disconnectCalled bool @@ -39,7 +41,7 @@ type fakeClient struct { callOrder []string } -func (c *fakeClient) IsConnected() bool { panic("not used by buildForceReconnectFn") } +func (c *fakeClient) IsConnected() bool { return c.isConnected } func (c *fakeClient) IsConnectionOpen() bool { c.callOrder = append(c.callOrder, "IsConnectionOpen") return c.isConnectionOpen diff --git a/cmd/ingestor/mqtt_watchdog_102_103_test.go b/cmd/ingestor/mqtt_watchdog_102_103_test.go new file mode 100644 index 000000000..313174a0e --- /dev/null +++ b/cmd/ingestor/mqtt_watchdog_102_103_test.go @@ -0,0 +1,221 @@ +package main + +import ( + "errors" + "fmt" + "runtime" + "strings" + "sync" + "sync/atomic" + "testing" + "time" + + mqtt "github.com/eclipse/paho.mqtt.golang" +) + +// Issue #102: while paho is in its initial ConnectRetry loop (status +// connecting), the watchdog's force-reconnect must not log the expected +// Connect() status error as a failure. Issue #103: a force-reconnect still in +// flight when the watchdog stops must not panic with "send on closed channel". + +const ( + connectFailedLine102 = "WATCHDOG force-reconnect Connect() failed" + retryInProgress102 = "retry already in progress" +) + +// newNeverConnectedTestClient builds the client as main does, against a +// broker that is down from the start, so paho stays in its initial +// ConnectRetry loop (status connecting). It waits for the first CONNECT. +func newNeverConnectedTestClient(t *testing.T, b *forceReconnectTestBroker, tag string, retryInterval time.Duration) mqtt.Client { + t.Helper() + b.goDown() + opts := buildMQTTOpts(MQTTSource{Broker: b.url(), Name: tag}). + SetConnectTimeout(500 * time.Millisecond). + SetWriteTimeout(500 * time.Millisecond). + SetMaxReconnectInterval(forceReconnectTestRetryInterval). + SetConnectRetryInterval(retryInterval) + client := mqtt.NewClient(opts) + client.Connect() + t.Cleanup(func() { client.Disconnect(0) }) + waitForConnects(t, b, 0, "test setup: initial connect") + return client +} + +// #102: five triggers during the initial ConnectRetry loop (the #29 review +// measured 5/5). Connect() returns paho's errStatusMustBeDisconnected each +// time; paho keeps retrying, so this is not a failure. The client must not +// log "Connect() failed", must log that a retry is in progress, and paho's +// loop must keep going and connect once the broker is up. +func TestForceReconnect_RealPaho_InitialRetryLoopIsNotAConnectFailure_102(t *testing.T) { + b := newForceReconnectTestBroker(t) + logs := captureLog118(t) + client := newNeverConnectedTestClient(t, b, "force-reconnect-initial", 200*time.Millisecond) + + // Premise: status connecting looks like reconnecting to the watchdog. + if !client.IsConnected() || client.IsConnectionOpen() { + t.Fatalf("test setup: IsConnected=%v IsConnectionOpen=%v, want true/false (initial ConnectRetry loop)", client.IsConnected(), client.IsConnectionOpen()) + } + fn := buildForceReconnectFn(client, "force-reconnect-initial") + for i := 0; i < 5; i++ { + fn() + } + if n := strings.Count(logs.String(), connectFailedLine102); n != 0 { + t.Errorf("%d %q lines while paho was in its initial retry loop:\n%s", n, connectFailedLine102, logs) + } + if n := strings.Count(logs.String(), retryInProgress102); n != 5 { + t.Errorf("%d %q lines, want 5:\n%s", n, retryInProgress102, logs) + } + + n, _ := b.counts() + waitForConnects(t, b, n, "paho's retry loop after the triggers") + b.goUp() + if !pollUntil(forceReconnectTestDeadline, client.IsConnectionOpen) { + t.Fatal("client did not connect after the broker came back") + } +} + +// #102: the same paho error is a genuine failure when nothing will retry. +// Disconnect() during the initial retry loop moves paho to disconnecting +// without a reconnect, and blocks there until the retry sleep ends. A +// force-reconnect in that window gets errStatusMustBeDisconnected too, but +// IsConnected() is false: no attempt will follow, so it stays an error. +func TestForceReconnect_RealPaho_ConnectErrorWhileDisconnectingIsLogged_102(t *testing.T) { + b := newForceReconnectTestBroker(t) + logs := captureLog118(t) + client := newNeverConnectedTestClient(t, b, "force-reconnect-disconnecting", 3*time.Second) + + go client.Disconnect(0) + if !pollUntil(time.Second, func() bool { return !client.IsConnected() }) { + t.Fatal("test setup: client did not leave status connecting") + } + buildForceReconnectFn(client, "force-reconnect-disconnecting")() + + if n := strings.Count(logs.String(), connectFailedLine102); n != 1 { + t.Errorf("%d %q lines, want 1:\n%s", n, connectFailedLine102, logs) + } + if n := strings.Count(logs.String(), retryInProgress102); n != 0 { + t.Errorf("a Connect() error with no retry pending was logged as %q:\n%s", retryInProgress102, logs) + } +} + +// #102: only paho's status error is recognised. Any other Connect() error +// is logged as a failure, whatever IsConnected() reports. +func TestBuildForceReconnectFn_OtherConnectErrorIsLogged_102(t *testing.T) { + logs := captureLog118(t) + c := &fakeClient{isConnected: true, connectErr: errors.New("network unreachable")} + buildForceReconnectFn(c, "other-error-source")() + if n := strings.Count(logs.String(), connectFailedLine102+": network unreachable"); n != 1 { + t.Errorf("Connect() error not logged as a failure:\n%s", logs) + } + if strings.Count(logs.String(), retryInProgress102) != 0 { + t.Errorf("an unrelated Connect() error was logged as %q:\n%s", retryInProgress102, logs) + } +} + +// #103: emit after stop is a no-op. The force-reconnect goroutine can emit +// after the watchdog has stopped; that must not panic. +func TestNewAsyncEmit_EmitAfterStopDoesNotPanic_103(t *testing.T) { + var delivered atomic.Int64 + emit, stop := newAsyncEmit(func(...any) { delivered.Add(1) }) + emit("before stop") + stop() + func() { + defer func() { + if r := recover(); r != nil { + t.Fatalf("emit after stop panicked: %v", r) + } + }() + emit("after stop") + }() + if got := delivered.Load(); got != 1 { + t.Fatalf("delivered %d lines, want 1 (the line before stop)", got) + } +} + +// #103: emits racing stop neither panic nor race (run with -race). +func TestNewAsyncEmit_ConcurrentEmitDuringStop_103(t *testing.T) { + emit, stop := newAsyncEmit(func(...any) {}) + start := make(chan struct{}) + panics := make(chan any, 8) + var wg sync.WaitGroup + for g := 0; g < 8; g++ { + wg.Add(1) + go func() { + defer wg.Done() + defer func() { + if r := recover(); r != nil { + panics <- r + } + }() + <-start + for i := 0; i < 2000; i++ { + emit("line", i) + } + }() + } + close(start) + stop() + wg.Wait() + close(panics) + for r := range panics { + t.Fatalf("emit racing stop panicked: %v", r) + } +} + +// #103, the scenario from the issue: the production watchdog fires a +// force-reconnect whose ForceReconnectFn blocks, the watchdog is stopped +// (SIGTERM path), then ForceReconnectFn returns and its goroutine emits +// "reconnect attempt issued". No panic; stop does not wait for the blocked +// reconnect. +func TestRunLivenessWatchdog_ForceReconnectInFlightAtShutdown_103(t *testing.T) { + defer snapshotAndResetRegistry(t)() + entered := make(chan struct{}, 1) + release := make(chan struct{}) + s := &SourceLivenessState{ + Tag: "shutdown-103", + Broker: "test", + IsConnectedFn: func() bool { return true }, + ForceReconnectFn: func() { + select { + case entered <- struct{}{}: + default: + } + <-release + }, + } + atomic.StoreInt64(&s.LastMessageUnix, time.Now().Add(-time.Hour).Unix()) + if err := registerLivenessState(s); err != nil { + t.Fatal(err) + } + + stop := runLivenessWatchdog(5*time.Millisecond, time.Second) + select { + case <-entered: + case <-time.After(2 * time.Second): + stop() + close(release) + t.Fatal("the watchdog never called ForceReconnectFn") + } + stopped := make(chan struct{}) + go func() { + stop() + close(stopped) + }() + select { + case <-stopped: + case <-time.After(2 * time.Second): + close(release) + t.Fatal("stop() waited for a blocked ForceReconnectFn") + } + close(release) + if !pollUntil(2*time.Second, func() bool { return !goroutineRunning("maybeForceReconnect") }) { + t.Fatal("the force-reconnect goroutine did not finish") + } +} + +// goroutineRunning reports whether any goroutine's stack mentions fn. +func goroutineRunning(fn string) bool { + buf := make([]byte, 1<<20) + n := runtime.Stack(buf, true) + return strings.Contains(string(buf[:n]), fmt.Sprintf(".%s.", fn)) +} From 4db36ab870f348f1abd14d8de7861f19db060b5c Mon Sep 17 00:00:00 2001 From: dborup Date: Sun, 4 Oct 2026 10:41:15 +0200 Subject: [PATCH 2/4] fix(ingestor): quiet watchdog retry noise and make force-reconnect shutdown-safe (#102, #103) #102: while paho is in its initial ConnectRetry loop (status connecting), or disconnecting with a reconnect pending, Connect() returns paho's errStatusMustBeDisconnected and paho keeps retrying on its own (paho v1.5.0 client.go Connect, status.go Connecting). buildForceReconnectFn now recognises that error together with IsConnected()==true and logs "retry already in progress in paho" instead of "Connect() failed". The same error with no retry pending (disconnecting after a Disconnect), and any other error, are still logged as failures. The doc comment no longer claims Connect() is a no-op in every retrying state. #103: the force-reconnect goroutine is not joined at shutdown, so its "reconnect attempt issued" emit could run after stop had closed the queue and panic with "send on closed channel". newAsyncEmit's emit now checks a stopped flag under a mutex and drops lines after stop. stop still does not wait for a blocked ForceReconnectFn. No change in normal operation: before stop, emit queues or counts a drop exactly as before. Relates to #102, #103 Co-Authored-By: Claude Opus 5.5 --- cmd/ingestor/main.go | 39 ++++++++++++++++--- .../mqtt_force_reconnect_race_test.go | 3 +- cmd/ingestor/mqtt_watchdog.go | 21 ++++++++-- 3 files changed, 53 insertions(+), 10 deletions(-) diff --git a/cmd/ingestor/main.go b/cmd/ingestor/main.go index 9f98b23b0..c6fbc3ad5 100644 --- a/cmd/ingestor/main.go +++ b/cmd/ingestor/main.go @@ -568,9 +568,18 @@ func buildMQTTOpts(source MQTTSource) *mqtt.ClientOptions { // Disconnect then Connect; Disconnecting() does not need to wait on any // in-flight retry loop from status connected, so it completes well within // the 250ms quiesce) from "paho is already retrying on its own" (must NOT -// call Disconnect; Connect() alone is a safe no-op per paho when a retry is -// already under way, and properly starts a fresh attempt on the rare -// occasion status has actually settled to disconnected). +// call Disconnect). +// +// What Connect() then does depends on paho's status (paho.mqtt.golang +// v1.5.0, client.go Connect and status.go Connecting): +// - reconnecting (AutoReconnect loop): a no-op that returns a success token; +// - connecting (the initial ConnectRetry loop) or disconnecting with a +// reconnect pending: an error token, errStatusMustBeDisconnected, while +// paho's own loop keeps retrying. Logged as info, not as a failure (#102); +// - disconnected (paho gave up): starts a fresh attempt. +// +// The same status error with no retry pending (disconnecting after a +// Disconnect) is a genuine failure and is logged as one. func buildForceReconnectFn(client mqtt.Client, tag string) func() { return func() { if client.IsConnectionOpen() { @@ -581,12 +590,32 @@ func buildForceReconnectFn(client mqtt.Client, tag string) func() { // retrying, treated as a safe no-op" success case — only a genuine // fresh connection attempt leaves the token pending in the // background, and we must not block this call on that. - if token := client.Connect(); token.Error() != nil { - log.Printf("MQTT [%s] WATCHDOG force-reconnect Connect() failed: %v", tag, token.Error()) + err := client.Connect().Error() + switch { + case err == nil: + case connectRetryInProgress(client, err): + log.Printf("MQTT [%s] WATCHDOG force-reconnect: retry already in progress in paho, no new attempt needed", tag) + default: + log.Printf("MQTT [%s] WATCHDOG force-reconnect Connect() failed: %v", tag, err) } } } +// pahoErrStatusMustBeDisconnected is the text of paho's unexported +// errStatusMustBeDisconnected (paho.mqtt.golang v1.5.0 status.go), which +// Connect() returns when status is connecting or disconnecting. The paho +// test TestForceReconnect_RealPaho_InitialRetryLoopIsNotAConnectFailure_102 +// fails if an upgrade changes it. +const pahoErrStatusMustBeDisconnected = "status can only transition to connecting from disconnected" + +// connectRetryInProgress reports whether a Connect() error only means paho is +// already retrying (#102). IsConnected() is true in status connecting only +// with ConnectRetry, and in disconnecting only when a reconnect will follow; +// buildMQTTOpts sets ConnectRetry and AutoReconnect. +func connectRetryInProgress(client mqtt.Client, err error) bool { + return err.Error() == pahoErrStatusMustBeDisconnected && client.IsConnected() +} + func handleMessage(store *Store, tag string, source MQTTSource, m mqtt.Message, channelKeys map[string]string, regionKeys map[string][]byte, cfg *Config) { // Liveness watchdog (#1212): record receipt before any processing so a // slow handler still counts as "source is alive". Cheap atomic store. diff --git a/cmd/ingestor/mqtt_force_reconnect_race_test.go b/cmd/ingestor/mqtt_force_reconnect_race_test.go index 76df03e21..c61f3bed7 100644 --- a/cmd/ingestor/mqtt_force_reconnect_race_test.go +++ b/cmd/ingestor/mqtt_force_reconnect_race_test.go @@ -128,7 +128,8 @@ func TestBuildForceReconnectFn_SkipsDisconnectWhenAlreadyRetrying(t *testing.T) // class this whole bug hinged on — Connect() called while paho's status is // transitionally "disconnecting"). The old code discarded the returned token // entirely; the fix must surface it via log output instead of silently -// dropping it. +// dropping it. IsConnected()==false: no reconnect is pending, so it is a +// genuine failure (#102 logs it as info only while paho is retrying). func TestBuildForceReconnectFn_LogsConnectError(t *testing.T) { c := &fakeClient{isConnectionOpen: false, connectErr: errors.New("status can only transition to connecting from disconnected")} fn := buildForceReconnectFn(c, "erroring-source") diff --git a/cmd/ingestor/mqtt_watchdog.go b/cmd/ingestor/mqtt_watchdog.go index 96ed25459..e39a3335b 100644 --- a/cmd/ingestor/mqtt_watchdog.go +++ b/cmd/ingestor/mqtt_watchdog.go @@ -111,12 +111,16 @@ func WatchdogLogDropCount() int64 { // watchdogLogDropCount and returns immediately rather than waiting for // the drain goroutine to catch up. // -// stop closes the queue (no further sends may be attempted after -// calling stop — see runLivenessWatchdog's shutdown ordering) and blocks -// until the drain goroutine has flushed everything already queued. +// stop closes the queue and blocks until the drain goroutine has flushed +// everything already queued. emit after stop is a no-op (#103): the +// force-reconnect goroutine (maybeForceReconnect) is not joined at +// shutdown and can emit "reconnect attempt issued" after stop. The mutex +// only orders emit against the close; nothing blocks while holding it. func newAsyncEmit(realEmit func(...any)) (emit func(...any), stop func()) { queue := make(chan []any, asyncEmitQueueSize) drained := make(chan struct{}) + var mu sync.Mutex + stopped := false go func() { defer close(drained) for args := range queue { @@ -124,6 +128,11 @@ func newAsyncEmit(realEmit func(...any)) (emit func(...any), stop func()) { } }() emit = func(args ...any) { + mu.Lock() + defer mu.Unlock() + if stopped { + return + } select { case queue <- args: default: @@ -131,7 +140,10 @@ func newAsyncEmit(realEmit func(...any)) (emit func(...any), stop func()) { } } stop = func() { + mu.Lock() + stopped = true close(queue) + mu.Unlock() <-drained } return emit, stop @@ -672,7 +684,8 @@ func maybeForceReconnect(s *SourceLivenessState, now time.Time, emit func(...any // client.Disconnect(250) which blocks up to 250ms, then // client.Connect() which can block on the connect timeout. The // watchdog goroutine must not stall a per-tick scan over a single - // slow source. + // slow source. Shutdown does not wait for it, so its emit may come + // after the watchdog stopped; newAsyncEmit drops that line (#103). go func() { s.ForceReconnectFn() emit(fmt.Sprintf("MQTT [%s] WATCHDOG reconnect attempt issued", s.Tag)) From c9f9f693e5168f36730b25f90e450b532f527375 Mon Sep 17 00:00:00 2001 From: dborup Date: Sun, 4 Oct 2026 10:43:36 +0200 Subject: [PATCH 3/4] test(ingestor): keep emitters running across stop in the #103 race test The emitters now start before stop and keep emitting until after it has returned, so -race sees unsynchronised reads both sides of the stopped write. A mutant that reads the flag without the mutex is caught in 18 of 20 runs (was 2 of 10); master still fails 20 of 20. Relates to #103 Co-Authored-By: Claude Opus 5.5 --- cmd/ingestor/mqtt_watchdog_102_103_test.go | 27 ++++++++++++++++------ 1 file changed, 20 insertions(+), 7 deletions(-) diff --git a/cmd/ingestor/mqtt_watchdog_102_103_test.go b/cmd/ingestor/mqtt_watchdog_102_103_test.go index 313174a0e..98397041f 100644 --- a/cmd/ingestor/mqtt_watchdog_102_103_test.go +++ b/cmd/ingestor/mqtt_watchdog_102_103_test.go @@ -132,13 +132,16 @@ func TestNewAsyncEmit_EmitAfterStopDoesNotPanic_103(t *testing.T) { } } -// #103: emits racing stop neither panic nor race (run with -race). +// #103: emits racing stop neither panic nor race (run with -race). Every +// emitter is already emitting when stop runs and keeps going until after it +// has returned. func TestNewAsyncEmit_ConcurrentEmitDuringStop_103(t *testing.T) { emit, stop := newAsyncEmit(func(...any) {}) - start := make(chan struct{}) + stopped := make(chan struct{}) panics := make(chan any, 8) - var wg sync.WaitGroup + var running, wg sync.WaitGroup for g := 0; g < 8; g++ { + running.Add(1) wg.Add(1) go func() { defer wg.Done() @@ -147,14 +150,24 @@ func TestNewAsyncEmit_ConcurrentEmitDuringStop_103(t *testing.T) { panics <- r } }() - <-start - for i := 0; i < 2000; i++ { - emit("line", i) + emit("first line") + running.Done() + for { + select { + case <-stopped: + for i := 0; i < 100; i++ { + emit("after stop", i) + } + return + default: + emit("line") + } } }() } - close(start) + running.Wait() stop() + close(stopped) wg.Wait() close(panics) for r := range panics { From 29c2d4d9ffbda9b9cf662098053f673b5b719d40 Mon Sep 17 00:00:00 2001 From: dborup Date: Sun, 4 Oct 2026 12:17:24 +0200 Subject: [PATCH 4/4] reviewfix(ingestor): keep paho's error in the retry-pending line; pin emit-after-stop (#102, #103) Review feedback on #210: - F1: IsConnected() in status disconnecting reflects paho's sticky willReconnect flag, so it is a strong hint that a retry follows, not a guarantee. The info line now keeps paho's error and says only what paho reports: "Connect() returned ; paho reports a retry pending (IsConnected=true), not starting a new attempt". The comments on buildForceReconnectFn and connectRetryInProgress say so. A new real-paho test covers Disconnect() while paho is reconnecting: disconnecting, IsConnected()=true, no retry follows, and the line keeps the error. - F2: a new test pins that emit after stop is not counted in watchdogLogDropCount. - F3: the concurrent emit/stop test runs 200 rounds, so a stop that closes the queue before it marks itself stopped is caught reliably. Relates to #102, #103 Co-Authored-By: Claude Opus 5.5 --- cmd/ingestor/main.go | 29 ++++-- cmd/ingestor/mqtt_watchdog_102_103_test.go | 106 ++++++++++++++++++--- 2 files changed, 112 insertions(+), 23 deletions(-) diff --git a/cmd/ingestor/main.go b/cmd/ingestor/main.go index c6fbc3ad5..51511235b 100644 --- a/cmd/ingestor/main.go +++ b/cmd/ingestor/main.go @@ -573,13 +573,17 @@ func buildMQTTOpts(source MQTTSource) *mqtt.ClientOptions { // What Connect() then does depends on paho's status (paho.mqtt.golang // v1.5.0, client.go Connect and status.go Connecting): // - reconnecting (AutoReconnect loop): a no-op that returns a success token; -// - connecting (the initial ConnectRetry loop) or disconnecting with a -// reconnect pending: an error token, errStatusMustBeDisconnected, while -// paho's own loop keeps retrying. Logged as info, not as a failure (#102); +// - connecting (the initial ConnectRetry loop) or disconnecting: an error +// token, errStatusMustBeDisconnected; // - disconnected (paho gave up): starts a fresh attempt. // -// The same status error with no retry pending (disconnecting after a -// Disconnect) is a genuine failure and is logged as one. +// For that status error, IsConnected() tells whether paho reports a retry +// pending. In connecting it does (ConnectRetry keeps retrying), so the line +// is info, not a failure (#102). In disconnecting it reflects paho's +// willReconnect flag, which is sticky after any earlier auto-reconnect and is +// left set by a Disconnect(), so there it is a strong hint, not a guarantee. +// The info line therefore keeps paho's error and claims no more than paho +// reports. With IsConnected() false the error is logged as a failure. func buildForceReconnectFn(client mqtt.Client, tag string) func() { return func() { if client.IsConnectionOpen() { @@ -594,7 +598,7 @@ func buildForceReconnectFn(client mqtt.Client, tag string) func() { switch { case err == nil: case connectRetryInProgress(client, err): - log.Printf("MQTT [%s] WATCHDOG force-reconnect: retry already in progress in paho, no new attempt needed", tag) + log.Printf("MQTT [%s] WATCHDOG force-reconnect: Connect() returned %v; paho reports a retry pending (IsConnected=true), not starting a new attempt", tag, err) default: log.Printf("MQTT [%s] WATCHDOG force-reconnect Connect() failed: %v", tag, err) } @@ -608,10 +612,15 @@ func buildForceReconnectFn(client mqtt.Client, tag string) func() { // fails if an upgrade changes it. const pahoErrStatusMustBeDisconnected = "status can only transition to connecting from disconnected" -// connectRetryInProgress reports whether a Connect() error only means paho is -// already retrying (#102). IsConnected() is true in status connecting only -// with ConnectRetry, and in disconnecting only when a reconnect will follow; -// buildMQTTOpts sets ConnectRetry and AutoReconnect. +// connectRetryInProgress reports whether a Connect() error is paho's status +// error while paho reports a retry pending (#102). buildMQTTOpts sets +// ConnectRetry and AutoReconnect, so IsConnected() is true in status +// connecting (the initial retry loop keeps going) and in disconnecting when +// paho's willReconnect flag is set. That flag is sticky: an earlier +// auto-reconnect sets it, a successful reconnect never clears it, and a +// Disconnect() from connected or reconnecting leaves it alone. In +// disconnecting a true result is therefore a strong hint that a retry +// follows, not a guarantee. func connectRetryInProgress(client mqtt.Client, err error) bool { return err.Error() == pahoErrStatusMustBeDisconnected && client.IsConnected() } diff --git a/cmd/ingestor/mqtt_watchdog_102_103_test.go b/cmd/ingestor/mqtt_watchdog_102_103_test.go index 98397041f..1b994af26 100644 --- a/cmd/ingestor/mqtt_watchdog_102_103_test.go +++ b/cmd/ingestor/mqtt_watchdog_102_103_test.go @@ -20,7 +20,10 @@ import ( const ( connectFailedLine102 = "WATCHDOG force-reconnect Connect() failed" - retryInProgress102 = "retry already in progress" + retryPending102 = "paho reports a retry pending" + // retryPendingErr102 is the start of the retry-pending line: the line keeps + // paho's error, because IsConnected() is a hint and not a guarantee. + retryPendingErr102 = "WATCHDOG force-reconnect: Connect() returned " + pahoErrStatusMustBeDisconnected + "; " + retryPending102 ) // newNeverConnectedTestClient builds the client as main does, against a @@ -44,8 +47,8 @@ func newNeverConnectedTestClient(t *testing.T, b *forceReconnectTestBroker, tag // #102: five triggers during the initial ConnectRetry loop (the #29 review // measured 5/5). Connect() returns paho's errStatusMustBeDisconnected each // time; paho keeps retrying, so this is not a failure. The client must not -// log "Connect() failed", must log that a retry is in progress, and paho's -// loop must keep going and connect once the broker is up. +// log "Connect() failed", must log paho's error with the retry it reports +// pending, and paho's loop must keep going and connect once the broker is up. func TestForceReconnect_RealPaho_InitialRetryLoopIsNotAConnectFailure_102(t *testing.T) { b := newForceReconnectTestBroker(t) logs := captureLog118(t) @@ -62,8 +65,8 @@ func TestForceReconnect_RealPaho_InitialRetryLoopIsNotAConnectFailure_102(t *tes if n := strings.Count(logs.String(), connectFailedLine102); n != 0 { t.Errorf("%d %q lines while paho was in its initial retry loop:\n%s", n, connectFailedLine102, logs) } - if n := strings.Count(logs.String(), retryInProgress102); n != 5 { - t.Errorf("%d %q lines, want 5:\n%s", n, retryInProgress102, logs) + if n := strings.Count(logs.String(), retryPendingErr102); n != 5 { + t.Errorf("%d %q lines, want 5:\n%s", n, retryPendingErr102, logs) } n, _ := b.counts() @@ -93,8 +96,61 @@ func TestForceReconnect_RealPaho_ConnectErrorWhileDisconnectingIsLogged_102(t *t if n := strings.Count(logs.String(), connectFailedLine102); n != 1 { t.Errorf("%d %q lines, want 1:\n%s", n, connectFailedLine102, logs) } - if n := strings.Count(logs.String(), retryInProgress102); n != 0 { - t.Errorf("a Connect() error with no retry pending was logged as %q:\n%s", retryInProgress102, logs) + if n := strings.Count(logs.String(), retryPending102); n != 0 { + t.Errorf("a Connect() error with no retry pending was logged as %q:\n%s", retryPending102, logs) + } +} + +// #102 review F1: IsConnected() in status disconnecting reflects paho's +// willReconnect, which a Disconnect() leaves set. Here paho is reconnecting +// after a connection loss (willReconnect=true) when Disconnect() runs; status +// is then disconnecting, IsConnected() stays true, and Connect() returns +// errStatusMustBeDisconnected, yet no retry follows: paho ends disconnected. +// So the retry-pending line must keep paho's error and must not claim more. +func TestForceReconnect_RealPaho_DisconnectWhileReconnectingKeepsConnectError_102(t *testing.T) { + b := newForceReconnectTestBroker(t) + logs := captureLog118(t) + // A 3s cap keeps paho's sleep between attempts (1s, then 2s) far above + // the 5ms polling below, so the disconnecting window cannot be missed. + opts := buildMQTTOpts(MQTTSource{Broker: b.url(), Name: "force-reconnect-sticky"}). + SetConnectTimeout(500 * time.Millisecond). + SetWriteTimeout(500 * time.Millisecond). + SetMaxReconnectInterval(3 * time.Second) + client := mqtt.NewClient(opts) + client.Connect() + t.Cleanup(func() { client.Disconnect(0) }) + if !pollUntil(forceReconnectTestDeadline, client.IsConnectionOpen) { + t.Fatal("test setup: client never connected to the test broker") + } + + // Connection lost: paho's AutoReconnect loop starts (willReconnect=true) + // and sleeps after its first failed attempt. + n, _ := b.counts() + b.goDown() + waitForConnects(t, b, n, "test setup: paho's reconnect loop") + + // Disconnect() moves reconnecting to disconnecting and waits for the + // loop's sleep to end; willReconnect stays true. Connect() is a no-op + // probe here: a success token while reconnecting, the status error once + // disconnecting. + client.Disconnect(0) + if !pollUntil(time.Second, func() bool { return client.Connect().Error() != nil }) { + t.Fatal("test setup: paho never reached status disconnecting") + } + if !client.IsConnected() || client.IsConnectionOpen() { + t.Fatalf("test setup: IsConnected=%v IsConnectionOpen=%v, want true/false (disconnecting, willReconnect set)", client.IsConnected(), client.IsConnectionOpen()) + } + + buildForceReconnectFn(client, "force-reconnect-sticky")() + if n := strings.Count(logs.String(), retryPendingErr102); n != 1 { + t.Errorf("%d %q lines, want 1 (the line must keep paho's error):\n%s", n, retryPendingErr102, logs) + } + + // The reported retry never comes: with the broker back up, paho still + // ends disconnected rather than connected. + b.goUp() + if !pollUntil(forceReconnectTestDeadline, func() bool { return !client.IsConnected() }) { + t.Fatal("paho did not settle to disconnected after Disconnect(); the scenario no longer shows that IsConnected() is only a hint") } } @@ -107,8 +163,8 @@ func TestBuildForceReconnectFn_OtherConnectErrorIsLogged_102(t *testing.T) { if n := strings.Count(logs.String(), connectFailedLine102+": network unreachable"); n != 1 { t.Errorf("Connect() error not logged as a failure:\n%s", logs) } - if strings.Count(logs.String(), retryInProgress102) != 0 { - t.Errorf("an unrelated Connect() error was logged as %q:\n%s", retryInProgress102, logs) + if strings.Count(logs.String(), retryPending102) != 0 { + t.Errorf("an unrelated Connect() error was logged as %q:\n%s", retryPending102, logs) } } @@ -132,10 +188,36 @@ func TestNewAsyncEmit_EmitAfterStopDoesNotPanic_103(t *testing.T) { } } +// #103 review F2: a line dropped after stop is not a "queue full" drop, so it +// must not count in watchdogLogDropCount. +func TestNewAsyncEmit_EmitAfterStopIsNotCountedAsDrop_103(t *testing.T) { + emit, stop := newAsyncEmit(func(...any) {}) + stop() + before := WatchdogLogDropCount() + for i := 0; i < asyncEmitQueueSize+10; i++ { + emit("after stop", i) + } + if got := WatchdogLogDropCount() - before; got != 0 { + t.Fatalf("emit after stop counted %d drops in watchdogLogDropCount, want 0", got) + } +} + // #103: emits racing stop neither panic nor race (run with -race). Every // emitter is already emitting when stop runs and keeps going until after it -// has returned. +// has returned. A single round can miss a narrow window in stop (review F3: +// a stop that closes the queue before it marks itself stopped panics in +// about half the rounds without -race and fewer with it), so it runs many. func TestNewAsyncEmit_ConcurrentEmitDuringStop_103(t *testing.T) { + for round := 0; round < 200; round++ { + if r := concurrentEmitDuringStopRound(); r != nil { + t.Fatalf("round %d: emit racing stop panicked: %v", round, r) + } + } +} + +// concurrentEmitDuringStopRound runs one round of +// TestNewAsyncEmit_ConcurrentEmitDuringStop_103 and returns the first panic. +func concurrentEmitDuringStopRound() any { emit, stop := newAsyncEmit(func(...any) {}) stopped := make(chan struct{}) panics := make(chan any, 8) @@ -170,9 +252,7 @@ func TestNewAsyncEmit_ConcurrentEmitDuringStop_103(t *testing.T) { close(stopped) wg.Wait() close(panics) - for r := range panics { - t.Fatalf("emit racing stop panicked: %v", r) - } + return <-panics } // #103, the scenario from the issue: the production watchdog fires a