From 498fbc03219ffe98aa5a39f0ec8b416b5fb79a1b Mon Sep 17 00:00:00 2001 From: Marcel Verdult <37555613+marcelverdult@users.noreply.github.com> Date: Sat, 23 May 2026 20:22:51 +0200 Subject: [PATCH] fix: ingestor uses ingest-time now() instead of observer receive time (#1233) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## Problem The ingestor stamps every stored packet with its own ingest-time `time.Now()` (`BuildPacketData` in `db.go`; channel/DM paths in `main.go`), discarding the observer receive time the uploader already puts in the MQTT envelope's `timestamp` field. `MQTTPacketMessage` had no `Timestamp` field and `handleMessage` parsed every envelope field except that one. Observers that buffer packets offline and upload hours later get every buffered packet displayed at upload time, not receive time — a 5-hour deferred upload shows packets 5 hours late. Retained messages and broker backlog hit the same skew. ## Why the envelope timestamp is trustworthy Uploaders stamp `timestamp` when the radio receives the frame and freeze it; the MQTT *message* is published late, but the `timestamp` *field* is not re-stamped at publish. A buffered packet uploaded hours late still carries its true receive time. ## Fix New `resolveRxTime` helper reads `msg["timestamp"]` and falls back to `time.Now()` only when it is missing, unparseable, or implausibly in the future. Applied to all three ingest paths (raw packet, channel, DM). No wire-format change — the field already exists. Channel/DM dedup hashes intentionally stay on ingest time, since those bridge messages carry no real packet hash and need ingest-unique input. ## Observer/node last_seen correction Packet timestamps must reflect receive time, but observer/node `last_seen` must not. `InsertTransmission` fed `data.Timestamp` (now rxTime) into `observers.last_seen` and `UpsertNode`'s `last_seen`, so a buffered upload could drag both fields backwards, and retained-message replay on MQTT reconnect could flash long-offline observers as Online. - `UpsertObserverAt` takes an explicit `lastSeen`; the status-packet and BLE companion handlers pass the resolved rxTime. `UpsertObserver` keeps its wall-clock behaviour for other callers. - All three `last_seen` writes are guarded with `MAX(MIN(existing, ingestNow), rxTime)`: `last_seen` never moves backwards from a stale retained message, and never locks in a future value. ## Naive UTC+N timestamps `resolveRxTime` rejects a timestamp only when it is >14h ahead (UTC+14 is the maximum standard offset — anything further is a genuine clock error). A timestamp that is merely in the future is soft-clamped to ingest time: a future rxTime means a live packet from a UTC+N observer whose naive local clock parses as-if UTC, not a buffered packet, so ingest time is correct and no future timestamp reaches the DB. For buffered packets from naive-clock uploaders a bounded residual offset remains (equal to the observer's UTC offset); uploaders emitting zone-aware ISO8601 everywhere would be the full cure but is a separate format change. ## Test `cmd/ingestor/rxtime_test.go` covers `parseEnvelopeTime` (zone-aware, naive, microseconds, garbage, empty) and `resolveRxTime` (plausible past used verbatim, missing/garbage/future → ingest-time fallback). The existing `TestBuildPacketData` is updated to supply an envelope timestamp and assert it propagates, since `BuildPacketData` no longer self-stamps. --- cmd/ingestor/db.go | 63 +++++++++++++++++++--------- cmd/ingestor/db_test.go | 13 +++--- cmd/ingestor/main.go | 84 +++++++++++++++++++++++++++++++++---- cmd/ingestor/rxtime_test.go | 80 +++++++++++++++++++++++++++++++++++ 4 files changed, 206 insertions(+), 34 deletions(-) create mode 100644 cmd/ingestor/rxtime_test.go diff --git a/cmd/ingestor/db.go b/cmd/ingestor/db.go index 8d01044f..441d01d8 100644 --- a/cmd/ingestor/db.go +++ b/cmd/ingestor/db.go @@ -601,7 +601,7 @@ func (s *Store) prepareStatements() error { role = COALESCE(?, role), lat = COALESCE(?, lat), lon = COALESCE(?, lon), - last_seen = ? + last_seen = MAX(MIN(COALESCE(last_seen, ''), ?), ?) `) if err != nil { return err @@ -620,7 +620,7 @@ func (s *Store) prepareStatements() error { ON CONFLICT(id) DO UPDATE SET name = COALESCE(?, name), iata = COALESCE(?, iata), - last_seen = ?, + last_seen = MAX(MIN(COALESCE(last_seen, ''), ?), ?), packet_count = packet_count + 1, model = COALESCE(?, model), firmware = COALESCE(?, firmware), @@ -639,7 +639,14 @@ func (s *Store) prepareStatements() error { return err } - s.stmtUpdateObserverLastSeen, err = s.db.Prepare("UPDATE observers SET last_seen = ?, last_packet_at = ? WHERE rowid = ?") + // Args: ingestNow, rxTime, ingestNow, rxTime, rowid + // MIN(existing, ingestNow) clamps any future value already in the DB before + // taking MAX with rxTime, so the guard never locks in a past bug's stale future. + s.stmtUpdateObserverLastSeen, err = s.db.Prepare(` + UPDATE observers SET + last_seen = MAX(MIN(COALESCE(last_seen, ''), ?), ?), + last_packet_at = MAX(MIN(COALESCE(last_packet_at, ''), ?), ?) + WHERE rowid = ?`) if err != nil { return err } @@ -673,9 +680,10 @@ func (s *Store) InsertTransmission(data *PacketData) (bool, error) { return false, nil } - now := data.Timestamp - if now == "" { - now = time.Now().UTC().Format(time.RFC3339) + rxTime := data.Timestamp + ingestNow := time.Now().UTC().Format(time.RFC3339) + if rxTime == "" { + rxTime = ingestNow } var txID int64 @@ -688,14 +696,14 @@ func (s *Store) InsertTransmission(data *PacketData) (bool, error) { if err == nil { // Existing transmission txID = existingID - if now < existingFirstSeen { - _, _ = s.stmtUpdateTxFirstSeen.Exec(now, txID) + if rxTime < existingFirstSeen { + _, _ = s.stmtUpdateTxFirstSeen.Exec(rxTime, txID) } } else { // New transmission isNew = true result, err := s.stmtInsertTransmission.Exec( - data.RawHex, hash, now, + data.RawHex, hash, rxTime, data.RouteType, data.PayloadType, data.PayloadVersion, data.DecodedJSON, nilIfEmpty(data.ChannelHash), scopeNameForDB(data), @@ -722,13 +730,13 @@ func (s *Store) InsertTransmission(data *PacketData) (bool, error) { observerIdx = &rowid // Update observer last_seen and last_packet_at on every packet to prevent // low-traffic observers from appearing offline (#463) - _, _ = s.stmtUpdateObserverLastSeen.Exec(now, now, rowid) + _, _ = s.stmtUpdateObserverLastSeen.Exec(ingestNow, rxTime, ingestNow, rxTime, rowid) } } // Insert observation epochTs := time.Now().Unix() - if t, err := time.Parse(time.RFC3339, now); err == nil { + if t, err := time.Parse(time.RFC3339, rxTime); err == nil { epochTs = t.Unix() } @@ -753,13 +761,14 @@ func (s *Store) InsertTransmission(data *PacketData) (bool, error) { // UpsertNode inserts or updates a node. func (s *Store) UpsertNode(pubKey, name, role string, lat, lon *float64, lastSeen string) error { + ingestNow := time.Now().UTC().Format(time.RFC3339) now := lastSeen if now == "" { - now = time.Now().UTC().Format(time.RFC3339) + now = ingestNow } _, err := s.stmtUpsertNode.Exec( pubKey, name, role, lat, lon, now, now, - name, role, lat, lon, now, + name, role, lat, lon, ingestNow, now, ) if err != nil { s.Stats.WriteErrors.Add(1) @@ -822,9 +831,23 @@ type ObserverMeta struct { PacketsRecv *int // cumulative packets received since boot } -// UpsertObserver inserts or updates an observer with optional hardware metadata. +// UpsertObserver inserts or updates an observer using the current wall-clock +// time as last_seen. Use UpsertObserverAt when the message envelope provides +// an observer receive-time (e.g. MQTT status and data packet handlers). func (s *Store) UpsertObserver(id, name, iata string, meta *ObserverMeta) error { - now := time.Now().UTC().Format(time.RFC3339) + return s.UpsertObserverAt(id, name, iata, meta, time.Now().UTC().Format(time.RFC3339)) +} + +// UpsertObserverAt inserts or updates an observer with an explicit lastSeen +// timestamp (typically the observer receive-time from the MQTT envelope). The +// SQL uses MAX so last_seen never moves backwards — a retained or replayed +// message whose rxTime pre-dates the existing last_seen is a no-op for that +// field, preventing offline observers from flashing as Online on reconnect. +func (s *Store) UpsertObserverAt(id, name, iata string, meta *ObserverMeta, lastSeen string) error { + ingestNow := time.Now().UTC().Format(time.RFC3339) + if lastSeen == "" { + lastSeen = ingestNow + } normalizedIATA := strings.TrimSpace(strings.ToUpper(iata)) var model, firmware, clientVersion, radio interface{} @@ -854,8 +877,8 @@ func (s *Store) UpsertObserver(id, name, iata string, meta *ObserverMeta) error } _, err := s.stmtUpsertObserver.Exec( - id, name, normalizedIATA, now, now, model, firmware, clientVersion, radio, batteryMv, uptimeSecs, noiseFloor, - name, normalizedIATA, now, model, firmware, clientVersion, radio, batteryMv, uptimeSecs, noiseFloor, + id, name, normalizedIATA, lastSeen, lastSeen, model, firmware, clientVersion, radio, batteryMv, uptimeSecs, noiseFloor, + name, normalizedIATA, ingestNow, lastSeen, model, firmware, clientVersion, radio, batteryMv, uptimeSecs, noiseFloor, ) if err != nil { s.Stats.WriteErrors.Add(1) @@ -1329,7 +1352,8 @@ type MQTTPacketMessage struct { Score *float64 `json:"score"` Direction *string `json:"direction"` Origin string `json:"origin"` - Region string `json:"region,omitempty"` // optional region override (#788) + Region string `json:"region,omitempty"` // optional region override (#788) + Timestamp string `json:"timestamp,omitempty"` // observer receive time, resolved by handler } // BuildPacketData constructs a PacketData from a decoded packet and MQTT message. @@ -1337,7 +1361,6 @@ type MQTTPacketMessage struct { // to guarantee the stored path always matches the raw bytes. This matters for // TRACE packets where decoded.Path.Hops is overwritten with payload hops (#886). func BuildPacketData(msg *MQTTPacketMessage, decoded *DecodedPacket, observerID, region string, regionKeys map[string][]byte) *PacketData { - now := time.Now().UTC().Format(time.RFC3339) pathJSON := "[]" // For TRACE packets, path_json must be the payload-decoded route hops // (decoded.Path.Hops), NOT the raw_hex header bytes which are SNR values. @@ -1354,7 +1377,7 @@ func BuildPacketData(msg *MQTTPacketMessage, decoded *DecodedPacket, observerID, pd := &PacketData{ RawHex: msg.Raw, - Timestamp: now, + Timestamp: msg.Timestamp, ObserverID: observerID, ObserverName: msg.Origin, SNR: msg.SNR, diff --git a/cmd/ingestor/db_test.go b/cmd/ingestor/db_test.go index 642719da..8d45e31a 100644 --- a/cmd/ingestor/db_test.go +++ b/cmd/ingestor/db_test.go @@ -830,10 +830,11 @@ func TestBuildPacketData(t *testing.T) { snr := 5.0 rssi := -100.0 msg := &MQTTPacketMessage{ - Raw: rawHex, - SNR: &snr, - RSSI: &rssi, - Origin: "test-observer", + Raw: rawHex, + SNR: &snr, + RSSI: &rssi, + Origin: "test-observer", + Timestamp: "2026-05-16T10:00:00Z", } pkt := BuildPacketData(msg, decoded, "obs123", "SJC", nil) @@ -865,8 +866,8 @@ func TestBuildPacketData(t *testing.T) { if pkt.PayloadType != decoded.Header.PayloadType { t.Errorf("payloadType mismatch") } - if pkt.Timestamp == "" { - t.Error("timestamp should be set") + if pkt.Timestamp != "2026-05-16T10:00:00Z" { + t.Errorf("timestamp=%s, want 2026-05-16T10:00:00Z", pkt.Timestamp) } if pkt.DecodedJSON == "" || pkt.DecodedJSON == "{}" { t.Error("decodedJSON should be populated") diff --git a/cmd/ingestor/main.go b/cmd/ingestor/main.go index 54d0e684..fc3be6ad 100644 --- a/cmd/ingestor/main.go +++ b/cmd/ingestor/main.go @@ -447,7 +447,7 @@ func handleMessage(store *Store, tag string, source MQTTSource, m mqtt.Message, name, _ := msg["origin"].(string) iata := parts[1] meta := extractObserverMeta(msg) - if err := store.UpsertObserver(observerID, name, iata, meta); err != nil { + if err := store.UpsertObserverAt(observerID, name, iata, meta, resolveRxTime(msg, tag)); err != nil { log.Printf("MQTT [%s] observer status error: %v", tag, err) } // Insert metrics sample from status message @@ -531,6 +531,7 @@ func handleMessage(store *Store, tag string, source MQTTSource, m mqtt.Message, } mqttMsg := &MQTTPacketMessage{Raw: rawHex} + mqttMsg.Timestamp = resolveRxTime(msg, tag) // Parse optional region from JSON payload (#788) if v, ok := msg["region"].(string); ok && v != "" { mqttMsg.Region = v @@ -668,7 +669,7 @@ func handleMessage(store *Store, tag string, source MQTTSource, m mqtt.Message, if mqttMsg.Region != "" { effectiveRegion = mqttMsg.Region } - if err := store.UpsertObserver(observerID, origin, effectiveRegion, nil); err != nil { + if err := store.UpsertObserverAt(observerID, origin, effectiveRegion, nil, mqttMsg.Timestamp); err != nil { log.Printf("MQTT [%s] observer upsert error: %v", tag, err) } } @@ -712,8 +713,9 @@ func handleMessage(store *Store, tag string, source MQTTSource, m mqtt.Message, decodedJSON, _ := json.Marshal(channelMsg) - now := time.Now().UTC().Format(time.RFC3339) - hashInput := fmt.Sprintf("ch:%s:%s:%s", channelIdx, text, now) + ingestNow := time.Now().UTC().Format(time.RFC3339) + rxTime := resolveRxTime(msg, tag) + hashInput := fmt.Sprintf("ch:%s:%s:%s", channelIdx, text, ingestNow) h := sha256.Sum256([]byte(hashInput)) hash := hex.EncodeToString(h[:])[:16] @@ -753,7 +755,7 @@ func handleMessage(store *Store, tag string, source MQTTSource, m mqtt.Message, } pktData := &PacketData{ - Timestamp: now, + Timestamp: rxTime, ObserverID: "companion", ObserverName: "L1 Pro (BLE)", SNR: snr, @@ -805,8 +807,9 @@ func handleMessage(store *Store, tag string, source MQTTSource, m mqtt.Message, decodedJSON, _ := json.Marshal(dm) - now := time.Now().UTC().Format(time.RFC3339) - hashInput := fmt.Sprintf("dm:%s:%s", text, now) + ingestNow := time.Now().UTC().Format(time.RFC3339) + rxTime := resolveRxTime(msg, tag) + hashInput := fmt.Sprintf("dm:%s:%s", text, ingestNow) h := sha256.Sum256([]byte(hashInput)) hash := hex.EncodeToString(h[:])[:16] @@ -846,7 +849,7 @@ func handleMessage(store *Store, tag string, source MQTTSource, m mqtt.Message, } pktData := &PacketData{ - Timestamp: now, + Timestamp: rxTime, ObserverID: "companion", ObserverName: "L1 Pro (BLE)", SNR: snr, @@ -1030,6 +1033,71 @@ func firstNonEmpty(vals ...string) string { return "" } +// resolveRxTime returns the observer receive-time for a packet, taken from +// the MQTT envelope's "timestamp" field. Falls back to ingest time only when +// the field is missing, unparseable, or implausibly in the future (a +// clock-skewed observer). Result is always RFC3339 UTC. +// +// The envelope timestamp is stamped by the uploader when the radio receives +// the frame, not when the MQTT message is published — so a buffered packet +// uploaded hours late still carries its true receive time. Using ingest time +// (time.Now()) here mis-dated such packets by the upload delay. +func resolveRxTime(msg map[string]interface{}, tag string) string { + now := time.Now().UTC() + raw, _ := msg["timestamp"].(string) + if raw == "" { + return now.Format(time.RFC3339) + } + t, err := parseEnvelopeTime(raw) + if err != nil { + log.Printf("MQTT [%s] unparseable timestamp %q, using ingest time", tag, raw) + return now.Format(time.RFC3339) + } + // Hard reject: > 14h ahead is a genuine clock error (UTC+14 is the maximum + // standard offset, so nothing valid should be further ahead than that). + if t.After(now.Add(14 * time.Hour)) { + log.Printf("MQTT [%s] future timestamp %q, using ingest time", tag, raw) + return now.Format(time.RFC3339) + } + // Hard reject: > 30 days in the past is an RTC-reset node reporting a + // factory date (e.g. 2020-01-01). Such a value would permanently drag + // transmissions.first_seen backwards via stmtUpdateTxFirstSeen in + // InsertTransmission. No legitimate buffered upload is that stale. + if t.Before(now.Add(-30 * 24 * time.Hour)) { + log.Printf("MQTT [%s] stale timestamp %q (>30d old), using ingest time", tag, raw) + return now.Format(time.RFC3339) + } + // Soft clamp: naive local-clock timestamps from UTC+N observers are parsed + // as-if UTC, making them appear N hours in the future. A UTC+2 observer's + // live packet looks 2h ahead, but it is NOT a buffered packet — the whole + // point of using rxTime is to preserve the past timestamp for packets that + // were buffered offline. If rxTime is ahead of now, the packet is live and + // ingest time is the correct value. This also prevents storing future + // timestamps that would show ⚠️ in the UI for every packet from UTC+N nodes. + if t.After(now) { + return now.Format(time.RFC3339) + } + return t.UTC().Format(time.RFC3339) +} + +// parseEnvelopeTime parses the MQTT envelope timestamp. Two on-wire forms +// occur: zone-aware ISO8601 (RFC3339), and a naive local-clock ISO string +// with no zone (python datetime.isoformat()). Zone-aware layouts are tried +// first; naive layouts are assumed UTC, leaving a bounded residual offset +// equal to the observer's UTC offset for naive-timestamp uploaders. +func parseEnvelopeTime(s string) (time.Time, error) { + for _, layout := range []string{ + time.RFC3339, // 2026-05-16T10:00:00Z / +02:00 + "2006-01-02T15:04:05.999999", // python isoformat w/ microseconds + "2006-01-02T15:04:05", // naive ISO + } { + if t, err := time.Parse(layout, s); err == nil { + return t, nil + } + } + return time.Time{}, fmt.Errorf("unrecognized timestamp layout: %q", s) +} + // deriveHashtagChannelKey derives an AES-128 key from a channel name. // Same algorithm as Node.js: SHA-256(channelName) → first 32 hex chars (16 bytes). func deriveHashtagChannelKey(channelName string) string { diff --git a/cmd/ingestor/rxtime_test.go b/cmd/ingestor/rxtime_test.go new file mode 100644 index 00000000..58bf6860 --- /dev/null +++ b/cmd/ingestor/rxtime_test.go @@ -0,0 +1,80 @@ +package main + +import ( + "testing" + "time" +) + +func TestParseEnvelopeTime(t *testing.T) { + cases := []struct { + name string + in string + ok bool + }{ + {"rfc3339 utc", "2026-05-16T10:00:00Z", true}, + {"rfc3339 offset", "2026-05-16T12:00:00+02:00", true}, + {"naive iso", "2026-05-16T10:00:00", true}, + {"naive iso micros", "2026-05-16T10:00:00.123456", true}, + {"garbage", "not-a-time", false}, + {"empty", "", false}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + _, err := parseEnvelopeTime(c.in) + if (err == nil) != c.ok { + t.Fatalf("parseEnvelopeTime(%q): want ok=%v, got err=%v", c.in, c.ok, err) + } + }) + } +} + +func TestResolveRxTime(t *testing.T) { + now := time.Now().UTC() + + mustParse := func(s string) time.Time { + t.Helper() + parsed, err := time.Parse(time.RFC3339, s) + if err != nil { + t.Fatalf("result %q is not RFC3339: %v", s, err) + } + return parsed + } + nearNow := func(s string) bool { + d := mustParse(s).Sub(now) + if d < 0 { + d = -d + } + return d <= time.Minute + } + + rx := now.Add(-5 * time.Hour).Format(time.RFC3339) + if got := resolveRxTime(map[string]interface{}{"timestamp": rx}, "test"); got != rx { + t.Errorf("plausible past timestamp: got %q want %q", got, rx) + } + if got := resolveRxTime(map[string]interface{}{}, "test"); !nearNow(got) { + t.Errorf("missing timestamp: got %q, expected ~now", got) + } + if got := resolveRxTime(map[string]interface{}{"timestamp": "garbage"}, "test"); !nearNow(got) { + t.Errorf("garbage timestamp: got %q, expected ~now", got) + } + future := now.Add(48 * time.Hour).Format(time.RFC3339) + if got := resolveRxTime(map[string]interface{}{"timestamp": future}, "test"); !nearNow(got) { + t.Errorf("future timestamp: got %q, expected ~now (rejected)", got) + } + + // RTC-reset node reporting a factory date — must not drag first_seen back. + factory := "2020-01-01T00:00:00Z" + if got := resolveRxTime(map[string]interface{}{"timestamp": factory}, "test"); !nearNow(got) { + t.Errorf("stale factory timestamp: got %q, expected ~now (rejected)", got) + } + // Just past the 30-day floor → rejected. + stale := now.Add(-31 * 24 * time.Hour).Format(time.RFC3339) + if got := resolveRxTime(map[string]interface{}{"timestamp": stale}, "test"); !nearNow(got) { + t.Errorf("stale timestamp >30d: got %q, expected ~now (rejected)", got) + } + // Just inside the 30-day floor → used verbatim. + recent := now.Add(-29 * 24 * time.Hour).Format(time.RFC3339) + if got := resolveRxTime(map[string]interface{}{"timestamp": recent}, "test"); got != recent { + t.Errorf("recent timestamp <30d: got %q want %q", got, recent) + } +}