diff --git a/.github/workflows/deploy.yml b/.github/workflows/deploy.yml index 38b661f75..8fbdff0bc 100644 --- a/.github/workflows/deploy.yml +++ b/.github/workflows/deploy.yml @@ -401,6 +401,7 @@ jobs: - name: Run Playwright E2E tests (fail-fast) run: | + node test-node-liveness-e2e.js 2>&1 | tee node-liveness-e2e-output.txt BASE_URL=http://localhost:13581 node test-e2e-playwright.js 2>&1 | tee e2e-output.txt # M5+M6 of #1668 — axe-core CI gate. # M5: color-contrast on desktop dark+light. diff --git a/cmd/server/node_activity.go b/cmd/server/node_activity.go new file mode 100644 index 000000000..548fe0fd4 --- /dev/null +++ b/cmd/server/node_activity.go @@ -0,0 +1,371 @@ +package main + +import ( + "strings" + "time" +) + +// NodeHealthStats keeps packet analytics separate from identity-safe activity. +// LastHeard is own advert or unambiguous relay evidence; LastAdvert is strictly +// the node's own ADVERT. Unknown timestamps serialize as null, never last_seen: +// the ingestor may refresh last_seen from a safely resolved relay hop. +type NodeHealthStats struct { + TotalTransmissions int `json:"totalTransmissions"` + TotalObservations int `json:"totalObservations"` + TotalPackets int `json:"totalPackets"` + PacketsToday int `json:"packetsToday"` + AvgSnr *float64 `json:"avgSnr"` + AvgHops *int `json:"avgHops,omitempty"` + LastHeard *string `json:"lastHeard"` + LastAdvert *string `json:"lastAdvert"` +} + +// relayPrefixMapLocked reuses the existing bounded node cache. Missing node +// metadata fails closed for short hashes (e.g. during startup/schema failure). +// Caller holds s.mu. There is no per-node or per-packet SQL query here. +func (s *PacketStore) relayPrefixMapLocked() *prefixMap { + if s.db == nil || s.db.conn == nil { + return s.nodePM + } + _, pm := s.getCachedNodesAndPM() + return pm +} + +// confirmedRelayKey deliberately does not use graph/geo/fallback guesses. +// A full raw identity or a unique wire prefix is evidence; a colliding prefix +// remains unknown even if a heuristic previously indexed it under a full key. +func confirmedRelayKey(token string, pm *prefixMap) string { + token = strings.ToLower(token) + if !isHexLower(token) { + return "" + } + key := "" + if len(token) == 64 { // exact internal identity, not a wire prefix + key = token + } else if len(token) == 2 || len(token) == 4 || len(token) == 6 { + key = uniqueResolve(pm, token) + } + if pm != nil { + if _, listener := pm.nonRelay[key]; listener { + return "" + } + } + return key +} + +func txHasConfirmedRelay(tx *StoreTx, key string, pm *prefixMap) bool { + return txConfirmsRelay(tx, newRelayKeyMatcher(key, pm)) +} + +// relayPrefixHexLens are the supported wire hash sizes (1/2/3 bytes). +var relayPrefixHexLens = [...]int{2, 4, 6} + +// relayKeyMatcher answers confirmedRelayKey(token, pm) == key for one key +// without per-token allocation or prefix-map lookups: uniqueness of the key's +// own wire prefixes is decided once. Build one per key per request. +type relayKeyMatcher struct { + key string + unique [3]bool // indexed like relayPrefixHexLens + listener bool +} + +func newRelayKeyMatcher(key string, pm *prefixMap) relayKeyMatcher { + m := relayKeyMatcher{key: key} + if pm == nil { + return m + } + _, m.listener = pm.nonRelay[key] + for i, l := range relayPrefixHexLens { + if len(key) < l || !isHexLower(key[:l]) { + continue + } + cands := pm.m[key[:l]] + m.unique[i] = len(cands) == 1 && strings.ToLower(cands[0].PublicKey) == key + } + return m +} + +// possible reports whether any token can confirm key; listener-only nodes +// and empty keys never relay. +func (m relayKeyMatcher) possible() bool { + return m.key != "" && !m.listener +} + +func (m relayKeyMatcher) uniquePrefix(hexLen int) bool { + for i, l := range relayPrefixHexLens { + if l == hexLen { + return m.unique[i] + } + } + return false +} + +func (m relayKeyMatcher) matches(token string) bool { + if !m.possible() { + return false + } + switch len(token) { + case 64: // exact internal identity, not a wire prefix + return len(m.key) == 64 && hexFoldEqual(token, m.key) + case 2, 4, 6: + return m.uniquePrefix(len(token)) && hexFoldEqual(token, m.key[:len(token)]) + } + return false +} + +// hexFoldEqual reports whether token equals the lowercase hex string lower, +// ignoring ASCII case, i.e. strings.ToLower(token) == lower && isHexLower(lower). +func hexFoldEqual(token, lower string) bool { + if len(token) != len(lower) { + return false + } + for i := 0; i < len(lower); i++ { + c, want := token[i], lower[i] + if !((want >= '0' && want <= '9') || (want >= 'a' && want <= 'f')) { + return false + } + if c != want && !(want >= 'a' && c == want-'a'+'A') { + return false + } + } + return true +} + +// txConfirmsRelay reports whether any raw observed flood path of tx names +// m.key with identity-safe evidence. +func txConfirmsRelay(tx *StoreTx, m relayKeyMatcher) bool { + // Flood paths record already-observed hops. Direct paths list the + // remaining intended route, not nodes that forwarded this observation. + if !m.possible() || !txHasObservedFloodPath(tx) { + return false + } + found := false + forEachObservedRelayHop(tx, func(token string) bool { + found = m.matches(token) + return !found + }) + return found +} + +// relayTokenResolver memoizes confirmedRelayKey per distinct raw token for one +// request, so repeated hops cost a map lookup instead of ToLower + resolve. +// Not safe for concurrent use. +type relayTokenResolver struct { + pm *prefixMap + cache map[string]string +} + +func newRelayTokenResolver(pm *prefixMap) *relayTokenResolver { + return &relayTokenResolver{pm: pm, cache: make(map[string]string, 1024)} +} + +func (r *relayTokenResolver) key(token string) string { + if key, ok := r.cache[token]; ok { + return key + } + key := confirmedRelayKey(token, r.pm) + r.cache[token] = key + return key +} + +// forEachObservedRelayHop visits the display path hops, then the hops of every +// other raw observation path. The display observation is only the longest +// route, so a shorter observation can be the sole raw evidence for a relay that +// resolved indexing attached to the transmission. visit returns false to stop. +// Repeated identical paths are rescanned instead of deduplicated: the scanner +// itself does not allocate (visit may), a per-transmission set would. Caller +// holds s.mu. +func forEachObservedRelayHop(tx *StoreTx, visit func(token string) bool) { + for _, token := range txGetParsedPath(tx) { + if !visit(token) { + return + } + } + for _, obs := range tx.Observations { + if obs == nil || obs.PathJSON == tx.PathJSON { + continue + } + if !visitPathJSONHops(obs.PathJSON, visit) { + return + } + } +} + +// visitPathJSONHops visits the hops of a path_json array with the same accept +// set as parsePathJSON, without reflection or allocation on the common flat +// string-array shape. Malformed input yields no hops, exactly like the display +// path; escapes and null elements fall back to encoding/json. Returns false if +// visit stopped early. +func visitPathJSONHops(pathJSON string, visit func(string) bool) bool { + if !flatPathJSON(pathJSON) { + if strings.IndexByte(pathJSON, '\\') >= 0 || strings.Contains(pathJSON, "null") { + for _, token := range parsePathJSON(pathJSON) { + if !visit(token) { + return false + } + } + } + return true + } + for i := 0; ; { + open := strings.IndexByte(pathJSON[i:], '"') + if open < 0 { + return true + } + start := i + open + 1 + end := start + strings.IndexByte(pathJSON[start:], '"') + if !visit(pathJSON[start:end]) { + return false + } + i = end + 1 + } +} + +// flatPathJSON reports whether s is a JSON array of plain (escape-free) +// strings, allowing JSON whitespace. Validation precedes any visit so a +// malformed tail cannot leak earlier hops as evidence. +func flatPathJSON(s string) bool { + i, n := 0, len(s) + skip := func() { + for i < n && (s[i] == ' ' || s[i] == '\t' || s[i] == '\n' || s[i] == '\r') { + i++ + } + } + skip() + if i == n || s[i] != '[' { + return false + } + i++ + skip() + if i < n && s[i] == ']' { + i++ + skip() + return i == n + } + for { + if i == n || s[i] != '"' { + return false + } + i++ + for i < n && s[i] != '"' { + if s[i] == '\\' || s[i] < 0x20 { + return false + } + i++ + } + if i == n { + return false + } + i++ + skip() + if i == n { + return false + } + if s[i] == ']' { + i++ + skip() + return i == n + } + if s[i] != ',' { + return false + } + i++ + skip() + } +} + +func txHasObservedFloodPath(tx *StoreTx) bool { + return tx.RouteType != nil && (*tx.RouteType == 0 || *tx.RouteType == routeTypeFlood) +} + +func txIsOwnAdvert(tx *StoreTx, key string) bool { + if tx.PayloadType == nil || *tx.PayloadType != payloadTypeAdvert { + return false + } + decoded := tx.ParsedDecoded() + if valid, ok := decoded["signatureValid"].(bool); ok && !valid { + return false + } + pubkey, _ := decoded["pubKey"].(string) + return pubkey != "" && strings.EqualFold(pubkey, key) +} + +func txIsInvalidAdvert(tx *StoreTx) bool { + if tx.PayloadType == nil || *tx.PayloadType != payloadTypeAdvert { + return false + } + valid, ok := tx.ParsedDecoded()["signatureValid"].(bool) + return ok && !valid +} + +// updateNodeActivity is the single-transmission form of node activity: own +// advert, then relay evidence. Health endpoints use the two parts separately +// so each relay candidate is evaluated once per node. +func updateNodeActivity(tx *StoreTx, key string, pm *prefixMap, heard, advert *string) { + updateOwnAdvertActivity(tx, key, heard, advert) + heardAt, _ := parseRelayTS(*heard) + updateRelayActivity(tx, newRelayKeyMatcher(key, pm), heard, &heardAt) +} + +// updateOwnAdvertActivity advances advert and heard for a valid own ADVERT. +func updateOwnAdvertActivity(tx *StoreTx, key string, heard, advert *string) { + if !txIsOwnAdvert(tx, key) { + return + } + timestamp, valid := parseRelayTS(tx.FirstSeen) + if !valid { + return + } + if oldAdvert, _ := parseRelayTS(*advert); timestamp.After(oldAdvert) { + *advert = tx.FirstSeen + } + if oldHeard, _ := parseRelayTS(*heard); timestamp.After(oldHeard) { + *heard = tx.FirstSeen + } +} + +// updateRelayActivity advances heard (and its parsed heardAt) when tx is newer +// and confirms relay evidence. The timestamp gate runs first: path evidence +// can only matter for a newer transmission. Returns whether the evidence was +// evaluated. Never touches advert. +func updateRelayActivity(tx *StoreTx, m relayKeyMatcher, heard *string, heardAt *time.Time) bool { + timestamp, valid := parseRelayTS(tx.FirstSeen) + if !valid || !timestamp.After(*heardAt) { + return false + } + if txIsInvalidAdvert(tx) || !txConfirmsRelay(tx, m) { + return true + } + *heard, *heardAt = tx.FirstSeen, timestamp + return true +} + +// updateIndexedRelayActivityLocked folds the relay evidence behind +// GetRepeaterRelayInfo into health activity, over the same candidates +// (byPathHop full key, byNode, unique raw prefixes). Only heard changes: a +// relay is never an advert. seen is request-local scratch (cleared here) so a +// transmission's path evidence is evaluated at most once per node; candidates +// not newer than heard are skipped before any path work. Caller holds s.mu. +func (s *PacketStore) updateIndexedRelayActivityLocked(key string, pm *prefixMap, heard *string, seen map[int]struct{}) { + m := newRelayKeyMatcher(key, pm) + if !m.possible() { + return + } + clear(seen) + heardAt, _ := parseRelayTS(*heard) + forEachRelayCandidate(s.byPathHop, s.byNode, m, func(tx *StoreTx, _ bool) { + if _, done := seen[tx.ID]; done { + return + } + if updateRelayActivity(tx, m, heard, &heardAt) { + seen[tx.ID] = struct{}{} + } + }) +} + +func timestampPointer(ts string) *string { + if ts == "" { + return nil + } + return &ts +} diff --git a/cmd/server/node_liveness_regression_test.go b/cmd/server/node_liveness_regression_test.go new file mode 100644 index 000000000..aef934f67 --- /dev/null +++ b/cmd/server/node_liveness_regression_test.go @@ -0,0 +1,574 @@ +package main + +import ( + "encoding/json" + "fmt" + "reflect" + "strings" + "testing" + "time" +) + +func TestNodeHealth_LastAdvertIsOwnAdvertNotTraffic(t *testing.T) { + db := setupCapabilityTestDB(t) + defer db.conn.Close() + if _, err := db.conn.Exec("ALTER TABLE nodes ADD COLUMN foreign_advert INTEGER DEFAULT 0"); err != nil { + t.Fatal(err) + } + key := "aa" + strings.Repeat("11", 31) + other := "bb" + strings.Repeat("22", 31) + if _, err := db.conn.Exec("INSERT INTO nodes (public_key, name, role, last_seen) VALUES (?, 'Synthetic', 'repeater', ?)", key, recentTS(0)); err != nil { + t.Fatal(err) + } + for _, ownAdvert := range []bool{true, false} { + t.Run(map[bool]string{true: "own advert", false: "no advert"}[ownAdvert], func(t *testing.T) { + store := NewPacketStore(db, nil) + mk := func(id, pt int, source string, hours float64) *StoreTx { + return &StoreTx{ID: id, PayloadType: &pt, FirstSeen: time.Now().UTC().Add(-time.Duration(hours * float64(time.Hour))).Format(time.RFC3339), DecodedJSON: `{"pubKey":"` + source + `"}`, PathJSON: `[]`} + } + packets := []*StoreTx{mk(2, 4, other, 2), mk(3, 2, key, 1), mk(4, 3, key, 0.1)} + advert := mk(1, 4, strings.ToUpper(key), 72) + if ownAdvert { + packets = append([]*StoreTx{advert}, packets...) + } + store.byNode[key] = packets + health, err := store.GetNodeHealth(key) + if err != nil { + t.Fatal(err) + } + encoded, err := json.Marshal(health) + if err != nil { + t.Fatal(err) + } + var body struct { + Stats struct { + LastAdvert *string `json:"lastAdvert"` + LastHeard *string `json:"lastHeard"` + TotalPackets int `json:"totalPackets"` + } `json:"stats"` + } + if err := json.Unmarshal(encoded, &body); err != nil { + t.Fatal(err) + } + if !strings.Contains(string(encoded), `"lastAdvert":`) { + t.Fatal("lastAdvert contract missing (must be explicitly null when unknown)") + } + if ownAdvert && (body.Stats.LastAdvert == nil || *body.Stats.LastAdvert != advert.FirstSeen) { + t.Errorf("lastAdvert = %v, want own advert %s", body.Stats.LastAdvert, advert.FirstSeen) + } + if !ownAdvert && body.Stats.LastAdvert != nil { + t.Errorf("traffic or relay-touched last_seen became an advert: %v", body.Stats.LastAdvert) + } + if !reflect.DeepEqual(body.Stats.LastHeard, body.Stats.LastAdvert) || body.Stats.TotalPackets != len(packets) { + t.Fatalf("lastHeard must use certain activity while analytics counts stay intact: %+v", body.Stats) + } + bulk := store.GetBulkHealth(10, "", "") + if len(bulk) != 1 || !reflect.DeepEqual(health["stats"].(NodeHealthStats).LastHeard, bulk[0]["stats"].(NodeHealthStats).LastHeard) || !reflect.DeepEqual(health["stats"].(NodeHealthStats).LastAdvert, bulk[0]["stats"].(NodeHealthStats).LastAdvert) { + t.Fatalf("single/bulk health timestamp mismatch: %v", bulk) + } + }) + } +} + +func TestRepeaterRelayActivity_CollisionAndHeuristicFullKeyFailClosed(t *testing.T) { + db := setupCapabilityTestDB(t) + defer db.conn.Close() + a := "aa11" + strings.Repeat("11", 30) + b := "aa22" + strings.Repeat("22", 30) + for _, key := range []string{a, b} { + if _, err := db.conn.Exec("INSERT INTO nodes (public_key, role) VALUES (?, 'repeater')", key); err != nil { + t.Fatal(err) + } + } + store := NewPacketStore(db, nil) + pt, rt := 3, routeTypeFlood + ambiguous := &StoreTx{ID: 1, PayloadType: &pt, RouteType: &rt, FirstSeen: recentTS(0), PathJSON: `["AA"]`, ScopeName: "ambiguous"} + confirmed := &StoreTx{ID: 2, PayloadType: &pt, RouteType: &rt, FirstSeen: recentTS(0), PathJSON: `["AA22"]`, ScopeName: "confirmed"} + store.byPathHop["aa"] = []*StoreTx{ambiguous} + // A previous heuristic guess must not turn a colliding raw hop into certainty. + store.byPathHop[a] = []*StoreTx{ambiguous} + store.byPathHop[b] = []*StoreTx{confirmed} + store.byPathHop["aa22"] = []*StoreTx{confirmed} + bulk := store.computeRepeaterRelayInfoMap(24) + for _, key := range []string{a, b} { + one := store.GetRepeaterRelayInfo(key, 24) + if !reflect.DeepEqual(one, bulk[key]) { + t.Fatalf("bulk/single mismatch for %s: single %+v, bulk %+v", key, one, bulk[key]) + } + if key == a && (one.RelayActive || one.LastRelayed != "" || one.RelayCount24h != 0 || len(one.TransportedScopes) != 0) { + t.Fatalf("offline colliding node falsely active: %+v", one) + } + if key == b && (!one.RelayActive || one.RelayCount1h != 1 || one.RelayCount24h != 1 || one.LastRelayed != confirmed.FirstSeen || !reflect.DeepEqual(one.TransportedScopes, []string{"confirmed"})) { + t.Fatalf("unique 2-byte relay lost or ambiguous traffic added: %+v", one) + } + } +} + +func TestRepeaterRelayActivity_MissingEvidenceFailsClosed(t *testing.T) { + key := "cc" + strings.Repeat("33", 31) + pt, rt := 2, routeTypeFlood + for _, path := range []string{"", `[]`, `["cc"]`} { + store := &PacketStore{byPathHop: map[string][]*StoreTx{key: {{ID: 1, PayloadType: &pt, RouteType: &rt, FirstSeen: recentTS(0), PathJSON: path}}}} + if got := store.GetRepeaterRelayInfo(key, 24); got.LastRelayed != "" || got.RelayActive { + t.Fatalf("missing prefix map/path must fail closed, path=%q: %+v", path, got) + } + } + // An exact full-key path remains certain without a prefix map. + tx := &StoreTx{ID: 2, PayloadType: &pt, RouteType: &rt, FirstSeen: time.Now().UTC().Add(-10 * time.Minute).Format(time.RFC3339), PathJSON: `["` + key + `"]`} + store := &PacketStore{byPathHop: map[string][]*StoreTx{key: {tx}}} + if got := store.GetRepeaterRelayInfo(key, 24); !got.RelayActive || got.RelayCount1h != 1 { + t.Fatalf("exact key positive control lost: %+v", got) + } +} + +func TestNodeActivity_ObservedFloodNotPlannedDirectRoute(t *testing.T) { + key := "cc" + strings.Repeat("33", 31) + pt := 2 + for _, route := range []int{0, 1, 2, 3} { + t.Run(fmt.Sprintf("route%d", route), func(t *testing.T) { + tx := &StoreTx{ID: 1, PayloadType: &pt, RouteType: &route, FirstSeen: recentTS(0), PathJSON: `["` + key + `"]`} + store := &PacketStore{byPathHop: map[string][]*StoreTx{key: {tx}}} + info := store.GetRepeaterRelayInfo(key, 24) + var heard, advert string + updateNodeActivity(tx, key, nil, &heard, &advert) + flood := route == 0 || route == 1 + if info.RelayActive != flood || (heard != "") != flood || advert != "" { + t.Fatalf("planned direct route became observed activity: route=%d relay=%+v heard=%q advert=%q", route, info, heard, advert) + } + }) + } +} + +func TestNodeActivity_RejectsInvalidAdvertAndMalformedTimestamps(t *testing.T) { + key := "dd" + strings.Repeat("44", 31) + pt, direct := payloadTypeAdvert, 2 + makeAdvert := func(ts string, signature string) *StoreTx { + return &StoreTx{PayloadType: &pt, RouteType: &direct, FirstSeen: ts, DecodedJSON: `{"pubKey":"` + key + `"` + signature + `}`} + } + var heard, advert string + for _, tx := range []*StoreTx{makeAdvert("not-a-timestamp", ""), makeAdvert(recentTS(0), `,"signatureValid":false`)} { + updateNodeActivity(tx, key, nil, &heard, &advert) + } + flood := routeTypeFlood + forgedFlood := makeAdvert(recentTS(0), `,"signatureValid":false`) + forgedFlood.RouteType, forgedFlood.PathJSON = &flood, `["`+key+`"]` + updateNodeActivity(forgedFlood, key, nil, &heard, &advert) + if heard != "" || advert != "" { + t.Fatalf("invalid advert/timestamp became activity: heard=%q advert=%q", heard, advert) + } + // Lexicographically Z sorts after '.', but the fractional timestamp is newer. + older := "2026-01-01T00:00:00Z" + newer := "2026-01-01T00:00:00.500Z" + updateNodeActivity(makeAdvert(older, ""), key, nil, &heard, &advert) + updateNodeActivity(makeAdvert(newer, ""), key, nil, &heard, &advert) + if heard != newer || advert != newer { + t.Fatalf("mixed RFC3339 precision sorted incorrectly: heard=%q advert=%q", heard, advert) + } + for _, token := range []string{"dd444444", "dd4444444444", "gg", strings.Repeat("g", 64)} { + if got := confirmedRelayKey(token, nil); got != "" { + t.Fatalf("unsupported/malformed hop became identity: %q -> %q", token, got) + } + } +} + +// Realistic indexed scale, with colliding short hashes and unique 3-byte hops. +// Bounds are inherited from packet-store eviction; no per-node SQL is needed. +func BenchmarkConfirmedRelayBulk30KPackets2KNodes(b *testing.B) { + const nodeCount, packetCount = 2000, 30000 + nodes := make([]nodeInfo, nodeCount) + for i := range nodes { + nodes[i] = nodeInfo{PublicKey: fmt.Sprintf("%06x", i+1) + strings.Repeat("00", 29), Role: "repeater"} + } + store := &PacketStore{nodePM: buildPrefixMap(nodes), byPathHop: make(map[string][]*StoreTx)} + pt, rt := 3, routeTypeFlood + ts := time.Now().UTC().Add(-10 * time.Minute).Format(time.RFC3339) + for i := 0; i < packetCount; i++ { + keys := []string{nodes[i%nodeCount].PublicKey, nodes[(i+1)%nodeCount].PublicKey, nodes[(i+2)%nodeCount].PublicKey} + tx := &StoreTx{ID: i + 1, PayloadType: &pt, RouteType: &rt, FirstSeen: ts, PathJSON: fmt.Sprintf(`["%s","%s","%s"]`, keys[0][:6], keys[1][:6], keys[2][:6])} + for _, key := range keys { + store.byPathHop[key] = append(store.byPathHop[key], tx) + store.byPathHop[key[:6]] = append(store.byPathHop[key[:6]], tx) + } + } + b.ReportAllocs() + b.ResetTimer() + for i := 0; i < b.N; i++ { + got := store.computeRepeaterRelayInfoMap(24) + if got[nodes[0].PublicKey].RelayCount24h != 45 { + b.Fatal("unexpected dedup/attribution result") + } + } +} + +type activitySnapshot struct { + relay RepeaterRelayInfo + bulkRelay RepeaterRelayInfo + health NodeHealthStats + bulkHealth NodeHealthStats + bulkPresent bool +} + +func snapshotNodeActivity(t *testing.T, store *PacketStore, key string) activitySnapshot { + t.Helper() + snap := activitySnapshot{relay: store.GetRepeaterRelayInfo(key, 24)} + bulkRelay, indexed := store.computeRepeaterRelayInfoMap(24)[key] + if !indexed { + // Keys never indexed are absent from the bulk map; enrichment skips them. + bulkRelay = RepeaterRelayInfo{WindowHours: 24} + } + snap.bulkRelay = bulkRelay + health, err := store.GetNodeHealth(key) + if err != nil || health == nil { + t.Fatalf("health for %s: %v %v", key, health, err) + } + snap.health = health["stats"].(NodeHealthStats) + for _, row := range store.GetBulkHealth(100, "", "") { + if row["public_key"] == key { + snap.bulkHealth, snap.bulkPresent = row["stats"].(NodeHealthStats), true + } + } + if !snap.bulkPresent { + t.Fatalf("bulk health omitted %s", key) + } + if !reflect.DeepEqual(snap.relay, snap.bulkRelay) { + t.Fatalf("single/bulk relay mismatch for %s: single %+v, bulk %+v", key, snap.relay, snap.bulkRelay) + } + if !reflect.DeepEqual(snap.health.LastHeard, snap.bulkHealth.LastHeard) || !reflect.DeepEqual(snap.health.LastAdvert, snap.bulkHealth.LastAdvert) { + t.Fatalf("single/bulk health mismatch for %s: single %+v, bulk %+v", key, snap.health, snap.bulkHealth) + } + return snap +} + +func derefTS(ts *string) string { + if ts == nil { + return "null" + } + return *ts +} + +func activityTestStore(t *testing.T, repeaters ...string) (*DB, *PacketStore) { + t.Helper() + db := setupCapabilityTestDB(t) + if _, err := db.conn.Exec("ALTER TABLE nodes ADD COLUMN foreign_advert INTEGER DEFAULT 0"); err != nil { + t.Fatal(err) + } + for _, key := range repeaters { + if _, err := db.conn.Exec("INSERT INTO nodes (public_key, name, role, last_seen) VALUES (?, ?, 'repeater', ?)", key, "Synthetic "+key[:6], recentTS(0)); err != nil { + t.Fatal(err) + } + } + return db, NewPacketStore(db, nil) +} + +// storeObservedTx builds a transmission the way ingest does: one StoreObs per +// observer path, with the longest path selected as the display observation. +func storeObservedTx(store *PacketStore, id, pt int, ts, decoded string, paths ...string) *StoreTx { + rt := routeTypeFlood + tx := &StoreTx{ID: id, Hash: fmt.Sprintf("synthetic-%d", id), PayloadType: &pt, RouteType: &rt, FirstSeen: ts, LatestSeen: ts, DecodedJSON: decoded} + for i, path := range paths { + tx.Observations = append(tx.Observations, &StoreObs{ID: id*10 + i, TransmissionID: id, ObserverID: fmt.Sprintf("observer-%d", i), PathJSON: path, Timestamp: ts}) + } + tx.ObservationCount = len(tx.Observations) + pickBestObservation(tx) + store.packets = append(store.packets, tx) + store.byHash[tx.Hash] = tx + store.byTxID[tx.ID] = tx + addTxToPathHopIndex(store.byPathHop, tx) + store.indexByNode(tx) + return tx +} + +// Review fix A: the display observation is the longest path, but a shorter +// observation of the same flood can be the only raw evidence for a relay. +func TestNodeActivity_AlternateObservationPathConfirmsRelay(t *testing.T) { + relay := "a1b2c3" + strings.Repeat("11", 29) + hopA := "d4e5f6" + strings.Repeat("22", 29) + hopB := "e7f8a9" + strings.Repeat("33", 29) + guessed := "b0cafe" + strings.Repeat("44", 29) + twin := "b0beef" + strings.Repeat("55", 29) + db, store := activityTestStore(t, relay, hopA, hopB, guessed, twin) + defer db.conn.Close() + ts := time.Now().UTC().Add(-10 * time.Minute).Format(time.RFC3339) + hopsSeen := map[string]bool{} + + // Two observers heard the short route through relay (one lower-case), + // the display observation took a longer route that does not name relay. + confirmed := storeObservedTx(store, 1, 2, ts, `{"type":"TXT_MSG"}`, `["A1B2C3"]`, `["D4E5F6","E7F8A9"]`, `["a1b2c3"]`) + if confirmed.PathJSON != `["D4E5F6","E7F8A9"]` { + t.Fatalf("fixture must display the longer observation, got %s", confirmed.PathJSON) + } + store.indexResolvedPathHops(confirmed, mergeResolvedPubkeys([]*string{&relay}, []*string{&hopA, &hopB}), hopsSeen) + // A heuristic resolution of a colliding alternate hop is not evidence. + heuristic := storeObservedTx(store, 2, 2, ts, `{"type":"TXT_MSG"}`, `["B0"]`, `["D4E5F6","E7F8A9"]`) + store.indexResolvedPathHops(heuristic, mergeResolvedPubkeys([]*string{&guessed}, []*string{&hopA, &hopB}), hopsSeen) + // A planned DIRECT route in any observation is not observed relay activity. + planned := storeObservedTx(store, 3, 2, ts, `{"type":"TXT_MSG"}`, `["A1B2C3"]`, `["D4E5F6","E7F8A9"]`) + direct := 2 + planned.RouteType = &direct + store.indexResolvedPathHops(planned, mergeResolvedPubkeys([]*string{&relay}, []*string{&hopA, &hopB}), hopsSeen) + + got := snapshotNodeActivity(t, store, relay) + if !got.relay.RelayActive || got.relay.RelayCount1h != 1 || got.relay.RelayCount24h != 1 || got.relay.LastRelayed != ts { + t.Errorf("alternate observation relay lost or not deduplicated per transmission: %+v", got.relay) + } + if got.health.LastHeard == nil || *got.health.LastHeard != ts || got.health.LastAdvert != nil { + t.Errorf("health must reflect the alternate-observation relay without inventing an advert: lastHeard=%v lastAdvert=%v", derefTS(got.health.LastHeard), derefTS(got.health.LastAdvert)) + } + for _, key := range []string{guessed, twin} { + got := snapshotNodeActivity(t, store, key) + if got.relay.RelayActive || got.relay.LastRelayed != "" || got.health.LastHeard != nil || got.health.LastAdvert != nil { + t.Fatalf("colliding alternate hop became activity for %s: relay %+v health %+v", key[:6], got.relay, got.health) + } + } + if got := snapshotNodeActivity(t, store, hopA); got.relay.RelayCount24h != 2 || got.health.LastHeard == nil { + t.Fatalf("display-path positive control lost: relay %+v health %+v", got.relay, got.health) + } +} + +// Review fix B: relay status accepts a unique raw prefix bucket, so health +// must see the same evidence even without resolved full-key byNode membership. +func TestNodeHealth_RawPrefixOnlyRelayMatchesRelayActivity(t *testing.T) { + relay := "c3d4e5" + strings.Repeat("66", 29) + other := "f1a2b3" + strings.Repeat("77", 29) + for _, width := range []int{2, 4, 6} { + for _, ownAdvert := range []bool{false, true} { + t.Run(fmt.Sprintf("%dbyte/advert=%v", width/2, ownAdvert), func(t *testing.T) { + db, store := activityTestStore(t, relay, other) + defer db.conn.Close() + relayTS := time.Now().UTC().Add(-10 * time.Minute).Format(time.RFC3339) + path := `["` + strings.ToUpper(relay[:width]) + `","F1A2B3"]` + relayed := storeObservedTx(store, 1, 2, relayTS, `{"type":"TXT_MSG"}`, path) + var advert *StoreTx + if ownAdvert { + advertTS := time.Now().UTC().Add(-72 * time.Hour).Format(time.RFC3339) + advert = storeObservedTx(store, 2, payloadTypeAdvert, advertTS, `{"type":"ADVERT","pubKey":"`+relay+`","signatureValid":true}`, `[]`) + } + for _, tx := range store.byNode[relay] { + if tx == relayed { + t.Fatal("fixture must not have resolved full-key byNode membership") + } + } + got := snapshotNodeActivity(t, store, relay) + if !got.relay.RelayActive || got.relay.LastRelayed != relayTS { + t.Fatalf("precondition: unique raw prefix must be relay evidence: %+v", got.relay) + } + if got.health.LastHeard == nil || *got.health.LastHeard != relayTS { + t.Fatalf("RelayActive=true (lastRelayed %s) but health lastHeard=%s", relayTS, derefTS(got.health.LastHeard)) + } + if ownAdvert && (got.health.LastAdvert == nil || *got.health.LastAdvert != advert.FirstSeen) { + t.Fatalf("own advert timestamp lost: %v", got.health.LastAdvert) + } + if !ownAdvert && got.health.LastAdvert != nil { + t.Fatalf("relay became an advert: %v", *got.health.LastAdvert) + } + }) + } + } +} + +// Alternate observation paths use an allocation-free scanner; it must accept +// exactly what parsePathJSON accepts for the display path. +func TestVisitPathJSONHops_MatchesParsePathJSON(t *testing.T) { + for _, path := range []string{ + ``, `[]`, ` [ ] `, `null`, `["AA"]`, `["A1B2C3","d4"]`, " [ \"AA\" ,\n\t\"BB\" ] ", + `["AA",]`, `["AA" "BB"]`, `["AA",1]`, `[1,"AA"]`, `["AA"`, `"AA"`, `["AA"]]`, `["AA"]x`, + `[["AA"]]`, `{"a":"AA"}`, `["AA"]`, `["AA\"BB"]`, `["AA",null]`, "[\"A\nA\"]", + } { + var got []string + visitPathJSONHops(path, func(token string) bool { + got = append(got, token) + return true + }) + if want := parsePathJSON(path); !reflect.DeepEqual(got, want) && !(len(got) == 0 && len(want) == 0) { + t.Errorf("path %q: scanner %q, parsePathJSON %q", path, got, want) + } + } +} + +// Review fix F1: after a restart buildPathHopIndex re-indexes display paths +// only, so a relay named solely by a non-display observation is no longer in +// its full-key path-hop bucket. Relay status and health must still agree with +// the same transmission ingested live. +func TestNodeActivity_AlternateObservationRelaySurvivesRestart(t *testing.T) { + relay := "a1b2c3" + strings.Repeat("11", 29) + hopA := "d4e5f6" + strings.Repeat("22", 29) + hopB := "e7f8a9" + strings.Repeat("33", 29) + type activity struct { + relayActive bool + lastRelayed string + count1h, count24h int + heard, bulkHeard string + advert, bulkAdvert string + bulkRelayActive bool + bulkLastRelayed string + bulkCount24h int + totalPackets, obsCnt int + } + snapshot := func(t *testing.T, store *PacketStore) activity { + t.Helper() + got := snapshotNodeActivity(t, store, relay) + return activity{ + relayActive: got.relay.RelayActive, lastRelayed: got.relay.LastRelayed, + count1h: got.relay.RelayCount1h, count24h: got.relay.RelayCount24h, + heard: derefTS(got.health.LastHeard), bulkHeard: derefTS(got.bulkHealth.LastHeard), + advert: derefTS(got.health.LastAdvert), bulkAdvert: derefTS(got.bulkHealth.LastAdvert), + bulkRelayActive: got.bulkRelay.RelayActive, bulkLastRelayed: got.bulkRelay.LastRelayed, bulkCount24h: got.bulkRelay.RelayCount24h, + totalPackets: got.health.TotalPackets, obsCnt: got.health.TotalObservations, + } + } + for _, persisted := range []bool{false, true} { + t.Run(fmt.Sprintf("persisted_resolved_path=%v", persisted), func(t *testing.T) { + now := time.Now().UTC() + firstSeen := now.Add(-10 * time.Minute).Format(time.RFC3339) + seed := func(t *testing.T, db *DB) { + for _, key := range []string{relay, hopA, hopB} { + mustExec(t, db, `INSERT INTO nodes (public_key, name, role, last_seen, first_seen, advert_count) VALUES (?, ?, 'repeater', ?, '2026-01-01', 1)`, key, "Synthetic "+key[:6], firstSeen) + } + } + insertTx := func(t *testing.T, db *DB) { + mustExec(t, db, `INSERT INTO transmissions (id, raw_hex, hash, first_seen, route_type, payload_type, decoded_json) VALUES (7, 'CAFE', 'restart-alt-obs', ?, 1, 2, '{"type":"TXT_MSG"}')`, firstSeen) + var short, long interface{} + if persisted { + short, long = `["`+relay+`"]`, `["`+hopA+`","`+hopB+`"]` + } + mustExec(t, db, `INSERT INTO observations (transmission_id, observer_idx, path_json, timestamp, resolved_path) VALUES (7, NULL, '["D4","E7"]', ?, ?)`, now.Add(-10*time.Minute).Unix(), long) + mustExec(t, db, `INSERT INTO observations (transmission_id, observer_idx, path_json, timestamp, resolved_path) VALUES (7, NULL, '["A1"]', ?, ?)`, now.Add(-9*time.Minute).Unix(), short) + } + + // Restart: rows already persisted, loaded via Load + deferred index build. + loadDB := setupTestDB(t) + defer loadDB.conn.Close() + seed(t, loadDB) + insertTx(t, loadDB) + loaded := NewPacketStore(loadDB, nil) + if err := loaded.Load(); err != nil { + t.Fatal(err) + } + if !loaded.WaitIndexesReady(5 * time.Second) { + t.Fatal("indexes not ready") + } + tx := loaded.byTxID[7] + if tx == nil || len(tx.Observations) != 2 || tx.PathJSON != `["D4","E7"]` { + t.Fatalf("fixture must load two observations with the longer display path: %+v", tx) + } + for _, bucket := range []string{relay, "a1"} { + for _, indexed := range loaded.byPathHop[bucket] { + if indexed == tx { + t.Fatalf("fixture must reproduce the restart index shape; tx found in byPathHop[%s]", bucket) + } + } + } + + // Live: same rows ingested after startup. + liveDB := setupTestDB(t) + defer liveDB.conn.Close() + seed(t, liveDB) + live := NewPacketStore(liveDB, nil) + if err := live.Load(); err != nil { + t.Fatal(err) + } + if !live.WaitIndexesReady(5 * time.Second) { + t.Fatal("live indexes not ready") + } + insertTx(t, liveDB) + live.IngestNewFromDB(0, 100) + if live.byTxID[7] == nil || len(live.byTxID[7].Observations) != 2 { + t.Fatalf("live ingest fixture incomplete: %+v", live.byTxID[7]) + } + + want := activity{relayActive: true, lastRelayed: firstSeen, count1h: 1, count24h: 1, heard: firstSeen, bulkHeard: firstSeen, advert: "null", bulkAdvert: "null", + bulkRelayActive: true, bulkLastRelayed: firstSeen, bulkCount24h: 1, totalPackets: 1, obsCnt: 2} + gotLive := snapshot(t, live) + if !reflect.DeepEqual(gotLive, want) { + t.Errorf("live ingest activity:\n got %+v\nwant %+v", gotLive, want) + } + gotLoaded := snapshot(t, loaded) + if !reflect.DeepEqual(gotLoaded, gotLive) { + t.Errorf("restart activity differs from live ingest:\nrestart %+v\n live %+v", gotLoaded, gotLive) + } + }) + } +} + +// byNode membership is only a candidate, never identity proof, and must not +// widen transported-scope provenance. +func TestRelayActivity_ByNodeCandidatesKeepEvidenceRules(t *testing.T) { + relay := "a1b2c3" + strings.Repeat("11", 29) + listener := "c5d6e7" + strings.Repeat("22", 29) + guessed := "b0cafe" + strings.Repeat("44", 29) + twin := "b0beef" + strings.Repeat("55", 29) + hop := "d4e5f6" + strings.Repeat("66", 29) + var nodes []nodeInfo + for _, key := range []string{relay, listener, guessed, twin, hop} { + nodes = append(nodes, nodeInfo{PublicKey: key, Role: "repeater"}) + } + pm := buildPrefixMap(nodes) + pm.markNonRelay([]string{listener}) + ts := time.Now().UTC().Add(-10 * time.Minute).Format(time.RFC3339) + mk := func(id, route int, scope string, paths ...string) *StoreTx { + pt := 2 + tx := &StoreTx{ID: id, PayloadType: &pt, RouteType: &route, FirstSeen: ts, ScopeName: scope} + for _, path := range paths { + tx.Observations = append(tx.Observations, &StoreObs{TransmissionID: id, PathJSON: path, Timestamp: ts}) + } + pickBestObservation(tx) + return tx + } + alternate := mk(1, routeTypeFlood, "scoped", `["A1B2C3"]`, `["D4E5F6","D4E5F6"]`) + ambiguous := mk(2, routeTypeFlood, "", `["B0"]`, `["D4E5F6","D4E5F6"]`) + planned := mk(3, 2, "", `["A1B2C3"]`, `["D4E5F6","D4E5F6"]`) + quiet := mk(4, routeTypeFlood, "", `["C5D6E7"]`, `["D4E5F6","D4E5F6"]`) + store := &PacketStore{nodePM: pm, byPathHop: map[string][]*StoreTx{}, byNode: map[string][]*StoreTx{ + // Duplicated candidates must still count once. + relay: {alternate, planned, alternate}, + guessed: {ambiguous}, + listener: {quiet}, + }} + for _, tx := range []*StoreTx{alternate, ambiguous, planned, quiet} { + addTxToPathHopIndex(store.byPathHop, tx) + } + store.byPathHop["a1b2c3"] = append(store.byPathHop["a1b2c3"], alternate) + bulk := store.computeRepeaterRelayInfoMap(24) + for _, key := range []string{relay, listener, guessed, twin} { + single := store.GetRepeaterRelayInfo(key, 24) + if b, ok := bulk[key]; ok && !reflect.DeepEqual(single, b) { + t.Fatalf("single/bulk mismatch for %s: %+v vs %+v", key[:6], single, b) + } else if !ok && (single.RelayActive || single.LastRelayed != "") { + t.Fatalf("bulk omitted active key %s: %+v", key[:6], single) + } + var heard string + store.updateIndexedRelayActivityLocked(key, pm, &heard, map[int]struct{}{}) + if key == relay { + if !single.RelayActive || single.RelayCount24h != 1 || len(single.TransportedScopes) != 0 || heard != ts { + t.Fatalf("byNode alternate-observation relay: relay %+v heard %q (want once, no scope from candidate)", single, heard) + } + continue + } + if single.RelayActive || single.LastRelayed != "" || heard != "" { + t.Fatalf("%s: candidate membership became evidence: relay %+v heard %q", key[:6], single, heard) + } + } +} + +// relayKeyMatcher replaces confirmedRelayKey(token) == key on hot paths. +func TestRelayKeyMatcher_MatchesConfirmedRelayKey(t *testing.T) { + unique := "a1b2c3" + strings.Repeat("11", 29) + collideA := "b0cafe" + strings.Repeat("44", 29) + collideB := "b0cbee" + strings.Repeat("55", 29) + listener := "c5d6e7" + strings.Repeat("22", 29) + short := "e9f0" + pm := buildPrefixMap([]nodeInfo{{PublicKey: strings.ToUpper(unique), Role: "repeater"}, {PublicKey: collideA, Role: "repeater"}, {PublicKey: collideB, Role: "repeater"}, {PublicKey: listener, Role: "repeater"}, {PublicKey: short, Role: "room"}, {PublicKey: "f1" + strings.Repeat("00", 31), Role: "companion"}}) + pm.markNonRelay([]string{listener}) + keys := []string{unique, collideA, collideB, listener, short, "f1" + strings.Repeat("00", 31), "", "a1", "zz" + strings.Repeat("11", 31)} + tokens := []string{"", "a", "A1", "a1", "A1B2", "a1B2", "A1B2C3", "a1b2c4", "A1B2C3D4", "B0", "B0CA", "b0cb", "C5", "c5d6e7", "E9", "e9f0", "F1", "g1", "A1B2C", strings.ToUpper(unique), unique, collideA, listener, "zz" + strings.Repeat("11", 31), "K1"} + for _, withPM := range []*prefixMap{pm, nil} { + for _, key := range keys { + m := newRelayKeyMatcher(key, withPM) + for _, token := range tokens { + if got, want := m.matches(token), key != "" && confirmedRelayKey(token, withPM) == key; got != want { + t.Errorf("pm=%v key=%q token=%q: matcher %v, confirmedRelayKey %v", withPM != nil, key, token, got, want) + } + } + } + } +} diff --git a/cmd/server/repeater_enrich_bulk.go b/cmd/server/repeater_enrich_bulk.go index 37dc65bb9..c3794c6f2 100644 --- a/cmd/server/repeater_enrich_bulk.go +++ b/cmd/server/repeater_enrich_bulk.go @@ -52,176 +52,145 @@ func (s *PacketStore) GetRepeaterRelayInfoMap(windowHours float64) map[string]Re return cached } -// computeRepeaterRelayInfoMap walks byPathHop once under a single RLock, -// pre-parses every FirstSeen timestamp once (not once-per-pubkey-bucket), -// and emits one RepeaterRelayInfo per hop key. +// computeRepeaterRelayInfoMap walks the relay candidate indexes under a single +// RLock, pre-parses every FirstSeen timestamp and relay identity once (not +// once-per-pubkey-bucket), and emits one RepeaterRelayInfo per hop key and per +// confirmed relay key. Only aggregation runs after the lock is released. // -// Time-complexity invariant: O(unique-tx-in-byPathHop + total-key-bucket -// entries). Memory: one map entry per byPathHop key. Both are bounded by -// the same eviction policy that bounds byPathHop itself. +// Time-complexity invariant: O(unique candidate tx + observation-path hops + +// total candidate-bucket entries): each tx's identity-safe keys and timestamp +// are computed once (hop tokens memoized per request), then every key walks +// only its own candidate buckets (see forEachRelayCandidate) with O(1) +// generation-stamped dedupe. A unique raw prefix belongs to only one full key, +// avoiding collided-prefix fanout. Memory: O(unique tx + keys + distinct hop +// tokens), bounded by store eviction. func (s *PacketStore) computeRepeaterRelayInfoMap(windowHours float64) map[string]RepeaterRelayInfo { + // Everything that reads store data happens under the read lock: ingest + // appends to and eviction compacts byPathHop/byNode slices in place, so a + // copied slice header is not a stable snapshot, and StoreTx observations + // are mutable. The unlocked phase reads only values owned by this call. s.mu.RLock() + pm := s.relayPrefixMapLocked() - // Snapshot the slices (header copy) so we can release the lock before - // the expensive parse pass. Slice headers point at the live underlying - // arrays but those are append-only-by-id; the worst-case race here is - // that ingest grows a slice we already snapshotted (we miss the new - // tail), which is acceptable for a 15s-TTL status read. - snap := make(map[string][]*StoreTx, len(s.byPathHop)) - for k, list := range s.byPathHop { - snap[k] = list + // Per transmission: its relay entry values and the identity-safe relay + // keys confirmed by its raw observed flood paths, computed once. + type bulkRelayTx struct { + entry relayEntry + keys []string + gen int } - - // Build a tx-id-keyed pre-parsed cache so the inner loop doesn't - // re-parse the same FirstSeen N times when the same tx is indexed - // under multiple hop keys (very common — every hop on a path indexes - // the tx). - type parsedTx struct { - t time.Time - ok bool - pt int - } - parseCache := make(map[int]parsedTx, 1<<14) - for _, list := range snap { + resolve := newRelayTokenResolver(pm) + index := make(map[int]int, 1<<14) + txs := make([]bulkRelayTx, 0, 1<<14) + confirmed := make(map[string]struct{}) + add := func(list []*StoreTx) { for _, tx := range list { if tx == nil { continue } - if _, ok := parseCache[tx.ID]; ok { + if _, ok := index[tx.ID]; ok { continue } - pt := -1 - if tx.PayloadType != nil { - pt = *tx.PayloadType + b := bulkRelayTx{entry: newRelayEntry(tx, false)} + b.entry.parsed = true + b.entry.t, b.entry.valid = parseRelayTS(tx.FirstSeen) + if txHasObservedFloodPath(tx) { + forEachObservedRelayHop(tx, func(token string) bool { + key := resolve.key(token) + if key == "" { + return true + } + for _, have := range b.keys { + if have == key { + return true + } + } + b.keys = append(b.keys, key) + confirmed[key] = struct{}{} + return true + }) } - t, ok := parseRelayTS(tx.FirstSeen) - parseCache[tx.ID] = parsedTx{t: t, ok: ok, pt: pt} + index[tx.ID] = len(txs) + txs = append(txs, b) } } - s.mu.RUnlock() - - now := time.Now().UTC() - cutoff1h := now.Add(-1 * time.Hour) - cutoff24h := now.Add(-24 * time.Hour) - var windowCutoff time.Time - if windowHours > 0 { - windowCutoff = now.Add(-time.Duration(windowHours * float64(time.Hour))) + for _, list := range s.byPathHop { + add(list) + } + // byNode candidates only matter for keys a raw hop can confirm: known + // relay prefixes/keys or exact 64-hex identities. + for k, list := range s.byNode { + if len(k) == 64 || (pm != nil && len(pm.m[k]) > 0) { + add(list) + } } - out := make(map[string]RepeaterRelayInfo, len(snap)) - for key, list := range snap { - info := RepeaterRelayInfo{WindowHours: windowHours} - // #1751: accumulate the set of region scope names carried by this - // hop key across every non-advert path-hop tx (NOT time-windowed). - // Captured by the visit closure below — lazily allocated on the first - // scope hit so hosts without scope_name pay nothing per key; converted - // to a sorted, capped slice before this key's info is stored. - var scopeSet map[string]struct{} - // Map scope-filter parity follow-up: latest parseable timestamp - // per scope, mirroring computeRelayInfoFromEntries's scopeLatest - // — kept in exact parity per this file's existing convention. - var scopeLatest map[string]time.Time - // When key looks like a full pubkey (>= 2 hex chars), also fold - // in the matching 1-byte raw-prefix bucket to mirror - // GetRepeaterRelayInfo's behavior. We dedup by tx ID. - var seen map[int]bool - if len(key) >= 2 { - prefix := key[:2] - if prefix != key { - if extra := snap[prefix]; len(extra) > 0 { - seen = make(map[int]bool, len(list)+len(extra)) - } - } + keys := make(map[string]struct{}, len(s.byPathHop)+len(confirmed)) + for key := range s.byPathHop { + keys[key] = struct{}{} + // Unique raw-only prefixes need a full-key result too. + if resolved := resolve.key(key); resolved != "" { + keys[resolved] = struct{}{} } - // includeScope gates TransportedScopes accumulation: the 1-byte - // prefix-bucket fallback below folds in transmissions whose hop - // hash was NEVER resolved to this specific pubkey — MeshCore - // firmware (examples/simple_repeater/MyMesh.cpp allowPacketForward) - // only relays a TRANSPORT_FLOOD/DIRECT packet when the repeater's - // own locally configured region matches, so crediting a scope to a - // node based on nothing but a shared 1-byte hash prefix produces - // claims the protocol itself would never allow (e.g. a repeater - // hundreds of km away "transporting" a hyper-local town scope). - // RelayCount/LastRelayed/RelayActive keep the fallback — those are - // intentionally approximate "is this node active" signals, not a - // specific factual claim about which region it carried. - visit := func(txs []*StoreTx, includeScope bool) { - for _, tx := range txs { - if tx == nil { - continue - } - if seen != nil { - if seen[tx.ID] { - continue - } - seen[tx.ID] = true - } - p, ok := parseCache[tx.ID] + } + // Keys with evidence only via byNode (e.g. non-display observations after + // a restart) need a result too; otherwise list and detail disagree. + for key := range confirmed { + keys[key] = struct{}{} + } + + // Per key: indexes into txs of its deduplicated, confirmed candidates. + // A candidate missing from index cannot confirm key (every candidate list + // was indexed above under the same lock); it is skipped, never mapped to + // another transmission. + type relayRef struct { + tx int + fromPrefix bool + } + type keyRefs struct { + key string + start, end int + } + var refs []relayRef + perKey := make([]keyRefs, 0, len(keys)) + gen := 0 + for key := range keys { + gen++ + start := len(refs) + m := newRelayKeyMatcher(key, pm) + if m.possible() { + forEachRelayCandidate(s.byPathHop, s.byNode, m, func(tx *StoreTx, fromPrefix bool) { + i, ok := index[tx.ID] if !ok { - continue + return } - if p.pt == payloadTypeAdvert { - continue + b := &txs[i] + if b.gen == gen { + return } - // #1751 (tightened, see includeScope doc above): scope - // accumulation is intentionally NOT gated on p.ok - // (timestamp parseability) — a packet with an unparseable - // first_seen still proves the repeater transported that - // scope, as long as the hop resolved unambiguously to it. - if includeScope && tx.ScopeName != "" { - if scopeSet == nil { - scopeSet = map[string]struct{}{} + b.gen = gen + for _, have := range b.keys { + if have == key { + refs = append(refs, relayRef{tx: i, fromPrefix: fromPrefix}) + return } - scopeSet[tx.ScopeName] = struct{}{} } - if !p.ok { - continue - } - if includeScope && tx.ScopeName != "" && windowHours > 0 { - if scopeLatest == nil { - scopeLatest = map[string]time.Time{} - } - if p.t.After(scopeLatest[tx.ScopeName]) { - scopeLatest[tx.ScopeName] = p.t - } - } - if p.t.After(cutoff24h) { - info.RelayCount24h++ - if tx.RouteType != nil && *tx.RouteType == routeTypeFlood { - info.UnscopedRelayCount24h++ - } - if p.t.After(cutoff1h) { - info.RelayCount1h++ - } - } - if info.LastRelayed == "" || tx.FirstSeen > info.LastRelayed { - info.LastRelayed = tx.FirstSeen - if windowHours > 0 && p.t.After(windowCutoff) { - info.RelayActive = true - } else if windowHours > 0 { - info.RelayActive = false - } - } - } - } - visit(list, true) - if seen != nil { - prefix := key[:2] - if prefix != key { - visit(snap[prefix], false) - } + }) } - info.TransportedScopes = sortedCappedScopes(scopeSet) - if windowHours > 0 && scopeLatest != nil { - recentSet := make(map[string]struct{}, len(scopeLatest)) - for scope, t := range scopeLatest { - if t.After(windowCutoff) { - recentSet[scope] = struct{}{} - } - } - info.TransportedScopesRecent = sortedCappedScopes(recentSet) + perKey = append(perKey, keyRefs{key: key, start: start, end: len(refs)}) + } + s.mu.RUnlock() + + out := make(map[string]RepeaterRelayInfo, len(perKey)) + var entries []relayEntry + for _, k := range perKey { + entries = entries[:0] + for _, ref := range refs[k.start:k.end] { + e := txs[ref.tx].entry + e.fromPrefix = ref.fromPrefix + entries = append(entries, e) } - out[key] = info + out[k.key] = computeRelayInfoFromEntries(entries, windowHours) } return out } diff --git a/cmd/server/repeater_liveness.go b/cmd/server/repeater_liveness.go index d1a7c0726..ee8f187fb 100644 --- a/cmd/server/repeater_liveness.go +++ b/cmd/server/repeater_liveness.go @@ -123,68 +123,80 @@ type relayEntry struct { // scope is the tx's region scope name (transmissions.scope_name). // Empty when absent / on older schemas. Used for TransportedScopes (#1751). scope string - // fromPrefix marks entries that came from the 1-byte raw-prefix - // fallback bucket rather than this exact (resolved-pubkey) key — see - // collectRelayEntriesLocked. computeRelayInfoFromEntries must not - // accumulate `scope` for these: MeshCore firmware only relays a - // TRANSPORT_FLOOD/DIRECT packet when the repeater's own configured - // region matches (allowPacketForward in examples/simple_repeater/ - // MyMesh.cpp), so crediting a scope based on nothing but a shared - // 1-byte hash prefix asserts something the protocol wouldn't allow. + // fromPrefix preserves the existing scope-field contract: only exact + // full-key path-hop entries contribute transported scopes. Raw-prefix and + // byNode candidates can advance relay activity, but never broaden scope + // claims. fromPrefix bool + parsed bool + t time.Time + valid bool } -// collectRelayEntriesLocked returns deduplicated relayEntry snapshots for -// all StoreTx entries indexed under key (full pubkey) and its 1-byte wire -// prefix. Caller MUST hold s.mu at least for reading. -// -// byPathHop is keyed by both full resolved pubkey AND raw 1-byte hop -// prefix (e.g. "a3"). Many ingested non-advert packets only carry the -// raw hop on the wire — resolution to the full pubkey happens later via -// neighbor affinity. Looking up both keys and de-duping by tx ID matches -// what the "Paths seen through node" view shows. +func newRelayEntry(tx *StoreTx, fromPrefix bool) relayEntry { + pt, rt := -1, -1 + if tx.PayloadType != nil { + pt = *tx.PayloadType + } + if tx.RouteType != nil { + rt = *tx.RouteType + } + return relayEntry{ts: tx.FirstSeen, pt: pt, rt: rt, scope: tx.ScopeName, fromPrefix: fromPrefix} +} + +// collectRelayEntriesLocked returns one relayEntry per transmission that is a +// relay candidate for key (see forEachRelayCandidate) and whose raw observed +// flood paths confirm key. Caller MUST hold s.mu at least for reading. // -// The 1-byte prefix lookup CAN over-count when multiple nodes share the -// same first byte. This trades a possible over-count for clearly false -// zeros (issue #662). +// Candidate membership is never proof: persisted/live resolution can include +// heuristic guesses. Every candidate is verified against its raw observation +// paths, without target-biased resolve. Each transmission is checked once. func (s *PacketStore) collectRelayEntriesLocked(key string) []relayEntry { - txList := s.byPathHop[key] - var prefixList []*StoreTx - if len(key) >= 2 { - // key[:2] is the first 2 hex characters — exactly 1 byte of raw - // hop data, matching addTxToPathHopIndex for wire-level hops. - prefix := key[:2] - if prefix != key { - prefixList = s.byPathHop[prefix] - } + m := newRelayKeyMatcher(key, s.relayPrefixMapLocked()) + var entries []relayEntry + if !m.possible() { + return entries } + seen := make(map[int]struct{}) + forEachRelayCandidate(s.byPathHop, s.byNode, m, func(tx *StoreTx, fromPrefix bool) { + if _, done := seen[tx.ID]; done { + return + } + seen[tx.ID] = struct{}{} + if txConfirmsRelay(tx, m) { + entries = append(entries, newRelayEntry(tx, fromPrefix)) + } + }) + return entries +} - // Capacity hint: upper-bound is len(txList)+len(prefixList). The - // collect() pass below uses `seen` for true dedup, so we don't need - // a separate prepass (PR #1164 CR item 3: dead `uniq` map removed). - hint := len(txList) + len(prefixList) - entries := make([]relayEntry, 0, hint) - seen := make(map[int]bool, hint) - collect := func(list []*StoreTx, fromPrefix bool) { - for _, tx := range list { - if tx == nil || seen[tx.ID] { - continue - } - seen[tx.ID] = true - pt := -1 - if tx.PayloadType != nil { - pt = *tx.PayloadType - } - rt := -1 - if tx.RouteType != nil { - rt = *tx.RouteType +// forEachRelayCandidate visits key's candidate transmissions, possibly more +// than once across buckets (callers deduplicate by ID). Within a bucket it +// walks from the most recently indexed entry backwards; index order is not +// chronological, so callers must not stop early based on timestamps. Bucket +// order matters for scope provenance: the full-key path-hop bucket comes +// first and is the only one with fromPrefix=false. Reads live index slices: +// caller holds s.mu. +// - byPathHop[key]: raw full-key hops and live resolved-path indexing. +// - byNode[key]: decoded and resolved-path membership from every +// observation. Needed after a restart, where buildPathHopIndex re-indexes +// display paths only and drops non-display resolved hops. +// - byPathHop[unique 1/2/3-byte prefix]: raw wire hops. +func forEachRelayCandidate(pathHop, byNode map[string][]*StoreTx, m relayKeyMatcher, visit func(tx *StoreTx, fromPrefix bool)) { + each := func(list []*StoreTx, fromPrefix bool) { + for i := len(list) - 1; i >= 0; i-- { + if list[i] != nil { + visit(list[i], fromPrefix) } - entries = append(entries, relayEntry{ts: tx.FirstSeen, pt: pt, rt: rt, scope: tx.ScopeName, fromPrefix: fromPrefix}) } } - collect(txList, false) - collect(prefixList, true) - return entries + each(pathHop[m.key], false) + each(byNode[m.key], true) + for _, l := range relayPrefixHexLens { + if m.uniquePrefix(l) && m.key[:l] != m.key { + each(pathHop[m.key[:l]], true) + } + } } // computeRelayInfoFromEntries derives RepeaterRelayInfo from pre-snapshotted @@ -215,7 +227,10 @@ func computeRelayInfoFromEntries(entries []relayEntry, windowHours float64) Repe } scopeSet[e.scope] = struct{}{} } - t, ok := parseRelayTS(e.ts) + t, ok := e.t, e.valid + if !e.parsed { + t, ok = parseRelayTS(e.ts) + } if !ok { continue } @@ -271,8 +286,9 @@ func computeRelayInfoFromEntries(entries []relayEntry, windowHours float64) Repe } // GetRepeaterRelayInfo returns relay-activity information for a node by -// scanning the byPathHop index for non-advert packets that name the -// pubkey as a hop. It computes the most recent appearance timestamp, +// scanning byPathHop for non-advert flood packets with identity-safe observed +// hop evidence. Direct routes describe planned hops and are not relay evidence. +// It computes the most recent appearance timestamp, // 1h/24h hop counts, and whether the latest appearance falls within // windowHours. // diff --git a/cmd/server/repeater_liveness_test.go b/cmd/server/repeater_liveness_test.go index 07f108812..3d6652426 100644 --- a/cmd/server/repeater_liveness_test.go +++ b/cmd/server/repeater_liveness_test.go @@ -21,10 +21,11 @@ func TestRepeaterRelayActivity_Active(t *testing.T) { // A non-advert packet (payload_type=1, TXT_MSG) with the repeater pubkey // indexed as a path hop. Index by lowercase pubkey directly to mirror // the resolved-path entries that decode-window writes. - pt := 1 + pt, rt := 1, routeTypeFlood relayed := &StoreTx{ RawHex: "0100", PayloadType: &pt, + RouteType: &rt, PathJSON: `["aa"]`, FirstSeen: recentTS(2), } @@ -82,11 +83,11 @@ func seedUnscopedRelayFixture(t *testing.T, hashPrefix string) (*PacketStore, st } // assertUnscopedCounts pins the contract both lookups share: FLOOD hops count -// as unscoped, DIRECT hops only as plain relays. +// as unscoped; DIRECT planned hops do not establish observed relays. func assertUnscopedCounts(t *testing.T, info RepeaterRelayInfo) { t.Helper() - if info.RelayCount24h != 2 { - t.Errorf("expected RelayCount24h=2 (both hops), got %d", info.RelayCount24h) + if info.RelayCount24h != 1 { + t.Errorf("expected RelayCount24h=1 (only observed flood hop, not planned direct route), got %d", info.RelayCount24h) } if info.UnscopedRelayCount24h != 1 { t.Errorf("expected UnscopedRelayCount24h=1 (only the FLOOD hop), got %d", info.UnscopedRelayCount24h) @@ -139,11 +140,12 @@ func TestRepeaterRelayActivity_Stale(t *testing.T) { store := NewPacketStore(db, nil) - pt := 1 + pt, rt := 1, routeTypeFlood staleTS := time.Now().UTC().Add(-48 * time.Hour).Format("2006-01-02T15:04:05.000Z") old := &StoreTx{ RawHex: "0100", PayloadType: &pt, + RouteType: &rt, PathJSON: `["11"]`, FirstSeen: staleTS, } @@ -235,10 +237,11 @@ func TestRepeaterRelayActivity_PrefixHop(t *testing.T) { // Non-advert packet with a single raw 1-byte hop matching the target // pubkey's first byte ("a3"). Index it the way addTxToPathHopIndex // does — under the raw hop key only, not the full pubkey. - pt := 1 + pt, rt := 1, routeTypeFlood tx := &StoreTx{ RawHex: "0100", PayloadType: &pt, + RouteType: &rt, PathJSON: `["a3"]`, FirstSeen: recentTS(2), } @@ -262,6 +265,10 @@ func TestRepeaterRelayActivity_PrefixHop(t *testing.T) { if !info.RelayActive { t.Errorf("expected RelayActive=true within 24h window, got false (LastRelayed=%s)", info.LastRelayed) } + bulk := store.computeRepeaterRelayInfoMap(24)[pubkey] + if bulk.RelayCount24h != info.RelayCount24h || bulk.LastRelayed != info.LastRelayed || bulk.RelayActive != info.RelayActive { + t.Fatalf("unique-prefix bulk/single mismatch: single=%+v bulk=%+v", info, bulk) + } } // TestRepeaterRelayActivity_DedupAcrossPrefixAndFullKey verifies that when @@ -279,10 +286,11 @@ func TestRepeaterRelayActivity_DedupAcrossPrefixAndFullKey(t *testing.T) { store := NewPacketStore(db, nil) - pt := 1 + pt, rt := 1, routeTypeFlood tx := &StoreTx{ RawHex: "0100", PayloadType: &pt, + RouteType: &rt, PathJSON: `["a3"]`, FirstSeen: recentTS(2), } diff --git a/cmd/server/repeater_relay_concurrency_test.go b/cmd/server/repeater_relay_concurrency_test.go new file mode 100644 index 000000000..388309fcd --- /dev/null +++ b/cmd/server/repeater_relay_concurrency_test.go @@ -0,0 +1,164 @@ +package main + +import ( + "reflect" + "strings" + "sync" + "testing" + "time" +) + +// Bulk relay info runs concurrently with ingest and eviction, which compacts +// byNode/byPathHop slices in place. Every bulk result must be the exact +// single-path result of one store state the reader could have observed: +// a candidate from a later state must never be credited (or deduplicated) +// as some other transmission. Run with -race. +func TestRepeaterRelayInfoMap_ConcurrentIngestAndEviction(t *testing.T) { + relayA := "a1b2c3" + strings.Repeat("11", 29) + relayB := "b7c8d9" + strings.Repeat("22", 29) + hop := "d4e5f6" + strings.Repeat("33", 29) + keys := []string{relayA, relayB, hop} + db, store := activityTestStore(t, keys...) + defer db.conn.Close() + + // Timestamps sit hours away from the 1h/24h windows, so results do not + // depend on when each reader samples time.Now. + base := time.Now().UTC().Add(-20 * time.Hour).Truncate(time.Second) + const initialTx, rounds, txPerRound, keepTx = 200, 60, 20, 40 + nextID := 1 + resolved := map[int][]string{} + hopsSeen := map[string]bool{} + addLocked := func() { + id := nextID + nextID++ + ts := base.Add(time.Duration(id) * time.Second).Format(time.RFC3339) + // Relay identity comes only from a shorter non-display observation; + // alternate transmissions credit different relays. + relay, alt := relayA, `["A1"]` + if id%2 == 1 { + relay, alt = relayB, `["B7C8"]` + } + tx := storeObservedTx(store, id, 2, ts, `{"type":"TXT_MSG"}`, alt, `["D4E5F6","D4E5F6","D4E5F6"]`) + resolved[id] = []string{relay, hop} + // Restart shape: the relay is reachable only through byNode, whose + // slices eviction compacts in place, so a stale candidate would + // change the relay result instead of being hidden by a full-key + // path-hop entry. + store.addToByNode(tx, relay) + store.indexResolvedPathHops(tx, []string{hop}, hopsSeen) + } + // oracle[round] is the single-path relay info for every key in the store + // state published as round, computed under the writer's lock. + var oracleMu sync.Mutex + var oracle []map[string]RepeaterRelayInfo + round := 0 + publishLocked := func() { + want := make(map[string]RepeaterRelayInfo, len(keys)) + for _, key := range keys { + want[key] = computeRelayInfoFromEntries(store.collectRelayEntriesLocked(key), 24) + } + oracleMu.Lock() + oracle = append(oracle, want) + oracleMu.Unlock() + } + store.mu.Lock() + for i := 0; i < initialTx; i++ { + addLocked() + } + publishLocked() + store.mu.Unlock() + + readRound := func() int { + store.mu.RLock() + defer store.mu.RUnlock() + return round + } + matchesSomeRound := func(got map[string]RepeaterRelayInfo, from, to int) bool { + oracleMu.Lock() + defer oracleMu.Unlock() + for r := from; r <= to; r++ { + same := true + for _, key := range keys { + if !reflect.DeepEqual(got[key], oracle[r][key]) { + same = false + break + } + } + if same { + return true + } + } + return false + } + + const readers = 3 + var wg sync.WaitGroup + ready := make(chan struct{}, readers) + stop := make(chan struct{}) + failures := make(chan string, readers) + for r := 0; r < readers; r++ { + wg.Add(1) + go func(r int) { + defer wg.Done() + ready <- struct{}{} + for first := true; ; first = false { + select { + case <-stop: + if !first { + return + } + default: + } + from := readRound() + if r == readers-1 { + // Health shares the candidate buckets; exercised for races. + _ = store.GetBulkHealth(10, "", "") + continue + } + got := store.computeRepeaterRelayInfoMap(24) + if to := readRound(); !matchesSomeRound(got, from, to) { + failures <- "bulk relay info matches no store state between rounds" + return + } + } + }(r) + } + for r := 0; r < readers; r++ { + <-ready + } + + for i := 0; i < rounds; i++ { + store.mu.Lock() + for j := 0; j < txPerRound; j++ { + addLocked() + } + // Evict older transmissions (1s apart). The cutoff sits mid-second so + // clock progress during the round cannot move it; eviction formats it + // with whole seconds, so keepTx+1 transmissions remain. + keepFrom := base.Add(time.Duration(nextID-keepTx)*time.Second - 500*time.Millisecond) + store.retentionHours = time.Since(keepFrom).Hours() + store.EvictStaleWithRP(resolved) + round++ + publishLocked() + store.mu.Unlock() + } + close(stop) + wg.Wait() + close(failures) + for failure := range failures { + t.Fatal(failure) + } + + // Final synchronization point: no writer is active. + got := store.computeRepeaterRelayInfoMap(24) + if !matchesSomeRound(got, round, round) { + t.Fatalf("final bulk relay info differs from single path: %+v", got) + } + // Each retained transmission credits exactly one relay. (The shared hop + // is also indexed under its full key, which eviction leaves in place; + // that policy is out of scope here.) + final := oracle[round] + if final[relayA].RelayCount24h == 0 || final[relayB].RelayCount24h == 0 || final[relayA].RelayCount24h+final[relayB].RelayCount24h != len(store.packets) || len(store.packets) != keepTx+1 { + t.Fatalf("fixture relay evidence inconsistent: packets=%d oracle=%+v", len(store.packets), final) + } +} diff --git a/cmd/server/routes_test.go b/cmd/server/routes_test.go index 0f8f4fe45..9eae00520 100644 --- a/cmd/server/routes_test.go +++ b/cmd/server/routes_test.go @@ -1023,8 +1023,8 @@ func TestNodeHealthPartialFromPackets(t *testing.T) { if stats["totalPackets"] != 1.0 { // JSON numbers are float64 t.Errorf("expected totalPackets=1, got %v", stats["totalPackets"]) } - if stats["lastHeard"] == nil { - t.Error("expected lastHeard to be set") + if stats["lastHeard"] != nil || stats["lastAdvert"] != nil { + t.Error("unattributed packet analytics must not imply confirmed node activity or advert") } } @@ -4883,12 +4883,15 @@ func TestHandleScopeStats_RepeatersByRegion(t *testing.T) { } pt5 := 5 // GRP_TXT — non-advert, so it counts toward TransportedScopes + flood := routeTypeFlood tx := &StoreTx{ ID: 1, Hash: "txhash1", FirstSeen: time.Now().UTC().Add(-5 * time.Minute).Format(time.RFC3339Nano), PayloadType: &pt5, ScopeName: "#belgium", + RouteType: &flood, + PathJSON: `["aa"]`, } // #1751 follow-up regression: byPathHop also indexes short hex-prefix // "bucket" keys (ambiguous-hop resolution fallback) alongside full @@ -4902,6 +4905,7 @@ func TestHandleScopeStats_RepeatersByRegion(t *testing.T) { ScopeName: "#belgium", } srv.store = &PacketStore{ + nodePM: buildPrefixMap([]nodeInfo{{PublicKey: "aabbccdd0011", Role: "repeater"}}), byPathHop: map[string][]*StoreTx{ "aabbccdd0011": {tx}, "aabb": {bucketTx}, @@ -4957,8 +4961,14 @@ func TestHandleScopeStats_BridgeRepeaters(t *testing.T) { txBelgium := &StoreTx{ID: 1, Hash: "tx1", FirstSeen: time.Now().UTC().Add(-5 * time.Minute).Format(time.RFC3339Nano), PayloadType: &pt5, ScopeName: "#belgium"} txFrance := &StoreTx{ID: 2, Hash: "tx2", FirstSeen: time.Now().UTC().Add(-5 * time.Minute).Format(time.RFC3339Nano), PayloadType: &pt5, ScopeName: "#france"} txBelgium2 := &StoreTx{ID: 3, Hash: "tx3", FirstSeen: time.Now().UTC().Add(-5 * time.Minute).Format(time.RFC3339Nano), PayloadType: &pt5, ScopeName: "#belgium"} + flood := routeTypeFlood + for _, tx := range []*StoreTx{txBelgium, txFrance} { + tx.RouteType, tx.PathJSON = &flood, `["bb"]` + } + txBelgium2.RouteType, txBelgium2.PathJSON = &flood, `["cc"]` srv.store = &PacketStore{ + nodePM: buildPrefixMap([]nodeInfo{{PublicKey: "bbbbccdd0011", Role: "repeater"}, {PublicKey: "ccccccdd0011", Role: "repeater"}}), byPathHop: map[string][]*StoreTx{ "bbbbccdd0011": {txBelgium, txFrance}, // relayed BOTH regions — a bridge "ccccccdd0011": {txBelgium2}, // relayed only #belgium — not a bridge diff --git a/cmd/server/store.go b/cmd/server/store.go index 7a4a6ba86..3c180119d 100644 --- a/cmd/server/store.go +++ b/cmd/server/store.go @@ -9309,13 +9309,16 @@ func (s *PacketStore) GetBulkHealth(limit int, region, area string) []map[string todayStart := time.Now().UTC().Truncate(24 * time.Hour).Format(time.RFC3339) results := make([]map[string]interface{}, 0, len(nodes)) + pm := s.relayPrefixMapLocked() + relaySeen := make(map[int]struct{}) // request-local relay-evidence dedupe scratch for _, n := range nodes { packets := s.byNode[n.pk] + activityKey := strings.ToLower(n.pk) var packetsToday int var snrSum float64 var snrCount int - var lastHeard string + var lastHeard, lastAdvert string observerStats := map[string]*struct { name string snrSum, rssiSum float64 @@ -9335,9 +9338,7 @@ func (s *PacketStore) GetBulkHealth(limit int, region, area string) []map[string snrSum += *pkt.SNR snrCount++ } - if lastHeard == "" || pkt.FirstSeen > lastHeard { - lastHeard = pkt.FirstSeen - } + updateOwnAdvertActivity(pkt, activityKey, &lastHeard, &lastAdvert) obsID := pkt.ObserverID if obsID != "" { obs := observerStats[obsID] @@ -9360,6 +9361,7 @@ func (s *PacketStore) GetBulkHealth(limit int, region, area string) []map[string } } } + s.updateIndexedRelayActivityLocked(activityKey, pm, &lastHeard, relaySeen) observerRows := make([]map[string]interface{}, 0) for id, o := range observerStats { @@ -9379,13 +9381,10 @@ func (s *PacketStore) GetBulkHealth(limit int, region, area string) []map[string return observerRows[i]["packetCount"].(int) > observerRows[j]["packetCount"].(int) }) - var avgSnr interface{} + var avgSnr *float64 if snrCount > 0 { - avgSnr = snrSum / float64(snrCount) - } - var lhVal interface{} - if lastHeard != "" { - lhVal = lastHeard + v := snrSum / float64(snrCount) + avgSnr = &v } results = append(results, map[string]interface{}{ @@ -9394,13 +9393,10 @@ func (s *PacketStore) GetBulkHealth(limit int, region, area string) []map[string "role": nilIfEmpty(n.role), "lat": n.lat, "lon": n.lon, - "stats": map[string]interface{}{ - "totalTransmissions": len(packets), - "totalObservations": totalObservations, - "totalPackets": len(packets), - "packetsToday": packetsToday, - "avgSnr": avgSnr, - "lastHeard": lhVal, + "stats": NodeHealthStats{ + TotalTransmissions: len(packets), TotalObservations: totalObservations, + TotalPackets: len(packets), PacketsToday: packetsToday, AvgSnr: avgSnr, + LastHeard: timestampPointer(lastHeard), LastAdvert: timestampPointer(lastAdvert), }, "observers": observerRows, }) @@ -9439,13 +9435,15 @@ func (s *PacketStore) GetNodeHealth(pubkey string) (map[string]interface{}, erro defer s.mu.RUnlock() packets := s.byNode[pubkey] + activityKey := strings.ToLower(pubkey) todayStart := time.Now().UTC().Truncate(24 * time.Hour).Format(time.RFC3339) var packetsToday int var snrSum float64 var snrCount int var totalHops, hopCount int - var lastHeard string + var lastHeard, lastAdvert string + pm := s.relayPrefixMapLocked() totalObservations := 0 observerStats := map[string]*struct { @@ -9463,9 +9461,7 @@ func (s *PacketStore) GetNodeHealth(pubkey string) (map[string]interface{}, erro snrSum += *pkt.SNR snrCount++ } - if lastHeard == "" || pkt.FirstSeen > lastHeard { - lastHeard = pkt.FirstSeen - } + updateOwnAdvertActivity(pkt, activityKey, &lastHeard, &lastAdvert) // Hop counting hops := txGetParsedPath(pkt) if len(hops) > 0 { @@ -9495,6 +9491,7 @@ func (s *PacketStore) GetNodeHealth(pubkey string) (map[string]interface{}, erro } } } + s.updateIndexedRelayActivityLocked(activityKey, pm, &lastHeard, make(map[int]struct{})) observerRows := make([]map[string]interface{}, 0) // Issue #1290: surface listener/repeater hint on node detail by @@ -9552,18 +9549,15 @@ func (s *PacketStore) GetNodeHealth(pubkey string) (map[string]interface{}, erro return observerRows[i]["packetCount"].(int) > observerRows[j]["packetCount"].(int) }) - var avgSnr interface{} + var avgSnr *float64 if snrCount > 0 { - avgSnr = snrSum / float64(snrCount) + v := snrSum / float64(snrCount) + avgSnr = &v } avgHops := 0 if hopCount > 0 { avgHops = int(math.Round(float64(totalHops) / float64(hopCount))) } - var lhVal interface{} - if lastHeard != "" { - lhVal = lastHeard - } // Recent packets (up to 20, newest first — read from tail of oldest-first slice) recentLimit := 20 @@ -9580,14 +9574,10 @@ func (s *PacketStore) GetNodeHealth(pubkey string) (map[string]interface{}, erro return map[string]interface{}{ "node": node, "observers": observerRows, - "stats": map[string]interface{}{ - "totalTransmissions": len(packets), - "totalObservations": totalObservations, - "totalPackets": len(packets), - "packetsToday": packetsToday, - "avgSnr": avgSnr, - "avgHops": avgHops, - "lastHeard": lhVal, + "stats": NodeHealthStats{ + TotalTransmissions: len(packets), TotalObservations: totalObservations, + TotalPackets: len(packets), PacketsToday: packetsToday, AvgSnr: avgSnr, + AvgHops: &avgHops, LastHeard: timestampPointer(lastHeard), LastAdvert: timestampPointer(lastAdvert), }, "recentPackets": recentPackets, }, nil diff --git a/cmd/server/transported_scopes_1751_test.go b/cmd/server/transported_scopes_1751_test.go index 4e9770da6..cb627dc96 100644 --- a/cmd/server/transported_scopes_1751_test.go +++ b/cmd/server/transported_scopes_1751_test.go @@ -20,16 +20,18 @@ import ( // so the field stays in parity for /api/nodes (bulk) and the single-node // detail endpoint (per-node). -const scope1751Key = "aabbccdd11223344" +const scope1751Key = "aabbccdd11223344000000000000000000000000000000000000000000000000" // scopeTx builds a path-hop StoreTx with the given payload type, scope name, // and an in-window FirstSeen. func scopeTx(id int, payloadType int, scope string) *StoreTx { - pt := payloadType + pt, rt := payloadType, routeTypeFlood return &StoreTx{ ID: id, Hash: "scope-tx-" + scope + "-" + strconv.Itoa(id), PayloadType: &pt, + RouteType: &rt, + PathJSON: `["` + scope1751Key + `"]`, ScopeName: scope, FirstSeen: time.Now().UTC().Add(-10 * time.Minute).Format(time.RFC3339Nano), } @@ -112,7 +114,8 @@ func TestTransportedScopes_EmptyWhenNoScope(t *testing.T) { func TestTransportedScopes_PrefixBucketExcludedFromScope(t *testing.T) { full := scopeTx(1, 2, "region-direct") // only in the full-key bucket — must count prefixOnly := scopeTx(2, 2, "region-via-prefix") // only in the 1-byte bucket — must NOT count - shared := scopeTx(3, 2, "region-shared") // in BOTH buckets — must count (present in the full-key list) + prefixOnly.PathJSON = `["aa"]` + shared := scopeTx(3, 2, "region-shared") // in BOTH buckets — must count (present in the full-key list) store := &PacketStore{ byPathHop: map[string][]*StoreTx{ @@ -137,6 +140,7 @@ func TestTransportedScopes_PrefixBucketExcludedFromScope(t *testing.T) { func TestTransportedScopes_PerNodePrefixBucketExcludedFromScope(t *testing.T) { full := scopeTx(1, 2, "region-direct") prefixOnly := scopeTx(2, 2, "region-via-prefix") + prefixOnly.PathJSON = `["aa"]` shared := scopeTx(3, 2, "region-shared") store := &PacketStore{ diff --git a/cmd/server/transported_scopes_recent_test.go b/cmd/server/transported_scopes_recent_test.go index a58569833..cafd32f2c 100644 --- a/cmd/server/transported_scopes_recent_test.go +++ b/cmd/server/transported_scopes_recent_test.go @@ -21,15 +21,17 @@ import ( // - computeRepeaterRelayInfoMap (bulk, repeater_enrich_bulk.go) // - GetRepeaterRelayInfo (per-node, repeater_liveness.go) -const scopeRecentKey = "ffeeddcc55667788" +const scopeRecentKey = "ffeeddcc55667788000000000000000000000000000000000000000000000000" // scopeTxAt builds a path-hop StoreTx with an explicit FirstSeen age. func scopeTxAt(id int, payloadType int, scope string, age time.Duration) *StoreTx { - pt := payloadType + pt, rt := payloadType, routeTypeFlood return &StoreTx{ ID: id, Hash: "scope-recent-tx-" + scope + "-" + strconv.Itoa(id), PayloadType: &pt, + RouteType: &rt, + PathJSON: `["` + scopeRecentKey + `"]`, ScopeName: scope, FirstSeen: time.Now().UTC().Add(-age).Format(time.RFC3339Nano), } @@ -104,6 +106,7 @@ func TestTransportedScopesRecent_DisabledWhenWindowHoursZero(t *testing.T) { // region matches, so an unresolved hop can't be credited with any scope). func TestTransportedScopesRecent_PrefixBucketExcluded(t *testing.T) { prefixOnly := scopeTxAt(1, 2, "region-via-prefix", 10*time.Minute) + prefixOnly.PathJSON = `["ff"]` store := &PacketStore{ byPathHop: map[string][]*StoreTx{ diff --git a/public/nodes.js b/public/nodes.js index db73c94b4..94e6ead44 100644 --- a/public/nodes.js +++ b/public/nodes.js @@ -203,6 +203,47 @@ /* === Shared helper functions for node detail rendering === */ + function getLastAdvert(n, stats, adverts) { + // New health APIs distinguish own adverts from arbitrary packet analytics. + // Null is authoritative: never relabel relay-touched last_seen as an advert. + if (Object.prototype.hasOwnProperty.call(stats, 'lastAdvert')) return stats.lastAdvert || null; + let latest = null; + // Compatibility with older APIs: exact origin, never a path/prefix match. + for (const packet of adverts || []) { + if (Number(packet.payload_type) !== 4) continue; + let decoded = packet.decoded_json; + if (typeof decoded === 'string') { + try { decoded = JSON.parse(decoded); } catch (_) { continue; } + } + if (decoded && decoded.signatureValid === false) continue; + const source = decoded && (decoded.pubKey || decoded.publicKey); + if (typeof source !== 'string' || source.toLowerCase() !== String(n.public_key || '').toLowerCase()) continue; + const ts = packet.first_seen || packet.timestamp; + if (ts && Number.isFinite(Date.parse(ts)) && (!latest || Date.parse(ts) > Date.parse(latest))) latest = ts; + } + return latest; + } + + function nodeWithHealthActivity(n, stats, adverts) { + // Directory last_seen/last_heard can include historical heuristic relay + // inference. Do not let them outrank the identity-safe health contract. + const activity = Object.prototype.hasOwnProperty.call(stats, 'lastAdvert') + ? stats.lastHeard || null + : getLastAdvert(n, stats, adverts); + return Object.assign({}, n, { _lastHeard: activity, last_heard: null, last_seen: null, _liveSeen: null }); + } + + function renderRelayActivity(n) { + let html = n.last_relayed ? renderNodeTimestampHtml(n.last_relayed) + ' ' : ''; + html += n.last_relayed && n.relay_active + ? 'actively relaying' + : 'no recent confirmed relay'; + if (n.relay_count_1h != null || n.relay_count_24h != null) { + html += ` (${n.relay_count_1h || 0} relays/hr, ${n.relay_count_24h || 0} relays/24h)`; + } + return html; + } + function getStatusTooltip(role, status) { const isInfra = role === 'repeater' || role === 'room'; const threshMs = isInfra ? HEALTH_THRESHOLDS.infraSilentMs : HEALTH_THRESHOLDS.nodeSilentMs; @@ -660,7 +701,7 @@ const nodeData = await fetchNodeDetail(pubkey); if (viewSeq !== detailViewSeq || !body.isConnected) return; const healthData = nodeData.healthData; - const n = nodeData.node; + const n = nodeWithHealthActivity(nodeData.node, (healthData && healthData.stats) || {}, nodeData.recentAdverts || []); const adverts = (nodeData.recentAdverts || []).sort((a, b) => new Date(b.timestamp) - new Date(a.timestamp)); const title = document.querySelector('.node-full-title'); if (title) title.textContent = n.name || pubkey.slice(0, 12); @@ -683,10 +724,9 @@ const stats = h.stats || {}; const observers = h.observers || []; const recent = h.recentPackets || []; - const lastHeard = stats.lastHeard; + const lastHeard = getLastAdvert(n, stats, adverts); - // Attach health lastHeard for shared helpers - n._lastHeard = lastHeard || n.last_seen; + // General activity was supplied separately by nodeWithHealthActivity. const si = getStatusInfo(n); const roleColor = si.roleColor; const statusLabel = si.statusLabel; @@ -724,8 +764,8 @@ - - ${(n.role === 'repeater' || n.role === 'room') ? `` : ''} + + ${(n.role === 'repeater' || n.role === 'room') ? `` : ''} ${(n.role === 'repeater' || n.role === 'room') && (n.traffic_share_score != null || n.usefulness_score != null) ? (() => { // #1456: prefer the new traffic_share_score field; fall back // to legacy usefulness_score for graceful degradation @@ -1788,7 +1828,7 @@ } function renderDetail(panel, data) { - const n = data.node; + const n = nodeWithHealthActivity(data.node, (data.healthData && data.healthData.stats) || {}, data.recentAdverts || []); const adverts = (data.recentAdverts || []).sort((a, b) => new Date(b.timestamp) - new Date(a.timestamp)); const h = data.healthData || {}; const stats = h.stats || {}; @@ -1801,8 +1841,7 @@ const nodeUrl = location.origin + '/#/nodes/' + encodeURIComponent(n.public_key); // Status calculation via shared helper - const lastHeard = stats.lastHeard; - n._lastHeard = lastHeard || n.last_seen; + const lastHeard = getLastAdvert(n, stats, adverts); const si = getStatusInfo(n); const roleColor = si.roleColor; const totalPackets = stats.totalTransmissions || stats.totalPackets || n.advert_count || 0; @@ -1834,7 +1873,8 @@

