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
79 changes: 63 additions & 16 deletions cmd/ingestor/db.go
Original file line number Diff line number Diff line change
Expand Up @@ -1413,8 +1413,68 @@ func (s *Store) UpsertObserverAt(id, name, iata string, meta *ObserverMeta, last
}
normalizedIATA := strings.TrimSpace(strings.ToUpper(iata))

var model, firmware, clientVersion, radio interface{}
var batteryMv, uptimeSecs, noiseFloor, canRelay interface{}
model, firmware, clientVersion, radio, batteryMv, uptimeSecs, noiseFloor, canRelay := observerMetaColumns(meta)

_, err := s.stmtUpsertObserver.Exec(
id, name, normalizedIATA, lastSeen, lastSeen, model, firmware, clientVersion, radio, batteryMv, uptimeSecs, noiseFloor, canRelay, canRelay,
name, normalizedIATA, ingestNow, lastSeen, model, firmware, clientVersion, radio, batteryMv, uptimeSecs, noiseFloor, canRelay, canRelay,
)
if err != nil {
s.Stats.WriteErrors.Add(1)
return err
}
s.Stats.ObserverUpserts.Add(1)

// Reactivate if this observer was previously marked inactive
s.db.Exec(`UPDATE observers SET inactive = 0 WHERE id = ? AND inactive = 1`, id)
return nil
}

// UpsertObserverRetained applies the metadata from a RETAINED status message
// without treating it as a sign of life. The broker replays retained messages
// on every subscribe, so an ingestor restart would otherwise stamp last_seen
// with the restart time for every observer that ever published one — making
// dead observers permanently un-ageable by RemoveStaleObservers, and undoing
// any inactive flag it did manage to set.
//
// Deliberately narrower than UpsertObserverAt: it updates metadata columns on
// an existing row only. It does not advance last_seen, does not clear
// inactive, does not bump packet_count, and does not INSERT — a retained-only
// observer the analyzer has never heard from live describes a past that may be
// months old and does not belong in the list. A live message from the same
// observer arrives moments later and creates the row through the normal path.
func (s *Store) UpsertObserverRetained(id, name, iata string, meta *ObserverMeta) error {
normalizedIATA := strings.TrimSpace(strings.ToUpper(iata))
model, firmware, clientVersion, radio, batteryMv, uptimeSecs, noiseFloor, canRelay := observerMetaColumns(meta)

_, err := s.db.Exec(`
UPDATE observers SET
name = COALESCE(?, name),
iata = COALESCE(?, iata),
model = COALESCE(?, model),
firmware = COALESCE(?, firmware),
client_version = COALESCE(?, client_version),
radio = COALESCE(?, radio),
battery_mv = COALESCE(?, battery_mv),
uptime_secs = COALESCE(?, uptime_secs),
noise_floor = COALESCE(?, noise_floor),
can_relay = COALESCE(?, can_relay),
can_relay_seen = CASE WHEN ? IS NULL THEN can_relay_seen ELSE 1 END
WHERE id = ?`,
name, normalizedIATA, model, firmware, clientVersion, radio,
batteryMv, uptimeSecs, noiseFloor, canRelay, canRelay, id,
)
if err != nil {
s.Stats.WriteErrors.Add(1)
return err
}
return nil
}

