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
48 changes: 43 additions & 5 deletions cmd/ingestor/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -568,9 +568,22 @@ 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: an error
// token, errStatusMustBeDisconnected;
// - disconnected (paho gave up): starts a fresh attempt.
//
// 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() {
Expand All @@ -581,12 +594,37 @@ 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: 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)
}
}
}

// 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 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()
}

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.
Expand Down
11 changes: 7 additions & 4 deletions cmd/ingestor/mqtt_force_reconnect_race_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -25,21 +25,23 @@ 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
connectCalled bool
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
Expand Down Expand Up @@ -126,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")
Expand Down
21 changes: 17 additions & 4 deletions cmd/ingestor/mqtt_watchdog.go
Original file line number Diff line number Diff line change
Expand Up @@ -111,27 +111,39 @@ 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 {
realEmit(args...)
}
}()
emit = func(args ...any) {
mu.Lock()
defer mu.Unlock()
if stopped {
return
}
select {
case queue <- args:
default:
watchdogLogDropCount.Add(1)
}
}
stop = func() {
mu.Lock()
stopped = true
close(queue)
mu.Unlock()
<-drained
}
return emit, stop
Expand Down Expand Up @@ -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))
Expand Down
Loading
Loading