fix(ingestor): quiet watchdog retry noise and make force-reconnect shutdown-safe (#102, #103) - #210
Conversation
…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 <noreply@anthropic.com>
…utdown-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 <noreply@anthropic.com>
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 <noreply@anthropic.com>
Rapport — CS-Minimax PR#210 #102+#103 — head c9f9f69Status: ready for review (draft): both issues fixed in Base
Paho behaviour comes from the paho.mqtt.golang v1.5.0 source in the module cache ( Evidence tags: [T] test or CI, [A] analysis, [K] known, not re-run. Acceptance criteria
Tests
MutantsEach mutant was applied to the fix commit, run under
All other tests stay green under each mutant, so each mutant is caught by its own test. [T] PerformanceThe watchdog emit path is not a hot path (about one line per source per 60 s tick). [A] In a scratch benchmark (not committed, n=6), uncontended Rules
CI (run 37189990781, head
|
| Job | Result |
|---|---|
| Go Build & Test | success (ingestor package ok, 118 s with coverage) [T] |
| Playwright E2E Tests | success [T] |
| Build & Publish Docker Image | success [T] |
| Release Artifacts | skipped (PR) [T] |
| Deploy Staging | skipped (PR) [T] |
| Publish Badges & Summary | skipped (PR) [T] |
Remaining
- The error is matched on the text of paho's unexported
errStatusMustBeDisconnected. A paho upgrade that changes the text fails…InitialRetryLoopIsNotAConnectFailure_102; it is not silently missed. [A] - The "reconnect attempt issued" line of a force-reconnect still in flight at shutdown is now dropped, not logged. It is not counted in
WatchdogLogDropCount, which keeps meaning "queue full". [A] origin/masterhas moved to8bafcdf2(fix(ui): 48px navigation and channel buttons (port of upstream 2078) #201, UI only) since the base. GitHub reports the PR as mergeable. [T]- Not run: browser validation (backend-only change) and staging. [A]
Generated by Claude Code
Review — CS-Macmini PR#210 mqtt-watchdog — head c9f9f69Dom: APPROVE med nits This is an independent, read-only review of head Evidence tags: [T] run here, [A] analysis of the source, [K] taken from the author's report or CI, not re-run. Findings
F1 in detailThese are review-only probe tests, not committed. They use a real paho client against the in-test broker and read paho's internal status through reflection.
How far this reaches [A]:
Suggested fix, small and in this PR or a follow-up:
1. #102: paho behaviour
2. #103: shutdown
3. Normal operationBefore 4. Tests and mutants
My own mutants were each applied to a copy of head and run with
Separately, the author reports mutants M1–M5 as killed [K]. 5. PerformanceNot a hot path. 6. Rules
Not verified
The head was Generated by Claude Code |
… 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 <err>; 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 <noreply@anthropic.com>
Rapport — CS-Minimax PR#210 runde 2 — head 29c2d4dStatus: review nits F1–F3 addressed in one commit in Review feedback addressed (commit Evidence tags: [T] test or run here, [A] analysis of the source (paho.mqtt.golang v1.5.0
Tests
MutantsEach mutant was applied to a copy of
M1–M5 from round 1 and M6/M8 from the review were not re-run. [K] The code they target is unchanged except for the log text. Scope
CI and remaining
Generated by Claude Code |
main() is not unit-testable, so a source-text guard (as for the route mask backfill) pins its share of the wiring: prepareMQTTSource, then setup.attachClient(client) before the first Connect(), and the initial connect error logged with setup.secrets. Red at fe40aae. Kills the review's mutant M5 and its variants, and guards the coming merge with #210 in the same loop. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Conflict in cmd/ingestor/main.go buildForceReconnectFn, resolved by keeping both sides: - from #210: the doc comment on paho's Connect() statuses, the switch on err with connectRetryInProgress, and the retry-pending info line; - from #213: the secrets ...string parameter, and errForLog(err, secrets...) in both the retry-pending and the Connect() failed line. connectRetryInProgress still reads the raw error, whose text it compares with paho's. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
#159) #210 added a retry-pending info line next to the Connect() failed line. The test drives it with paho's status error and IsConnected()=true, and a configured password that occurs in that text, and asserts the line is masked and still classified as retry-pending (classification reads the raw error). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Relates to #102, #103
Two fixes to the MQTT watchdog's force-reconnect in
cmd/ingestor.Plan
fc0c068f, red on master), incmd/ingestor/mqtt_watchdog_102_103_test.go:mqtt_force_reconnect_paho_test.go:ConnectRetryloop (statusconnecting): no "Connect() failed" line; an info line with paho's error and "paho reports a retry pending" is logged instead;emitafterstop, andemitracingstop, must not panic;ForceReconnectFn, then shutdown, then unblock.4db36ab8) inmain.go(buildForceReconnectFn) andmqtt_watchdog.go(newAsyncEmit).c9f9f693): the race test keeps the emitters running acrossstop.29c2d4d9): the retry-pending line keeps paho's error and claims less (IsConnected()indisconnectingis a hint, not a guarantee); emit afterstopis pinned as not counted inwatchdogLogDropCount; the emit/stop race test runs 200 rounds.#102: retry noise during the initial ConnectRetry loop
What paho v1.5.0 does (
client.goConnect,status.goConnecting):Connect()returnsIsConnected()reconnectingconnecting(initial ConnectRetry loop)errStatusMustBeDisconnecteddisconnecting,willReconnectseterrStatusMustBeDisconnectedwillReconnect)disconnecting,willReconnectclearerrStatusMustBeDisconnecteddisconnectedbuildForceReconnectFnstill callsConnect()exactly as before. It now classifies the result:errStatusMustBeDisconnectedtogether withIsConnected() == trueis logged as info:WATCHDOG force-reconnect: Connect() returned <err>; paho reports a retry pending (IsConnected=true), not starting a new attempt. paho's errors are unexported, so the error is matched on its text, kept in one constant. The real-paho test fails if a paho upgrade changes that text.connectingthis is reliable: paho'sConnectRetryloop keeps going. Indisconnecting,IsConnected()reflects paho'swillReconnectflag, which is sticky: an earlier auto-reconnect sets it, a successful reconnect never clears it, and aDisconnect()fromconnectedorreconnectingleaves it alone. So after aDisconnect()paho can report a retry pending that never comes. That makesIsConnected()a strong hint, not a guarantee, so the line keeps paho's error and claims no more than paho reports.IsConnected() == false(nothing will retry), and any other error, are still logged asWATCHDOG force-reconnect Connect() failed: ….The doc comment no longer says that
Connect()is a no-op in every retrying state. It now lists the cases from the table.#103:
send on closed channelat shutdownmaybeForceReconnectrunsForceReconnectFnin a goroutine that is not joined, and then emits "reconnect attempt issued". On SIGTERM,runLivenessWatchdog's stop waits for the loop goroutine and then closes the emit queue. A force-reconnect still blocked inConnect()/Disconnect()then emits on the closed channel and panics.Fix:
newAsyncEmit'semitchecks astoppedflag under a mutex, andstopsets the flag and closes the queue under the same mutex. A line emitted afterstopis dropped.stopstill does not wait for a blockedForceReconnectFn, so shutdown is not delayed by a hanging connect. I chose this over joining the goroutines for that reason, and because any late emitter is now safe, not just this one.No change in normal operation: before
stop,emitqueues the line, or counts a drop when the queue is full, exactly as before.Tests
TestForceReconnect_RealPaho_InitialRetryLoopIsNotAConnectFailure_102TestForceReconnect_RealPaho_ConnectErrorWhileDisconnectingIsLogged_102TestForceReconnect_RealPaho_DisconnectWhileReconnectingKeepsConnectError_102disconnectingwithwillReconnectset afterDisconnect(); no retry follows)TestBuildForceReconnectFn_OtherConnectErrorIsLogged_102TestNewAsyncEmit_EmitAfterStopDoesNotPanic_103send on closed channelTestNewAsyncEmit_EmitAfterStopIsNotCountedAsDrop_103stopis not a "queue full" drop)TestNewAsyncEmit_ConcurrentEmitDuringStop_103(200 rounds)TestRunLivenessWatchdog_ForceReconnectInFlightAtShutdown_103panic: send on closed channel(test binary crashes)fakeClient.IsConnected(mqtt_force_reconnect_race_test.go) now returns a field instead of panicking, because the fix calls it after aConnect()error.Performance
The watchdog emit path is not a hot path: it runs about once per source per 60 s tick, plus one line per force-reconnect. The uncontended mutex adds about 15 ns per
emit. In a scratch benchmark (not committed, n=6, darwin/arm64), the oldemittook a median of 46 ns/op and the new one 62 ns/op, with the same 16 B and 1 alloc. The packet ingest path is untouched.Rules
cmd/ingestor;cmd/server,internal/and.githubare untouched (fork guards indeploy.yml: 9).map[string]interface{}.🤖 Generated with Claude Code