fix(mqtt): escalate persistent paho disconnect + recover from emit panic + expose watchdog tick (#1749) (#1810)

# Partial fix for #1749 — MQTT watchdog escalation + panic recovery +
tick exposure

Red commit: 9912bbb3e9 (CI run:
https://github.com/Kpa-clawbot/CoreScope/actions/runs?branch=fix%2Fissue-1749)

## Problem
Production CoreScope v3.9.1 (and a more recent prod recurrence on
2026-06-30 with the wcmesh source on `ssl://mqtt2.wcmesh.com:8883`)
showed two distinct watchdog failure modes:

1. **Per-source paho machinery dies silently.** `IsConnectedFn` returns
false; paho's `SetAutoReconnect(true)` never retries. The watchdog's
`processLivenessTransition` deliberately stays silent on
`LivenessDisconnected`, trusting paho to recover — so there is no
escalation path when that trust is misplaced.
2. **Watchdog goroutine death.** Three sources went silent within ~60s
of each other; no `WATCHDOG` log lines for 75 min. The most plausible
single point of failure is a panic inside `emit` (e.g. a blocked log
pipe) killing the loop with no defer/recover.

## Changes

**`cmd/ingestor/mqtt_watchdog.go`**
- Added `disconnectedReconnectMultiplier = 5` constant.
- Added `DisconnectedSinceUnix int64` (atomic) on `SourceLivenessState`.
Stamped on the first tick the source is observed disconnected; cleared
on any non-disconnected tick.
- `processLivenessTransition` now escalates: when `(now -
DisconnectedSinceUnix) > multiplier × threshold`, emits a `WATCHDOG
ESCALATION` WARN and calls `maybeForceReconnect` (subject to existing
`forceReconnectThrottle`). Distinct from the existing `LivenessStalled`
path so operators can grep escalation events independently.
- Added package-level `watchdogLastTickUnix atomic.Int64` +
`WatchdogLastTickUnix()` getter. The loop stamps it BEFORE per-source
processing — a wedged source-handler does not freeze the clock for an
external observer.
- `runLivenessWatchdogLoop` wraps each per-source
`processLivenessTransition` call in `func() { defer recover; ... }()` so
a panic in `emit` (or in any per-source code path) is logged and
skipped, not fatal. The loop continues to the next source and the next
tick.

**`cmd/ingestor/stats_file.go`**
- Added `WatchdogLastTickUnix int64` field on `IngestorStatsSnapshot`
(additive, `omitempty`); populated from `WatchdogLastTickUnix()` each
stats tick.

**`cmd/server/mqtt_status.go`**
- `MqttStatusResponse` gains `WatchdogLastTickUnix int64` (additive,
`omitempty`) sourced from the ingestor stats file; surfaced via `GET
/api/mqtt/status`.

**`config.example.json`**
- No new config field added — the multiplier is a code constant (5×) per
the issue's "N×threshold (e.g. 5×)" recommendation. The active
per-source `threshold` is the 5-minute scan threshold hard-coded at the
sole `runLivenessWatchdog` callsite (`cmd/ingestor/main.go:460`:
`runLivenessWatchdog(60*time.Second, 5*time.Minute)`), so escalation
fires at ~25 minutes (5 × 5min) of continuous disconnect — plus a
deterministic per-source jitter of 0..30s (#1810 round-1, see Taleb #4)
to avoid synchronized escalation across N sources sharing an upstream
broker outage.

## Acceptance (#1749)
- [x] Persistent `LivenessDisconnected` > N×threshold → force-reconnect
+ WARN
- [x] Watchdog goroutine liveness clock exposed (`WatchdogLastTickUnix`
in `/api/mqtt/status`)
- [x] Test: `IsConnectedFn` false for >5×threshold → assert
`ForceReconnectFn` invoked at least once
- [x] Test: panic in `emit` → assert loop recovers and continues ticking

## Test plan
- 4 new tests in `cmd/ingestor/mqtt_watchdog_1749_test.go` (all RED on
master, GREEN on this PR).
- Existing watchdog tests (`mqtt_watchdog_force_reconnect_test.go`,
`mqtt_reconnect_test.go`, r1/r2/m1 suites) continue to pass — the
escalation path is additive.

## Preflight
- TDD: red commit pushed and asserted to fail BEFORE green commit
landed.
- PII grep: clean on diff and PR body.
- Worktree: `_wt-fix-1749` on branch `fix/issue-1749`.

---------

Co-authored-by: corescope-bot <bot@corescope.dev>
Co-authored-by: openclaw-bot <bot@openclaw.local>
Co-authored-by: bot <bot@local>
This commit is contained in:
Kpa-clawbot
2026-06-30 15:16:31 -07:00
committed by GitHub
co-authored by corescope-bot openclaw-bot bot
parent b74a64ccfa
commit 242c7c609b
7 changed files with 1043 additions and 21 deletions
+198 -4
View File
@@ -2,7 +2,9 @@ package main
import (
"fmt"
"hash/fnv"
"log"
"os"
"sync"
"sync/atomic"
"time"
@@ -18,6 +20,60 @@ const livenessHeartbeatInterval = time.Hour
// reconnects on the SAME source. See processLivenessTransition.
const forceReconnectThrottle = 60 * time.Second
// disconnectedReconnectMultiplier (#1749) governs how long a source may
// stay in LivenessDisconnected before the watchdog escalates with a
// forced reconnect. paho's SetAutoReconnect(true) normally recovers a
// dropped connection; in production we have observed paho's reconnect
// machinery silently dying for a single source while another source on
// the same binary reconnects fine (prod 2026-06-30: connectCount=1,
// disconnectCount=1, lastError="EOF", zero retries for 18h). When the
// source stays !IsConnected for more than `multiplier × threshold`,
// the watchdog forces a reconnect rather than trusting paho to recover.
const disconnectedReconnectMultiplier = 5
// watchdogLastTickUnix (#1749) is the wall-clock unix-seconds timestamp
// of the most recent runLivenessWatchdogLoop tick. The watchdog
// goroutine has itself gone silent in production (#1749: 3 sources
// stalled simultaneously for 75 min with zero WATCHDOG log lines),
// suggesting goroutine-level failure. Exposing this clock via
// /api/mqtt/status lets external monitoring assert that the watchdog
// is still ticking — a stale value (e.g. > 2× the scan interval)
// indicates the watchdog itself is dead.
//
// Style note (#1810 round-1, adv #5): new package-level counters in
// this file use atomic.Int64 (typed, method-based) while per-source
// state on SourceLivenessState uses plain int64 + atomic.StoreInt64.
// The struct fields stay int64 because they are accessed through
// pointer receivers all over the codebase and atomic.Int64 inside a
// struct breaks the "noCopy" semantics expected of value receivers
// in a few callsites; package-level vars have no such constraint and
// the typed form catches misuse at compile time. Kept-both is
// intentional, not drift.
var watchdogLastTickUnix atomic.Int64
// watchdogPanicCount (#1810 round-1, Taleb #3) counts recovered panics
// inside the per-source watchdog work IIFE. A loop that panic-loops
// every tick currently looks healthy by WatchdogLastTickUnix alone —
// the tick stamp lands BEFORE the per-source work, so a panic on every
// source still advances the clock. This counter, surfaced via
// /api/mqtt/status and the stats snapshot, lets external monitoring
// alarm on a rapidly-growing value (= the loop is alive but the work
// is broken).
var watchdogPanicCount atomic.Int64
// WatchdogLastTickUnix returns the unix-seconds timestamp of the most
// recent watchdog tick. Returns 0 if the watchdog has never ticked.
func WatchdogLastTickUnix() int64 {
return watchdogLastTickUnix.Load()
}
// WatchdogPanicCount returns the running total of recovered panics
// inside the watchdog per-source work IIFE (#1810). Monotonic across
// the process lifetime; never decreases.
func WatchdogPanicCount() int64 {
return watchdogPanicCount.Load()
}
// LivenessKind enumerates the watchdog verdicts for a source. Edge-triggered
// transitions use this to decide whether to emit (and what severity).
type LivenessKind int
@@ -55,8 +111,8 @@ const (
// window. r1's StartedAt-as-grace-clock conflated transient-stall
// suppression with cold-start grace; r2 separates them.
type SourceLivenessState struct {
Tag string
Broker string
Tag string
Broker string
LastMessageUnix int64 // atomic; unix seconds of last successfully WRITTEN MQTT message (handleMessage post-write)
// LastReceiptUnix (PR #1609 M1) is stamped at MQTT receipt time —
// BEFORE the message is handed to the buffer/writer. STUB: unused
@@ -88,6 +144,15 @@ type SourceLivenessState struct {
// recent forced reconnect for this source; the watchdog reads it
// to enforce forceReconnectThrottle. atomic.
LastForceReconnectUnix int64
// DisconnectedSinceUnix (#1749) is the unix-seconds timestamp of
// the FIRST tick on which the watchdog observed this source in
// LivenessDisconnected (paho reports !IsConnected). Cleared back
// to 0 on any tick where the source is NOT disconnected. When the
// gap (now - DisconnectedSinceUnix) exceeds
// disconnectedReconnectMultiplier × threshold, the watchdog
// escalates with a forced reconnect on the assumption that paho's
// own auto-reconnect machinery has silently died. atomic.
DisconnectedSinceUnix int64
// AttemptCount is incremented on every TCP/TLS connection attempt. Used
// by ConnectionAttemptHandler to log attempt # independent of paho's
// internal reconnect-loop state. atomic.
@@ -127,6 +192,14 @@ func (s *SourceLivenessState) MarkReconnected(now time.Time) {
atomic.StoreInt64(&s.LastMessageUnix, 0)
atomic.StoreInt64(&s.StartedAt, now.Unix())
atomic.StoreInt64(&s.LastAlertUnix, 0)
// #1810 round-1 (Taleb #5): clear DisconnectedSinceUnix so the
// next LivenessDisconnected observation is treated as a NEW outage
// and starts its escalation timer from scratch. Without this, a
// post-recovery disconnect immediately satisfies (now -
// DisconnectedSinceUnix > multiplier × threshold) and force-
// reconnects on the first tick — making the escalation a
// trigger-on-flap instead of trigger-on-persistent-failure.
atomic.StoreInt64(&s.DisconnectedSinceUnix, 0)
}
// checkSourceLiveness returns (message, kind) describing the source's
@@ -315,6 +388,13 @@ func runLivenessWatchdogLoop(tick <-chan time.Time, done <-chan struct{}, thresh
if !ok {
return
}
// #1749: stamp the watchdog clock BEFORE per-source work so
// a panic or hang inside processLivenessTransition does
// not freeze the heartbeat that /api/mqtt/status exposes.
// External monitoring on WatchdogLastTickUnix detects a
// wedged loop only if this clock is fresh-while-ticking
// and stale-when-dead.
watchdogLastTickUnix.Store(now.Unix())
livenessRegistryMu.RLock()
states := make([]*SourceLivenessState, 0, len(livenessRegistry))
for _, s := range livenessRegistry {
@@ -322,13 +402,128 @@ func runLivenessWatchdogLoop(tick <-chan time.Time, done <-chan struct{}, thresh
}
livenessRegistryMu.RUnlock()
for _, s := range states {
// #1749: handle disconnect-escalation bookkeeping
// BEFORE the transition dispatch. checkSourceLiveness
// returns LivenessDisconnected when paho reports
// !IsConnected; we track how long that has persisted
// and escalate with a forced reconnect when paho's
// own auto-reconnect machinery has clearly failed
// (multiplier × threshold without recovery).
msg, kind := checkSourceLiveness(s, threshold, now)
processLivenessTransition(s, kind, msg, now, emit)
// #1749: a panic in emit (blocked log pipe, full
// Docker JSON-file driver, etc.) MUST NOT kill the
// watchdog goroutine. Recover per-source so one bad
// source — or one bad log call — does not silence
// all monitoring across all sources. Both
// maybeEscalateDisconnected and processLivenessTransition
// call emit, so both must be inside the recover scope.
func(state *SourceLivenessState, k LivenessKind, m string) {
defer func() {
if r := recover(); r != nil {
// #1810 round-1 (Taleb #2 + #3):
// (a) Write to os.Stderr DIRECTLY rather than via
// log.Printf. The original log sink is the prime
// suspect for a panic in emit (blocked stderr pipe,
// full Docker JSON-file driver) — using the same
// sink to report the recovery risks a second
// panic-on-recover that kills the goroutine the
// recover was meant to save. os.Stderr.Write is a
// raw syscall and bypasses log's own mutex.
// (b) Increment watchdogPanicCount so a panic-per-
// tick loop is visible to external monitoring even
// though WatchdogLastTickUnix continues to advance.
watchdogPanicCount.Add(1)
fmt.Fprintf(os.Stderr, "[ingestor] WATCHDOG RECOVERED panic processing source %q: %v\n", state.Tag, r)
}
}()
maybeEscalateDisconnected(state, k, threshold, now, emit)
processLivenessTransition(state, k, m, now, emit)
}(s, kind, msg)
}
}
}
}
// maybeEscalateDisconnected (#1749) tracks how long a source has
// continuously been observed in LivenessDisconnected and triggers a
// forced reconnect (subject to forceReconnectThrottle) once the gap
// exceeds disconnectedReconnectMultiplier × threshold.
//
// Background: paho's SetAutoReconnect(true) is supposed to recover
// dropped connections on its own. In production we have observed paho
// silently giving up on one source — connectCount=1, disconnectCount=1,
// lastError="EOF", zero retries for 18h — while another source on the
// same binary reconnects fine. The existing watchdog's
// processLivenessTransition stays silent on LivenessDisconnected to
// avoid double-logging paho's own disconnect line; this escalation
// path is the recovery hook for the case where paho never tries again.
//
// Behavior:
// - kind != LivenessDisconnected: clear DisconnectedSinceUnix (reset
// the timer; recovery or other state).
// - kind == LivenessDisconnected, first observation: stamp
// DisconnectedSinceUnix = now.
// - kind == LivenessDisconnected, gap >= multiplier × threshold:
// emit a WARN and call maybeForceReconnect (throttled). The
// timestamp is NOT advanced beyond the original observation so the
// emit fires on every tick past the boundary; the throttle inside
// maybeForceReconnect handles the broker-hammering concern.
//
// emit is wrapped in defer/recover by the caller — a panic here does
// not kill the loop.
// disconnectedEscalationGap returns the gap (now - DisconnectedSinceUnix)
// at which the watchdog escalates a persistent LivenessDisconnected
// observation into a forced reconnect. Hoisted out of
// maybeEscalateDisconnected (#1810 round-1, adv #6) so it is computed
// once per call site rather than per-source-per-tick. The result
// includes a per-source jitter offset (#1810 round-1, Taleb #4) so
// that a shared-broker outage across N sources does NOT cause a
// synchronized thundering-herd reconnect at the multiplier × threshold
// boundary. Jitter range: 0..forceReconnectThrottle, deterministic per
// tag (hash-based, no RNG state) so retries do not phase-walk.
func disconnectedEscalationGap(tag string, threshold time.Duration) time.Duration {
base := disconnectedReconnectMultiplier * threshold
if forceReconnectThrottle <= 0 {
return base
}
h := fnv.New32a()
_, _ = h.Write([]byte(tag))
jitter := time.Duration(h.Sum32()%uint32(forceReconnectThrottle/time.Millisecond)) * time.Millisecond
return base + jitter
}
func maybeEscalateDisconnected(s *SourceLivenessState, kind LivenessKind, threshold time.Duration, now time.Time, emit func(...any)) {
if kind != LivenessDisconnected {
atomic.StoreInt64(&s.DisconnectedSinceUnix, 0)
return
}
disconnectedSince := atomic.LoadInt64(&s.DisconnectedSinceUnix)
if disconnectedSince == 0 {
atomic.StoreInt64(&s.DisconnectedSinceUnix, now.Unix())
return
}
gap := now.Sub(time.Unix(disconnectedSince, 0))
escalateAt := disconnectedEscalationGap(s.Tag, threshold)
if gap < escalateAt {
return
}
// #1810 round-1 (adv #2 + Taleb #1): only emit the ESCALATION WARN
// when we are actually going to issue a force-reconnect. Without
// this throttle, every tick past the boundary re-emits the WARN
// — for a 1m scan interval and a 1h outage, that's 55+ duplicate
// alert lines. maybeForceReconnect already enforces
// forceReconnectThrottle and writes its own "forcing reconnect"
// telemetry, so the operator-visible log surface is still
// complete; we just no longer drown them in pre-amble.
lastForce := atomic.LoadInt64(&s.LastForceReconnectUnix)
if lastForce != 0 && now.Sub(time.Unix(lastForce, 0)) < forceReconnectThrottle {
return
}
emit(fmt.Sprintf("MQTT [%s] WATCHDOG ESCALATION: paho disconnected for %s (>%d×threshold=%s) with no auto-reconnect — forcing reconnect (#1749)",
s.Tag, gap.Round(time.Second), disconnectedReconnectMultiplier, escalateAt))
maybeForceReconnect(s, now, emit)
}
// processLivenessTransition applies the edge-trigger rules and updates
// LastAlertUnix accordingly. Separated for testability and to keep the
// loop body small.
@@ -407,4 +602,3 @@ func maybeForceReconnect(s *SourceLivenessState, now time.Time, emit func(...any
emit(fmt.Sprintf("MQTT [%s] WATCHDOG reconnect attempt issued", s.Tag))
}()
}
+409
View File
@@ -0,0 +1,409 @@
package main
import (
"sync"
"sync/atomic"
"testing"
"time"
)
// Issue #1749 — production CoreScope v3.9.1 experienced a complete MQTT
// ingest stall lasting 75+ minutes during which the watchdog never
// fired: no LivenessStalled, no LivenessNeverReceived, no
// LivenessDisconnected log lines, no force-reconnect attempts. Two
// distinct failure modes are addressed here:
//
// 1. Per-source paho machinery dies silently (recurrence with prod
// wcmesh source 2026-06-30: connectCount=1, disconnectCount=1,
// lastError="EOF", zero retries for ~18h while another source on
// the same binary reconnected fine). The original watchdog returns
// silently on LivenessDisconnected, trusting SetAutoReconnect(true)
// to recover — when that trust is misplaced, there is no escalation
// path.
//
// 2. The watchdog goroutine itself dies (3 sources going silent within
// ~60s of each other strongly suggests a single shared dependency
// failed, not 3 independent paho clients failing simultaneously).
// A panic inside the emit callback (log pipe issues observed in
// prior incidents) would kill the loop without leaving a trace.
//
// Fixes asserted here:
// - Persistent LivenessDisconnected past disconnectedReconnectMultiplier
// × threshold MUST trigger a forced reconnect with WARN telemetry.
// - A panic inside emit MUST be recovered and the loop MUST continue
// ticking.
// - WatchdogLastTickUnix MUST advance with every tick so external
// monitoring can detect a wedged watchdog goroutine.
// TestMQTTStallWatchdog_EscalateOnPersistentDisconnect_1749 (RED on
// master): a source that stays !IsConnected for longer than
// disconnectedReconnectMultiplier × threshold MUST be force-reconnected
// at least once. On master, processLivenessTransition returns silently
// on LivenessDisconnected — no escalation — so ForceReconnectFn is
// never invoked.
func TestMQTTStallWatchdog_EscalateOnPersistentDisconnect_1749(t *testing.T) {
defer snapshotAndResetRegistry(t)()
threshold := 60 * time.Second
scanInterval := 5 * time.Millisecond
var reconnectCount atomic.Int32
s := &SourceLivenessState{
Tag: "silent-paho",
Broker: "ssl://mqtt2.example.com:8883",
IsConnectedFn: func() bool { return false }, // paho stuck disconnected
ForceReconnectFn: func() { reconnectCount.Add(1) },
}
if err := registerLivenessState(s); err != nil {
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)
}()
// Feed ticks spanning > (multiplier × threshold) of wall clock so
// the escalation path fires. We control the `now` parameter by
// sending fabricated timestamps down the tick channel.
base := time.Now()
totalSpan := time.Duration(disconnectedReconnectMultiplier+2) * threshold
for elapsed := time.Duration(0); elapsed <= totalSpan; elapsed += scanInterval * 200 {
select {
case tick <- base.Add(elapsed):
case <-time.After(time.Second):
t.Fatal("watchdog loop did not consume tick within 1s")
}
}
// ForceReconnectFn runs in a goroutine in production; poll for the
// counter to land.
deadline := time.Now().Add(2 * time.Second)
for time.Now().Before(deadline) && reconnectCount.Load() < 1 {
time.Sleep(5 * time.Millisecond)
}
if got := reconnectCount.Load(); got < 1 {
t.Fatalf("persistent LivenessDisconnected past %d×threshold MUST force-reconnect at least once (#1749); got %d invocations",
disconnectedReconnectMultiplier, got)
}
}
// TestMQTTStallWatchdog_DisconnectedEscalationThrottled_1749: once the
// watchdog has escalated, repeat escalations on the same source MUST be
// throttled by forceReconnectThrottle (no broker hammering during a
// prolonged outage).
func TestMQTTStallWatchdog_DisconnectedEscalationThrottled_1749(t *testing.T) {
defer snapshotAndResetRegistry(t)()
threshold := 60 * time.Second
var reconnectCount atomic.Int32
s := &SourceLivenessState{
Tag: "throttle-escalate",
Broker: "ssl://example.com:8883",
IsConnectedFn: func() bool { return false },
ForceReconnectFn: func() { reconnectCount.Add(1) },
}
if err := registerLivenessState(s); err != nil {
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) {})
base := time.Now()
// Pre-stamp DisconnectedSinceUnix so that the first tick is
// already past the multiplier×threshold escalation boundary.
// Without this we'd just observe the FIRST tick stamping the
// timestamp and subsequent ticks would be only seconds past it.
atomic.StoreInt64(&s.DisconnectedSinceUnix, base.Add(-time.Duration(disconnectedReconnectMultiplier+1)*threshold).Unix())
// Cross the escalation boundary multiple times within a single
// throttle window — expect ONE reconnect, not many.
for i := 0; i < 10; i++ {
select {
case tick <- base.Add(time.Duration(i) * time.Second):
case <-time.After(time.Second):
t.Fatal("tick blocked")
}
}
time.Sleep(200 * time.Millisecond)
if got := reconnectCount.Load(); got != 1 {
t.Fatalf("escalation must be throttled within %s; got %d invocations", forceReconnectThrottle, got)
}
}
// TestMQTTStallWatchdog_LoopRecoversFromPanicInEmit_1749 (RED on
// master): a panic inside the emit callback MUST NOT kill the watchdog
// loop. On master there is no defer/recover around the per-source
// processLivenessTransition call, so the first panic kills the
// goroutine and no further ticks are processed.
func TestMQTTStallWatchdog_LoopRecoversFromPanicInEmit_1749(t *testing.T) {
defer snapshotAndResetRegistry(t)()
threshold := 1 * time.Minute
s := &SourceLivenessState{
Tag: "panic-emit",
Broker: "tcp://x:1883",
IsConnectedFn: func() bool { return true },
}
atomic.StoreInt64(&s.LastMessageUnix, time.Now().Add(-10*time.Minute).Unix())
atomic.StoreInt64(&s.StartedAt, time.Now().Add(-20*time.Minute).Unix())
if err := registerLivenessState(s); err != nil {
t.Fatalf("setup: %v", err)
}
var mu sync.Mutex
var calls int
emit := func(args ...any) {
mu.Lock()
calls++
mu.Unlock()
panic("synthetic emit panic — simulates blocked log pipe (#1749 hypothesis 2)")
}
tick := make(chan time.Time)
done := make(chan struct{})
// Wrap the loop spawn with our own recover so that an unrecovered
// panic in the loop (the bug on master) does NOT crash the test
// process. The bug-under-test is whether the LOOP recovers; if it
// does not, the panic propagates up to OUR recover, this goroutine
// exits, `exited` is closed, and the loop is dead — at which point
// the second tick will block and we assert failure with a clear
// message rather than tearing down the whole test binary.
exited := make(chan struct{})
go func() {
defer func() {
_ = recover() // RED-mode safety net; production loop is what we are asserting on
close(exited)
}()
runLivenessWatchdogLoop(tick, done, threshold, emit)
}()
base := time.Now()
// Tick 1: should hit the WARN edge, panic in emit, be recovered.
select {
case tick <- base:
case <-time.After(time.Second):
t.Fatal("first tick blocked")
}
// Give the goroutine a moment to recover & re-loop (or to die from
// an unrecovered panic, which is the bug we are gating on).
time.Sleep(100 * time.Millisecond)
select {
case <-exited:
t.Fatal("watchdog loop died from panic in emit (#1749) — production loop lacks defer/recover")
default:
}
// Reset the stalled state's alert so the second tick triggers
// another emit (heartbeat suppression would otherwise mask the
// second call). Use MarkReconnected then re-arm staleness.
s.MarkReconnected(base.Add(50 * time.Millisecond))
atomic.StoreInt64(&s.LastMessageUnix, base.Add(-10*time.Minute).Unix())
atomic.StoreInt64(&s.StartedAt, base.Add(-20*time.Minute).Unix())
// Tick 2: if the loop survived, we should get a second emit call.
select {
case tick <- base.Add(time.Second):
case <-exited:
t.Fatal("watchdog loop died from panic in emit (#1749); second tick cannot be delivered because the goroutine exited")
case <-time.After(time.Second):
t.Fatal("watchdog loop did not survive panic in emit (#1749); second tick blocked because the goroutine died")
}
time.Sleep(100 * time.Millisecond)
mu.Lock()
got := calls
mu.Unlock()
if got < 2 {
t.Fatalf("watchdog loop must survive a panic in emit and continue ticking (#1749); emit calls=%d (expected ≥2)", got)
}
// Loop must still exit cleanly when signalled.
close(done)
select {
case <-exited:
case <-time.After(time.Second):
t.Fatal("loop did not exit after done signal")
}
}
// TestMQTTStallWatchdog_LastTickUnixExposed_1749 (RED on master):
// WatchdogLastTickUnix MUST advance with each tick so external monitoring
// can detect a wedged watchdog goroutine. On master no such clock is
// exposed.
func TestMQTTStallWatchdog_LastTickUnixExposed_1749(t *testing.T) {
defer snapshotAndResetRegistry(t)()
// Baseline: clock should be 0 (or stale) BEFORE the first tick of
// this test. We can't assert exactly 0 because prior tests in the
// same package may have ticked the loop, so just record the value
// and assert it ADVANCES.
before := WatchdogLastTickUnix()
// #1810 round-1 (adv #7): this test stamps a 48h-future value into
// the package-level watchdogLastTickUnix. Restore the prior value
// so a downstream test that asserts "tick advanced past 'before'"
// is not fooled by our leak.
t.Cleanup(func() {
watchdogLastTickUnix.Store(before)
})
tick := make(chan time.Time)
done := make(chan struct{})
defer close(done)
go runLivenessWatchdogLoop(tick, done, time.Minute, func(args ...any) {})
stamp := time.Now().Add(48 * time.Hour) // guaranteed > before
select {
case tick <- stamp:
case <-time.After(time.Second):
t.Fatal("tick blocked")
}
// The loop publishes the clock; poll for it to land.
deadline := time.Now().Add(2 * time.Second)
var got int64
for time.Now().Before(deadline) {
got = WatchdogLastTickUnix()
if got >= stamp.Unix() {
return
}
time.Sleep(5 * time.Millisecond)
}
t.Fatalf("WatchdogLastTickUnix() did not advance to tick timestamp (#1749); before=%d after=%d want≥%d",
before, got, stamp.Unix())
}
// TestMQTTStallWatchdog_LoopRecoversFromPanicInEmit_EscalationPath_1749
// (RED on the prior commit): a panic inside emit on the ESCALATION
// path (maybeEscalateDisconnected → emit) must NOT kill the watchdog
// loop. The fix in b3c75ca3 moved maybeEscalateDisconnected INSIDE the
// per-source IIFE that has defer/recover, so a panic on this path is
// recovered alongside panics on the processLivenessTransition path.
func TestMQTTStallWatchdog_LoopRecoversFromPanicInEmit_EscalationPath_1749(t *testing.T) {
defer snapshotAndResetRegistry(t)()
threshold := 1 * time.Minute
var reconnectCalls int32
s := &SourceLivenessState{
Tag: "panic-escalation",
Broker: "tcp://x:1883",
IsConnectedFn: func() bool { return false }, // disconnected
ForceReconnectFn: func() {
atomic.AddInt32(&reconnectCalls, 1)
},
}
// Pre-stamp DisconnectedSinceUnix past the escalation threshold
// (disconnectedReconnectMultiplier × threshold = 5 × 1min = 5min).
// Set it 10 minutes in the past so escalation fires immediately.
atomic.StoreInt64(&s.DisconnectedSinceUnix, time.Now().Add(-10*time.Minute).Unix())
atomic.StoreInt64(&s.StartedAt, time.Now().Add(-20*time.Minute).Unix())
if err := registerLivenessState(s); err != nil {
t.Fatalf("setup: %v", err)
}
// Track which tick we're on to control panic timing.
var tickNum int32
var mu sync.Mutex
var emitCalls int
emit := func(args ...any) {
mu.Lock()
emitCalls++
mu.Unlock()
// Tick 1 (tickNum==1): all emits succeed → ForceReconnectFn fires.
// Tick 2 (tickNum==2): panic on first emit in escalation path.
// Tick 3: proves loop survived.
if atomic.LoadInt32(&tickNum) == 2 {
panic("synthetic emit panic on ESCALATION path (#1749 round-1 finding)")
}
}
tick := make(chan time.Time)
done := make(chan struct{})
exited := make(chan struct{})
go func() {
defer func() {
_ = recover()
close(exited)
}()
runLivenessWatchdogLoop(tick, done, threshold, emit)
}()
base := time.Now()
// Tick 1: escalation fires, ForceReconnectFn is invoked, no panic.
atomic.StoreInt32(&tickNum, 1)
select {
case tick <- base:
case <-time.After(time.Second):
t.Fatal("tick 1 blocked")
}
time.Sleep(150 * time.Millisecond) // let goroutine with ForceReconnectFn run
if rc := atomic.LoadInt32(&reconnectCalls); rc < 1 {
t.Fatalf("ForceReconnectFn must be invoked ≥1 time after tick 1 (got %d)", rc)
}
// Reset state so tick 2 also triggers escalation.
atomic.StoreInt64(&s.DisconnectedSinceUnix, base.Add(-10*time.Minute).Unix())
atomic.StoreInt64(&s.LastForceReconnectUnix, 0)
// Tick 2: emit panics on escalation path. On unfixed code this
// kills the loop because maybeEscalateDisconnected is outside IIFE.
atomic.StoreInt32(&tickNum, 2)
select {
case tick <- base.Add(time.Second):
case <-time.After(time.Second):
t.Fatal("tick 2 blocked")
}
time.Sleep(100 * time.Millisecond)
select {
case <-exited:
t.Fatal("watchdog loop died from panic in emit on ESCALATION path (#1749 round-1) — maybeEscalateDisconnected is not panic-protected")
default:
}
// Tick 3: proves loop survived the panic on tick 2.
atomic.StoreInt32(&tickNum, 3)
atomic.StoreInt64(&s.DisconnectedSinceUnix, base.Add(-10*time.Minute).Unix())
atomic.StoreInt64(&s.LastForceReconnectUnix, 0)
select {
case tick <- base.Add(2 * time.Second):
case <-exited:
t.Fatal("watchdog loop died; tick 3 cannot be delivered")
case <-time.After(time.Second):
t.Fatal("tick 3 blocked — loop dead")
}
time.Sleep(100 * time.Millisecond)
mu.Lock()
got := emitCalls
mu.Unlock()
// Tick 1: 2 emits (escalation + "forcing reconnect") + 1 from goroutine = 3
// Tick 2: 1 emit (panic) = partial
// Tick 3: ≥1 emit
// Total should be ≥4 if loop survived
if got < 4 {
t.Fatalf("emit must be called ≥4 times across 3 ticks (got %d) — loop did not survive escalation panic", got)
}
close(done)
select {
case <-exited:
case <-time.After(time.Second):
t.Fatal("loop did not exit after done signal")
}
}
+276
View File
@@ -0,0 +1,276 @@
package main
import (
"bytes"
"encoding/json"
"io"
"log"
"os"
"strings"
"sync"
"sync/atomic"
"testing"
"time"
)
// PR #1810 round-1 follow-ups. Tests for the must-fix findings from
// adversarial / Taleb / Kent Beck reviews on the #1749 fix:
//
// - C5: panic-recover writes to os.Stderr, NOT log.Printf (Taleb #2).
// - C6: WatchdogPanicCount increments on recovered panic and is
// surfaced via the snapshot path (Taleb #3).
// - C7: per-source jitter on escalation prevents thundering-herd
// (Taleb #4) — two sources escalating on the same tick yield
// different escalation gaps.
// - B5/C3: MarkReconnected clears DisconnectedSinceUnix so the next
// disconnect starts its escalation timer from scratch (Taleb #5).
// - C4 rewrite: throttled escalation WARN — the escalation log line
// is bounded by forceReconnectThrottle even across many ticks
// past the boundary (adv #2 + Taleb #1).
// captureStderr runs fn while redirecting os.Stderr to a pipe; returns
// whatever fn wrote to stderr. log.Printf's default writer is also
// pointed at the pipe so the test can prove a string did NOT go through
// log.Printf vs DID go to os.Stderr.
func captureStderrAndLog(t *testing.T, fn func()) (stderrBytes string, logBytes string) {
t.Helper()
rStd, wStd, err := os.Pipe()
if err != nil {
t.Fatalf("pipe: %v", err)
}
rLog, wLog, err := os.Pipe()
if err != nil {
t.Fatalf("pipe: %v", err)
}
origStderr := os.Stderr
origLogOut := log.Writer()
os.Stderr = wStd
log.SetOutput(wLog)
t.Cleanup(func() {
os.Stderr = origStderr
log.SetOutput(origLogOut)
})
var stdBuf, logBuf bytes.Buffer
var wg sync.WaitGroup
wg.Add(2)
go func() { defer wg.Done(); _, _ = io.Copy(&stdBuf, rStd) }()
go func() { defer wg.Done(); _, _ = io.Copy(&logBuf, rLog) }()
fn()
_ = wStd.Close()
_ = wLog.Close()
wg.Wait()
return stdBuf.String(), logBuf.String()
}
// C5: panic recovery message MUST go to os.Stderr directly, NOT through
// log.Printf — log's writer is the suspected cause of the panic
// (blocked log pipe / full JSON-file driver) and using it for the
// recovery message risks re-panicking on the same broken sink.
func TestWatchdog_PanicRecoverWritesToStderrNotLog_1810(t *testing.T) {
defer snapshotAndResetRegistry(t)()
threshold := 1 * time.Minute
s := &SourceLivenessState{
Tag: "panic-sink-1810",
Broker: "tcp://x:1883",
IsConnectedFn: func() bool { return true },
}
atomic.StoreInt64(&s.LastMessageUnix, time.Now().Add(-10*time.Minute).Unix())
atomic.StoreInt64(&s.StartedAt, time.Now().Add(-20*time.Minute).Unix())
if err := registerLivenessState(s); err != nil {
t.Fatalf("setup: %v", err)
}
emit := func(args ...any) {
panic("synthetic panic to drive recover")
}
stderrOut, logOut := captureStderrAndLog(t, func() {
tick, done, exited := setupWatchdogTestLoop(t, threshold, emit)
sendTickOrFail(t, tick, time.Now(), time.Second, "panic-sink tick")
time.Sleep(150 * time.Millisecond)
close(done)
<-exited
})
if !strings.Contains(stderrOut, "WATCHDOG RECOVERED") {
t.Fatalf("expected WATCHDOG RECOVERED line on os.Stderr; stderr=%q log=%q", stderrOut, logOut)
}
if strings.Contains(logOut, "WATCHDOG RECOVERED") {
t.Fatalf("recovery message MUST NOT go through log.Printf (#1810 Taleb #2); log=%q", logOut)
}
}
// C6: WatchdogPanicCount must increment on a recovered panic and the
// value must be reachable via the package-level accessor (which the
// stats snapshot reads).
func TestWatchdog_PanicCountIncrementsOnRecover_1810(t *testing.T) {
defer snapshotAndResetRegistry(t)()
before := WatchdogPanicCount()
threshold := 1 * time.Minute
s := &SourceLivenessState{
Tag: "panic-count-1810",
Broker: "tcp://x:1883",
IsConnectedFn: func() bool { return true },
}
atomic.StoreInt64(&s.LastMessageUnix, time.Now().Add(-10*time.Minute).Unix())
atomic.StoreInt64(&s.StartedAt, time.Now().Add(-20*time.Minute).Unix())
if err := registerLivenessState(s); err != nil {
t.Fatalf("setup: %v", err)
}
emit := func(args ...any) { panic("boom") }
// Suppress stderr output during the recover so test output stays clean.
_, _ = captureStderrAndLog(t, func() {
tick, done, exited := setupWatchdogTestLoop(t, threshold, emit)
sendTickOrFail(t, tick, time.Now(), time.Second, "panic-count tick")
time.Sleep(150 * time.Millisecond)
close(done)
<-exited
})
after := WatchdogPanicCount()
if after <= before {
t.Fatalf("WatchdogPanicCount must advance on recovered panic (#1810 Taleb #3); before=%d after=%d", before, after)
}
}
// C7: per-source jitter — two sources crossing the escalation boundary
// on the same tick must NOT escalate at the same threshold gap value.
// This proves the hash-based jitter offset spreads escalations and
// avoids a thundering-herd reconnect when N sources share an upstream
// broker outage.
func TestWatchdog_EscalationJitterPerSource_1810(t *testing.T) {
threshold := 60 * time.Second
gapA := disconnectedEscalationGap("source-a", threshold)
gapB := disconnectedEscalationGap("source-b", threshold)
if gapA == gapB {
t.Fatalf("expected per-source jitter on escalation gap (#1810 Taleb #4); gapA=%s gapB=%s", gapA, gapB)
}
// Both must be at least the unjittered base.
base := time.Duration(disconnectedReconnectMultiplier) * threshold
if gapA < base || gapB < base {
t.Fatalf("escalation gap must be ≥ base (%s); gapA=%s gapB=%s", base, gapA, gapB)
}
// Jitter must be bounded (≤ base + 30s sanity).
maxJitter := base + 30*time.Second
if gapA > maxJitter || gapB > maxJitter {
t.Fatalf("escalation jitter exceeded bound; gapA=%s gapB=%s max=%s", gapA, gapB, maxJitter)
}
}
// C3 / B5: MarkReconnected must clear DisconnectedSinceUnix. On unfixed
// code, a post-recovery LivenessDisconnected immediately satisfies
// (now - DisconnectedSinceUnix > multiplier×threshold) and force-
// reconnects on the FIRST tick — making escalation a flap-trigger
// instead of a persistent-failure-trigger.
func TestMarkReconnected_ClearsDisconnectedSinceUnix_1810(t *testing.T) {
s := &SourceLivenessState{Tag: "recovery-clear-1810"}
atomic.StoreInt64(&s.DisconnectedSinceUnix, time.Now().Add(-1*time.Hour).Unix())
s.MarkReconnected(time.Now())
if got := atomic.LoadInt64(&s.DisconnectedSinceUnix); got != 0 {
t.Fatalf("MarkReconnected must clear DisconnectedSinceUnix (#1810 Taleb #5); got %d", got)
}
}
// C4 (rewrite): the ESCALATION WARN log line must be throttled. On
// unfixed code (round 0) every tick past the boundary emitted the WARN;
// for a 1m scan interval and a 1h outage that's 55+ duplicates. This
// test does NOT pre-stamp DisconnectedSinceUnix — it starts from the
// first-observation branch and drives ticks past the boundary,
// asserting both that reconnects are throttled AND that the escalation
// WARN log count stays bounded.
func TestWatchdog_EscalationWarnThrottled_1810(t *testing.T) {
defer snapshotAndResetRegistry(t)()
threshold := 1 * time.Second
var reconnectCount atomic.Int32
s := &SourceLivenessState{
Tag: "warn-throttle-1810",
Broker: "tcp://x:1883",
IsConnectedFn: func() bool { return false },
ForceReconnectFn: func() { reconnectCount.Add(1) },
}
if err := registerLivenessState(s); err != nil {
t.Fatalf("setup: %v", err)
}
var mu sync.Mutex
var warnCount int
emit := func(args ...any) {
mu.Lock()
defer mu.Unlock()
for _, a := range args {
if str, ok := a.(string); ok && strings.Contains(str, "WATCHDOG ESCALATION") {
warnCount++
}
}
}
tick := make(chan time.Time)
done := make(chan struct{})
defer close(done)
go runLivenessWatchdogLoop(tick, done, threshold, emit)
base := time.Now()
// First tick: stamps DisconnectedSinceUnix (no escalation yet).
tick <- base
// Subsequent ticks: drive 60 ticks at 1s each, far past the
// escalation boundary (5×threshold = 5s + ≤30s jitter). On
// unfixed code this would emit 50+ WARN lines.
for i := 1; i < 60; i++ {
tick <- base.Add(time.Duration(i) * time.Second)
}
time.Sleep(200 * time.Millisecond)
mu.Lock()
got := warnCount
mu.Unlock()
// One escalation per throttle window. With throttle=60s and the
// test spanning ~59s of fabricated wall clock, at most ~2 WARNs
// should fire (boundary cross + possibly one more if jitter is
// minimal). 5 is a generous upper bound; the original bug
// produced 50+.
if got > 5 {
t.Fatalf("escalation WARN must be throttled (#1810 adv #2 / Taleb #1); got %d emits across 60 ticks past the boundary", got)
}
if got < 1 {
t.Fatalf("escalation WARN must fire at least once when paho stays disconnected past boundary; got 0")
}
if rc := reconnectCount.Load(); rc < 1 {
t.Fatalf("ForceReconnectFn must be invoked at least once; got %d", rc)
}
}
// C2: round-trip — IngestorStatsSnapshot.WatchdogLastTickUnix and
// WatchdogPanicCount must serialize through JSON and deserialize via
// the server's envelope shape.
func TestIngestorStatsSnapshot_WatchdogFieldsRoundTrip_1810(t *testing.T) {
snap := IngestorStatsSnapshot{
SampledAt: time.Now().UTC().Format(time.RFC3339),
WatchdogLastTickUnix: 1700000000,
WatchdogPanicCount: 42,
}
b, err := json.Marshal(snap)
if err != nil {
t.Fatalf("marshal: %v", err)
}
if !bytes.Contains(b, []byte(`"watchdogLastTickUnix":1700000000`)) {
t.Fatalf("watchdogLastTickUnix missing from JSON: %s", string(b))
}
if !bytes.Contains(b, []byte(`"watchdogPanicCount":42`)) {
t.Fatalf("watchdogPanicCount missing from JSON: %s", string(b))
}
var back IngestorStatsSnapshot
if err := json.Unmarshal(b, &back); err != nil {
t.Fatalf("unmarshal: %v", err)
}
if back.WatchdogLastTickUnix != 1700000000 || back.WatchdogPanicCount != 42 {
t.Fatalf("round-trip mismatch: %+v", back)
}
}
@@ -0,0 +1,47 @@
package main
import (
"testing"
"time"
)
// setupWatchdogTestLoop spawns runLivenessWatchdogLoop with a panic
// safety net so an unrecovered panic in the loop (the historical bug
// shape #1749 / #1810 round-1 gates against) does NOT crash the test
// binary. Returns:
// - tick: send fabricated timestamps to drive the loop
// - done: caller closes to ask the loop to exit cleanly
// - exited: closed by the helper after the loop returns (whether by
// done-signal, normal channel close, or panic propagation)
//
// Extracted as part of #1810 round-1 (adv #4) — four tests in this
// package previously open-coded the same scaffolding with subtle
// variations.
func setupWatchdogTestLoop(t *testing.T, threshold time.Duration, emit func(...any)) (tick chan time.Time, done chan struct{}, exited chan struct{}) {
t.Helper()
tick = make(chan time.Time)
done = make(chan struct{})
exited = make(chan struct{})
go func() {
defer func() {
_ = recover() // safety net; the production loop is what we are asserting on
close(exited)
}()
runLivenessWatchdogLoop(tick, done, threshold, emit)
}()
return tick, done, exited
}
// sendTickOrFail pushes one fabricated timestamp into the loop's tick
// channel and fails the test if the loop does not consume it within
// timeout. Use this everywhere a test wants to step the watchdog by
// exactly one tick — open-coded select-default-Fatal patterns are
// repetitive and easy to get wrong (off-by-one timeouts).
func sendTickOrFail(t *testing.T, tick chan<- time.Time, stamp time.Time, timeout time.Duration, label string) {
t.Helper()
select {
case tick <- stamp:
case <-time.After(timeout):
t.Fatalf("%s: tick blocked after %s — loop dead?", label, timeout)
}
}
+36 -15
View File
@@ -62,6 +62,25 @@ type IngestorStatsSnapshot struct {
// counter view consumed by cmd/server's /api/mqtt/status handler.
// Additive; omitempty so older server builds ignore it.
SourceStatuses []SourceStatusSnapshot `json:"source_statuses,omitempty"`
// WatchdogLastTickUnix (#1749) is the unix-seconds timestamp of the
// most recent runLivenessWatchdogLoop tick. Surfaced via
// /api/mqtt/status so external monitoring can assert the watchdog
// goroutine is still alive — a value older than ~2× the watchdog
// scan interval (typically 60s) indicates the watchdog itself has
// wedged or panicked. 0 means the watchdog has never ticked
// (e.g. cold start before the first scan). Additive — omitempty
// so older server builds ignore it.
WatchdogLastTickUnix int64 `json:"watchdogLastTickUnix,omitempty"`
// WatchdogPanicCount (#1810 round-1) is the running total of
// recovered panics inside the watchdog per-source work IIFE.
// Exposed alongside WatchdogLastTickUnix because the tick clock is
// stamped BEFORE the per-source work — a loop that panics on every
// source still advances WatchdogLastTickUnix and looks healthy by
// that signal alone. A rapidly-growing WatchdogPanicCount means
// the loop is alive but per-source processing is broken (typically
// a panic in emit / log sink). Monotonic; 0 means no recovered
// panics yet. Additive — omitempty so older server builds ignore.
WatchdogPanicCount int64 `json:"watchdogPanicCount,omitempty"`
}
// SourceLivenessSnapshot is the per-source two-clock view exposed for
@@ -237,21 +256,23 @@ func StartStatsFileWriter(s *Store, interval time.Duration) {
ioRate := procIORate(prevIO, curIO, stamp)
prevIO = curIO
snap := IngestorStatsSnapshot{
SampledAt: stamp,
TxInserted: s.Stats.TransmissionsInserted.Load(),
ObsInserted: s.Stats.ObservationsInserted.Load(),
DuplicateTx: s.Stats.DuplicateTransmissions.Load(),
NodeUpserts: s.Stats.NodeUpserts.Load(),
ObserverUpserts: s.Stats.ObserverUpserts.Load(),
WriteErrors: s.Stats.WriteErrors.Load(),
SignatureDrops: s.Stats.SignatureDrops.Load(),
WALCommits: s.Stats.WALCommits.Load(),
GroupCommitFlushes: 0, // group commit reverted (refs #1129)
BackfillUpdates: s.Stats.SnapshotBackfills(),
ProcIO: ioRate,
WriterPerf: s.WriterStatsSnapshot(),
SourceLiveness: SnapshotLivenessClocks(),
SourceStatuses: SnapshotSourceStatuses(tickAt),
SampledAt: stamp,
TxInserted: s.Stats.TransmissionsInserted.Load(),
ObsInserted: s.Stats.ObservationsInserted.Load(),
DuplicateTx: s.Stats.DuplicateTransmissions.Load(),
NodeUpserts: s.Stats.NodeUpserts.Load(),
ObserverUpserts: s.Stats.ObserverUpserts.Load(),
WriteErrors: s.Stats.WriteErrors.Load(),
SignatureDrops: s.Stats.SignatureDrops.Load(),
WALCommits: s.Stats.WALCommits.Load(),
GroupCommitFlushes: 0, // group commit reverted (refs #1129)
BackfillUpdates: s.Stats.SnapshotBackfills(),
ProcIO: ioRate,
WriterPerf: s.WriterStatsSnapshot(),
SourceLiveness: SnapshotLivenessClocks(),
SourceStatuses: SnapshotSourceStatuses(tickAt),
WatchdogLastTickUnix: WatchdogLastTickUnix(),
WatchdogPanicCount: WatchdogPanicCount(),
}
buf.Reset()
if err := enc.Encode(&snap); err != nil {
+22 -2
View File
@@ -107,13 +107,31 @@ type MqttSourceStatus struct {
type MqttStatusResponse struct {
Sources []MqttSourceStatus `json:"sources"`
SampleAt string `json:"sampleAt"`
// WatchdogLastTickUnix (#1749) is the unix-seconds timestamp of the
// most recent ingestor watchdog tick. Surfaced so external monitoring
// can detect a wedged watchdog goroutine (value older than ~2× the
// scan interval — typically 60s — means the watchdog itself died,
// not just one source). 0 / omitted: ingestor has never ticked yet
// or is running an older build that did not publish this field.
WatchdogLastTickUnix int64 `json:"watchdogLastTickUnix,omitempty"`
// WatchdogPanicCount (#1810 round-1) is the running total of
// recovered panics inside the ingestor's watchdog per-source work
// IIFE. Surfaced alongside WatchdogLastTickUnix because the tick
// clock is stamped BEFORE per-source work — a loop panicking on
// every source still advances the tick and looks healthy by tick
// alone. A rapidly-growing value means the watchdog is alive but
// per-source processing is broken. 0 / omitted: no recovered
// panics OR older ingestor build.
WatchdogPanicCount int64 `json:"watchdogPanicCount,omitempty"`
}
// ingestorMqttStatusEnvelope is the partial shape the server decodes from
// the ingestor stats file (additive — older ingestors omit the field).
type ingestorMqttStatusEnvelope struct {
SampledAt string `json:"sampledAt"`
SourceStatuses []MqttSourceStatus `json:"source_statuses"`
SampledAt string `json:"sampledAt"`
SourceStatuses []MqttSourceStatus `json:"source_statuses"`
WatchdogLastTickUnix int64 `json:"watchdogLastTickUnix"`
WatchdogPanicCount int64 `json:"watchdogPanicCount"`
}
// handleMqttStatus serves GET /api/mqtt/status. Reads the ingestor stats
@@ -133,6 +151,8 @@ func (s *Server) handleMqttStatus(w http.ResponseWriter, r *http.Request) {
return
}
resp.SampleAt = env.SampledAt
resp.WatchdogLastTickUnix = env.WatchdogLastTickUnix
resp.WatchdogPanicCount = env.WatchdogPanicCount
for _, src := range env.SourceStatuses {
src.Broker = maskBrokerURL(src.Broker)
// Broker libraries occasionally quote the failing URL in the
+55
View File
@@ -0,0 +1,55 @@
package main
import (
"encoding/json"
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"testing"
)
// PR #1810 round-1: /api/mqtt/status must surface WatchdogLastTickUnix
// and WatchdogPanicCount so external monitoring can detect both a
// wedged watchdog goroutine and a panic-looping watchdog (the loop
// stamps WatchdogLastTickUnix BEFORE per-source work, so a panic on
// every source still advances the tick — the panic counter is the
// distinguishing signal).
func TestMqttStatus_ExposesWatchdogFields_1810(t *testing.T) {
tmp := t.TempDir()
statsPath := filepath.Join(tmp, "ingestor-stats.json")
t.Setenv("CORESCOPE_INGESTOR_STATS", statsPath)
stub := map[string]any{
"sampledAt": "2026-06-30T12:30:00Z",
"source_statuses": []map[string]any{},
"watchdogLastTickUnix": int64(1751313000),
"watchdogPanicCount": int64(7),
}
data, err := json.Marshal(stub)
if err != nil {
t.Fatalf("marshal stub: %v", err)
}
if err := os.WriteFile(statsPath, data, 0o600); err != nil {
t.Fatalf("write stub: %v", err)
}
srv := &Server{}
req := httptest.NewRequest(http.MethodGet, "/api/mqtt/status", nil)
rec := httptest.NewRecorder()
srv.handleMqttStatus(rec, req)
if rec.Code != http.StatusOK {
t.Fatalf("status = %d, want 200; body=%s", rec.Code, rec.Body.String())
}
var resp MqttStatusResponse
if err := json.Unmarshal(rec.Body.Bytes(), &resp); err != nil {
t.Fatalf("unmarshal: %v; body=%s", err, rec.Body.String())
}
if resp.WatchdogLastTickUnix != 1751313000 {
t.Errorf("WatchdogLastTickUnix = %d, want 1751313000; body=%s", resp.WatchdogLastTickUnix, rec.Body.String())
}
if resp.WatchdogPanicCount != 7 {
t.Errorf("WatchdogPanicCount = %d, want 7; body=%s", resp.WatchdogPanicCount, rec.Body.String())
}
}