// observerMetaColumns flattens an *ObserverMeta into the driver args the
// observer upserts bind. A nil field stays nil so the COALESCE in the SQL
// leaves the existing column untouched (#1290).
func observerMetaColumns(meta *ObserverMeta) (model, firmware, clientVersion, radio, batteryMv, uptimeSecs, noiseFloor, canRelay interface{}) {
if meta != nil {
if meta.Model != nil {
model = *meta.Model
Expand Down Expand Up @@ -1449,20 +1509,7 @@ func (s *Store) UpsertObserverAt(id, name, iata string, meta *ObserverMeta, last
}
}
}

_, err := s.stmtUpsertObserver.Exec(
id, name, normalizedIATA, lastSeen, lastSeen, model, firmware, clientVersion, radio, batteryMv, uptimeSecs, noiseFloor, canRelay, canRelay,
name, normalizedIATA, ingestNow, lastSeen, model, firmware, clientVersion, radio, batteryMv, uptimeSecs, noiseFloor, canRelay, canRelay,
)
if err != nil {
s.Stats.WriteErrors.Add(1)
return err
}
s.Stats.ObserverUpserts.Add(1)

// Reactivate if this observer was previously marked inactive
s.db.Exec(`UPDATE observers SET inactive = 0 WHERE id = ? AND inactive = 1`, id)
return nil
return model, firmware, clientVersion, radio, batteryMv, uptimeSecs, noiseFloor, canRelay
}