Overview

-
Last Heard
${renderNodeTimestampHtml(lastHeard || n.last_seen)}
+
Last Heard (advert)
${renderNodeTimestampHtml(lastHeard)}
+ ${(n.role === 'repeater' || n.role === 'room') ? `
Last Relayed
${renderRelayActivity(n)}
` : ''}
First Seen
${renderNodeTimestampHtml(n.first_seen)}
Total Packets
${totalPackets}
Packets Today
${stats.packetsToday || 0}
@@ -2079,6 +2119,9 @@ }; window._nodesSyncClaimedToFavorites = syncClaimedToFavorites; window._nodesRenderNodeTimestampHtml = renderNodeTimestampHtml; + window._nodesGetLastAdvert = getLastAdvert; + window._nodesWithHealthActivity = nodeWithHealthActivity; + window._nodesRenderRelayActivity = renderRelayActivity; window._nodesRenderNodeTimestampText = renderNodeTimestampText; window._nodesGetStatusInfo = getStatusInfo; window._nodesGetStatusTooltip = getStatusTooltip; diff --git a/test-frontend-helpers.js b/test-frontend-helpers.js index 5dbd47e57..2668c05c6 100644 --- a/test-frontend-helpers.js +++ b/test-frontend-helpers.js @@ -4490,6 +4490,69 @@ console.log('\n=== nodes.js: renderNodeTimestampHtml / renderNodeTimestampText = }); } +console.log('\n=== nodes.js: own advert freshness ==='); +{ + test('advert timestamp ignores recent arbitrary health traffic and relay-touched last_seen', () => { + const ctx = makeNodesSandbox(); + const key = 'aa' + '11'.repeat(31); + const stale = new Date(Date.now() - 96 * 3600000).toISOString(); + const recent = new Date().toISOString(); + const n = { public_key: key, role: 'repeater', last_seen: stale, last_heard: stale }; + const advert = { payload_type: 4, timestamp: stale, decoded_json: JSON.stringify({ pubKey: key }) }; + const lastAdvert = ctx.window._nodesGetLastAdvert(n, { lastHeard: recent, lastAdvert: stale }, [advert]); + assert.strictEqual(lastAdvert, stale); + n._lastHeard = lastAdvert; + assert.strictEqual(ctx.window._nodesGetStatusInfo(n).status, 'stale'); + }); + test('legacy API fallback uses only exact-owned adverts, never ACKs or another origin', () => { + const ctx = makeNodesSandbox(); + const key = 'aa' + '11'.repeat(31); + const old = new Date(Date.now() - 96 * 3600000).toISOString(); + const recent = new Date().toISOString(); + const adverts = [ + { payload_type: 4, timestamp: old, decoded_json: JSON.stringify({ pubKey: key.toUpperCase() }) }, + { payload_type: 4, timestamp: recent, decoded_json: JSON.stringify({ pubKey: 'bb' + '22'.repeat(31) }) }, + { payload_type: 3, timestamp: recent, decoded_json: JSON.stringify({ pubKey: key }) }, + { payload_type: 4, timestamp: recent, decoded_json: JSON.stringify({ pubKey: key, signatureValid: false }) }, + ]; + assert.strictEqual(ctx.window._nodesGetLastAdvert({ public_key: key, last_seen: recent }, { lastHeard: recent }, adverts), old); + }); + test('explicit unknown advert is null, even if legacy packet evidence or last_seen is fresh', () => { + const ctx = makeNodesSandbox(); + const key = 'aa' + '11'.repeat(31); + const recent = new Date().toISOString(); + assert.strictEqual(ctx.window._nodesGetLastAdvert({ public_key: key, last_seen: recent }, { lastAdvert: null, lastHeard: recent }, []), null); + assert.strictEqual(ctx.window._nodesGetLastAdvert({ public_key: key, last_seen: recent }, { lastHeard: recent }, []), null); + }); + test('safe health activity outranks polluted directory timestamps without mutating the node', () => { + const ctx = makeNodesSandbox(); + const old = new Date(Date.now() - 96 * 3600000).toISOString(); + const recent = new Date().toISOString(); + const node = { role: 'repeater', last_seen: recent, last_heard: recent, _liveSeen: Date.now() }; + const safe = ctx.window._nodesWithHealthActivity(node, { lastAdvert: old, lastHeard: old }, []); + assert.strictEqual(ctx.window._nodesGetStatusInfo(safe).status, 'stale'); + assert.strictEqual(node.last_seen, recent); + const unknown = ctx.window._nodesWithHealthActivity(node, { lastAdvert: null, lastHeard: null }, []); + assert.strictEqual(ctx.window._nodesGetStatusInfo(unknown).status, 'stale'); + }); + test('relay-only safe general activity is active but the advert timestamp remains unknown', () => { + const ctx = makeNodesSandbox(); + const stats = { lastAdvert: null, lastHeard: new Date().toISOString() }; + const node = { role: 'repeater' }; + const safe = ctx.window._nodesWithHealthActivity(node, stats, []); + assert.strictEqual(ctx.window._nodesGetStatusInfo(safe).status, 'active'); + assert.strictEqual(ctx.window._nodesGetLastAdvert(node, stats, []), null); + }); + test('absent or historical relays do not claim the node is alive', () => { + const ctx = makeNodesSandbox(); + for (const node of [{}, { last_relayed: new Date(Date.now() - 96 * 3600000).toISOString(), relay_active: false }]) { + const html = ctx.window._nodesRenderRelayActivity(node); + assert.ok(html.includes('no recent confirmed relay')); + assert.ok(!html.includes('alive')); + } + }); +} + // ===== NODES.JS: getStatusInfo edge cases (P0 coverage expansion) ===== console.log('\n=== nodes.js: getStatusInfo edge cases ==='); { diff --git a/test-node-liveness-e2e.js b/test-node-liveness-e2e.js new file mode 100644 index 000000000..e07d7b085 --- /dev/null +++ b/test-node-liveness-e2e.js @@ -0,0 +1,107 @@ +// Actual SPA/status rendering with synthetic API evidence; never contacts prod. +// Run: node test-node-liveness-e2e.js (Playwright Chromium required). +'use strict'; +const assert = require('node:assert/strict'); +const fs = require('node:fs'); +const http = require('node:http'); +const path = require('node:path'); +const { chromium } = require('playwright'); + +const publicDir = path.join(__dirname, 'public'); +const stale = new Date(Date.now() - 4 * 86400000).toISOString(); +const recent = new Date(Date.now() - 20 * 60000).toISOString(); +const cases = [ + { name: 'Offline collision', active: false, advert: stale, heard: stale }, + { name: 'Own advert', active: true, advert: recent, heard: recent }, + { name: 'Relay without advert', active: true, advert: null, heard: recent, relay: true }, + { name: 'Legacy unsafe health', active: false, advert: stale, heard: recent, legacy: true }, + { name: 'Unknown safe activity', active: false, advert: null, heard: null }, +].map((fixture, index) => ({ ...fixture, key: (index + 1).toString(16).padStart(2, '0') + '11'.repeat(31) })); +const nodes = cases.map(fixture => ({ + public_key: fixture.key, name: fixture.name, role: 'repeater', + lat: null, lon: null, last_seen: recent, last_heard: recent, + first_seen: stale, advert_count: fixture.advert ? 1 : 0, + last_relayed: fixture.relay ? recent : null, relay_active: !!fixture.relay, + relay_count_1h: fixture.relay ? 1 : 0, relay_count_24h: fixture.relay ? 1 : 0, +})); + +function json(res, data) { + res.writeHead(200, { 'Content-Type': 'application/json', 'Cache-Control': 'no-store' }); + res.end(JSON.stringify(data)); +} +const server = http.createServer((req, res) => { + const pathname = new URL(req.url, 'http://127.0.0.1').pathname; + if (pathname === '/api/nodes') return json(res, { nodes, total: nodes.length, counts: { all: nodes.length, repeater: nodes.length } }); + const match = pathname.match(/^\/api\/nodes\/([^/]+)(?:\/(.*))?$/); + if (match) { + const index = cases.findIndex(fixture => fixture.key === match[1]); + if (index >= 0) { + const fixture = cases[index], node = nodes[index]; + if (match[2] === 'health') return json(res, { + node, observers: [], recentPackets: [], stats: { + lastHeard: fixture.heard, ...(fixture.legacy ? {} : { lastAdvert: fixture.advert }), + totalPackets: 2, totalTransmissions: 2, totalObservations: 2, packetsToday: 1, + }, + }); + if (match[2] === 'neighbors') return json(res, { neighbors: [] }); + if (match[2] === 'paths') return json(res, { paths: [], totalTransmissions: 0 }); + if (!match[2]) return json(res, { + node, recentAdverts: fixture.advert ? [{ + id: 1, hash: 'synthetic-advert', payload_type: 4, route_type: 1, + first_seen: fixture.advert, timestamp: fixture.advert, + decoded_json: JSON.stringify({ type: 'ADVERT', pubKey: fixture.key, signatureValid: true }), observations: [], + }] : [], + }); + } + } + if (pathname === '/api/observers') return json(res, { observers: [] }); + if (pathname === '/api/channels') return json(res, { channels: [] }); + if (pathname.startsWith('/api/')) return json(res, {}); + const file = pathname === '/' ? path.join(publicDir, 'index.html') : path.resolve(publicDir, '.' + pathname); + if (!file.startsWith(publicDir + path.sep) || !fs.existsSync(file) || !fs.statSync(file).isFile()) { + res.writeHead(404); return res.end(); + } + const types = { '.js': 'application/javascript', '.css': 'text/css', '.html': 'text/html', '.svg': 'image/svg+xml' }; + res.writeHead(200, { 'Content-Type': types[path.extname(file)] || 'application/octet-stream' }); + res.end(fs.readFileSync(file)); +}); + +(async () => { + let browser; + let checks = 0; + try { + await new Promise(resolve => server.listen(0, '127.0.0.1', resolve)); + const origin = 'http://127.0.0.1:' + server.address().port; + browser = await chromium.launch({ headless: true, executablePath: process.env.CHROMIUM_PATH || undefined }); + const page = await browser.newPage({ viewport: { width: 1440, height: 1000 } }); + // Exclude CDN/fonts/map tiles and WebSocket reconnects from this local test. + await page.route('**/*', route => new URL(route.request().url()).origin === origin ? route.continue() : route.abort()); + for (const fixture of cases) { + await page.goto(origin + '/#/nodes/' + fixture.key, { waitUntil: 'domcontentloaded' }); + const body = page.locator('#nodeFullBody'); + await body.locator('tr').filter({ hasText: 'Last Heard (advert)' }).waitFor(); + const text = await body.innerText(); + assert.match(text, fixture.active ? /Active/ : /Stale/, fixture.name + ': full detail status'); checks++; + assert.doesNotMatch(text, fixture.active ? /Stale/ : /Active/, fixture.name + ': no contradictory full detail status'); checks++; + assert.ok(!text.includes('alive (idle)'), fixture.name + ': no unsupported alive claim'); checks++; + const advertRow = body.locator('tr').filter({ hasText: 'Last Heard (advert)' }); + assert.match(await advertRow.innerText(), fixture.advert ? (fixture.active ? /\d+m ago/ : /4d ago/) : /—/); checks++; + await page.goto(origin + '/#/nodes', { waitUntil: 'domcontentloaded' }); + await page.reload({ waitUntil: 'domcontentloaded' }); + await page.locator('tr[data-key="' + fixture.key + '"]').click(); + const pane = page.locator('#nodesRight'); + await pane.locator('dt').filter({ hasText: 'Last Heard (advert)' }).waitFor(); + const paneText = await pane.innerText(); + assert.match(paneText, fixture.active ? /Active/ : /Stale/, fixture.name + ': sidepane status'); checks++; + assert.doesNotMatch(paneText, fixture.active ? /Stale/ : /Active/, fixture.name + ': no contradictory sidepane status'); checks++; + const advertValue = pane.locator('dt').filter({ hasText: 'Last Heard (advert)' }).locator('xpath=following-sibling::dd[1]'); + assert.match(await advertValue.innerText(), fixture.advert ? (fixture.active ? /\d+m ago/ : /4d ago/) : /—/); checks++; + if (process.env.SCREENSHOT_DIR) await page.screenshot({ path: path.join(process.env.SCREENSHOT_DIR, 'node-liveness-' + fixture.key.slice(0, 2) + '.png') }); + console.log('PASS ' + fixture.name); + } + console.log(checks + ' checks passed'); + } finally { + if (browser) await browser.close(); + await new Promise(resolve => server.close(resolve)); + } +})().catch(error => { console.error(error); process.exitCode = 1; });
Status${statusLabel} ${statusExplanation}
Last Heard${renderNodeTimestampHtml(lastHeard || n.last_seen)}
Last Relayed${n.last_relayed ? renderNodeTimestampHtml(n.last_relayed) + ' ' + (n.relay_active ? ' actively relaying' : ' alive (idle)') : 'never observed as relay hop alive (idle)'}${(n.relay_count_1h != null || n.relay_count_24h != null) ? ` (${n.relay_count_1h || 0} relays/hr, ${n.relay_count_24h || 0} relays/24h)` : ''}
Last Heard (advert)${renderNodeTimestampHtml(lastHeard)}
Last Relayed${renderRelayActivity(n)}