// Close checkpoints the WAL and closes the database.
Expand Down
16 changes: 16 additions & 0 deletions cmd/ingestor/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -632,6 +632,22 @@ func handleMessage(store *Store, tag string, source MQTTSource, m mqtt.Message,
name, _ := msg["origin"].(string)
iata := parts[1]
meta := extractObserverMeta(msg)
// A replayed status message is the broker handing us the observer's
// last published snapshot — it is not evidence the observer is alive
// now. Stamping last_seen from it resurrects dead observers, so the
// replay path only refreshes metadata.
//
// Two ways to recognise one: the retain flag (our own subscribe), and
// the payload itself (a replay that reached us through the mosquitto
// bridge, where the flag does not survive the hop — see
// statusIsLiveness).
if m.Retained() || !statusIsLiveness(msg, time.Now().UTC()) {
if err := store.UpsertObserverRetained(observerID, name, iata, meta); err != nil {
log.Printf("MQTT [%s] retained observer status error: %v", tag, err)
}
log.Print(formatStatusLog(tag, firstNonEmpty(name, observerID), iata))
return
}
// observer.last_seen is "when did the analyzer last hear from this
// observer" — fundamentally an ingest-time question. Passing "" makes
// UpsertObserverAt use time.Now(), independent of the envelope timestamp
Expand Down
7 changes: 4 additions & 3 deletions cmd/ingestor/main_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -105,13 +105,14 @@ func TestUnixTime(t *testing.T) {

// mockMessage implements mqtt.Message for testing handleMessage
type mockMessage struct {
topic string
payload []byte
topic string
payload []byte
retained bool
}

func (m *mockMessage) Duplicate() bool { return false }
func (m *mockMessage) Qos() byte { return 0 }
func (m *mockMessage) Retained() bool { return false }
func (m *mockMessage) Retained() bool { return m.retained }
func (m *mockMessage) Topic() string { return m.topic }
func (m *mockMessage) MessageID() uint16 { return 0 }
func (m *mockMessage) Payload() []byte { return m.payload }
Expand Down
207 changes: 207 additions & 0 deletions cmd/ingestor/retained_status_gaps_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,207 @@
package main

import (
"strings"
"testing"
"time"
)

// Gap-closing tests for #1885. An independent mutation run against the
// original seven tests left nine mutants alive; these kill the five that
// matter. Each test names the mutation it kills so a future reader can tell
// what it is load-bearing for.
//
// The four remaining survivors are deliberate: two concern the 24h constant's
// exact value (any value between "a few hours" and "months" passes, and the
// constant is documented as a generous margin rather than a threshold), and
// two concern Stats.ObserverUpserts, which is a counter with no behavioural
// consequence on either path.

// Kills: re-adding `packet_count = packet_count + 1` to UpsertObserverRetained.
//
// The count is the observer's real traffic volume. A replay storm that bumped
// it would inflate every dead observer's count by one per ingestor restart,
// and nothing else in the suite looks at the column on this path.
func TestRetainedStatusDoesNotBumpPacketCount(t *testing.T) {
store := newTestStore(t)
seedObserver(t, store, "obs-count", 60)
if _, err := store.db.Exec(`UPDATE observers SET packet_count = 7 WHERE id = ?`, "obs-count"); err != nil {
t.Fatal(err)
}

handleMessage(store, "test", MQTTSource{Name: "test"},
retainedStatusMsg("meshcore/LAX/obs-count/status", `{"origin":"obs-count","noise_floor":-95.5}`),
nil, nil, &Config{})

var got int
if err := store.db.QueryRow(`SELECT packet_count FROM observers WHERE id = ?`, "obs-count").Scan(&got); err != nil {
t.Fatal(err)
}
if got != 7 {
t.Errorf("packet_count = %d, want 7 (a replay is not a packet)", got)
}

// The live path must still bump it, or the guard above is just a bug.
handleMessage(store, "test", MQTTSource{Name: "test"},
liveStatusMsg("meshcore/LAX/obs-count/status",
statusPayload("obs-count", "online", time.Now().UTC())),
nil, nil, &Config{})
if err := store.db.QueryRow(`SELECT packet_count FROM observers WHERE id = ?`, "obs-count").Scan(&got); err != nil {
t.Fatal(err)
}
if got != 8 {
t.Errorf("packet_count = %d after live status, want 8 (live traffic still counts)", got)
}
}

// Kills: making can_relay_seen regress 1 -> 0, and never setting it to 1.
//
// can_relay_seen is the tristate's "we have actually observed this" bit
// (#1290). A retained replay carrying no `repeat` field must not erase an
// answer the live path already recorded.
func TestRetainedStatusPreservesCanRelaySeen(t *testing.T) {
store := newTestStore(t)
seedObserver(t, store, "obs-relay", 60)
if _, err := store.db.Exec(
`UPDATE observers SET can_relay = 1, can_relay_seen = 1 WHERE id = ?`, "obs-relay"); err != nil {
t.Fatal(err)
}

// No `repeat` key: meta.CanRelay stays nil, so the CASE must leave the
// flag alone rather than reset it.
handleMessage(store, "test", MQTTSource{Name: "test"},
retainedStatusMsg("meshcore/LAX/obs-relay/status", `{"origin":"obs-relay","noise_floor":-95.5}`),
nil, nil, &Config{})

var seen, canRelay int
if err := store.db.QueryRow(
`SELECT can_relay_seen, can_relay FROM observers WHERE id = ?`, "obs-relay").
Scan(&seen, &canRelay); err != nil {
t.Fatal(err)
}
if seen != 1 {
t.Errorf("can_relay_seen = %d, want 1 (a replay without `repeat` must not un-observe it)", seen)
}
if canRelay != 1 {
t.Errorf("can_relay = %d, want 1 (unchanged by a replay that does not report it)", canRelay)
}

// A replay that DOES carry `repeat` still records the observation: the
// metadata half of the retained path is supposed to keep working.
store2 := newTestStore(t)
seedObserver(t, store2, "obs-relay2", 60)
if _, err := store2.db.Exec(
`UPDATE observers SET can_relay = NULL, can_relay_seen = 0 WHERE id = ?`, "obs-relay2"); err != nil {
t.Fatal(err)
}
handleMessage(store2, "test", MQTTSource{Name: "test"},
retainedStatusMsg("meshcore/LAX/obs-relay2/status",
`{"origin":"obs-relay2","repeat":"true"}`),
nil, nil, &Config{})
if err := store2.db.QueryRow(
`SELECT can_relay_seen, can_relay FROM observers WHERE id = ?`, "obs-relay2").
Scan(&seen, &canRelay); err != nil {
t.Fatal(err)
}
if seen != 1 || canRelay != 1 {
t.Errorf("can_relay_seen/can_relay = %d/%d, want 1/1 (a reported value is still recorded)", seen, canRelay)
}
}

// Kills: replacing strings.EqualFold with a case-sensitive ==.
//
// The whole guard turns on one string comparison. If it became
// case-sensitive, any firmware publishing "ONLINE" would have its liveness
// silently suppressed — the observer would age out despite being alive, which
// is the opposite of the bug this change fixes and far worse than it.
func TestStatusOnlineComparisonIsCaseInsensitive(t *testing.T) {
now := time.Now().UTC()
for _, s := range []string{"online", "ONLINE", "OnLiNe", "Online"} {
if !statusIsLiveness(map[string]interface{}{"status": s}, now) {
t.Errorf("statusIsLiveness(status=%q) = false, want true", s)
}
}
for _, s := range []string{"offline", "OFFLINE", "OffLine"} {
if statusIsLiveness(map[string]interface{}{"status": s}, now) {
t.Errorf("statusIsLiveness(status=%q) = true, want false", s)
}
}
}

// Kills: dropping the `s != ""` guard.
//
// An empty status string is not a claim of being offline — it is the absence
// of a claim, and the guard exists so a payload that carries the key but no
// value keeps the previous behaviour instead of being silently discarded.
func TestEmptyStatusStringIsNotAnOfflineClaim(t *testing.T) {
now := time.Now().UTC()
if !statusIsLiveness(map[string]interface{}{"status": ""}, now) {
t.Error(`statusIsLiveness(status="") = false, want true (no claim is not an offline claim)`)
}
// A non-string status is likewise unjudgeable and must not suppress.
if !statusIsLiveness(map[string]interface{}{"status": 1}, now) {
t.Error("statusIsLiveness(status=1) = false, want true (unjudgeable, keep previous behaviour)")
}
}

// Kills: flipping `!t.Before(cutoff)` to `t.After(cutoff)`.
//
// The two differ only exactly at the cutoff. The existing tests use 4-month
// staleness and "now", so nothing pinned which side of the boundary is
// inclusive.
func TestStatusLivenessAgeBoundaryIsInclusive(t *testing.T) {
// Truncate to whole seconds: the payload timestamp round-trips through
// RFC3339, which carries no sub-second part, so an untruncated `now`
// would put the "exactly at the cutoff" case a few hundred nanoseconds
// on the wrong side of the comparison and test the formatter instead of
// the boundary.
now := time.Now().UTC().Truncate(time.Second)
cutoff := now.Add(-statusLivenessMaxAge)

cases := []struct {
name string
ts time.Time
want bool
}{
{"exactly at the cutoff", cutoff, true},
{"one second inside", cutoff.Add(time.Second), true},
{"one second outside", cutoff.Add(-time.Second), false},
}
for _, c := range cases {
msg := map[string]interface{}{
"status": "online",
"timestamp": c.ts.Format(time.RFC3339),
}
if got := statusIsLiveness(msg, now); got != c.want {
t.Errorf("%s: statusIsLiveness = %v, want %v", c.name, got, c.want)
}
}
}

// Kills: dropping the TrimSpace/ToUpper normalisation of iata in
// UpsertObserverRetained.
//
// The live path normalises, so a replay writing a raw topic segment would
// make the same observer's region disagree with itself depending on which
// path last touched it.
func TestRetainedStatusNormalizesIATA(t *testing.T) {
store := newTestStore(t)
// Seed through the normal path so the row exists, then blank the column
// so the COALESCE has something to overwrite.
seedObserver(t, store, "obs-iata", 60)
if _, err := store.db.Exec(`UPDATE observers SET iata = 'XXX' WHERE id = ?`, "obs-iata"); err != nil {
t.Fatal(err)
}

handleMessage(store, "test", MQTTSource{Name: "test"},
retainedStatusMsg("meshcore/lax/obs-iata/status", `{"origin":"obs-iata"}`),
nil, nil, &Config{})

var got string
if err := store.db.QueryRow(`SELECT iata FROM observers WHERE id = ?`, "obs-iata").Scan(&got); err != nil {
t.Fatal(err)
}
if got != strings.ToUpper(got) || got != "LAX" {
t.Errorf("iata = %q, want %q (retained path must normalise like the live path)", got, "LAX")
}
}
Loading
Loading