diff --git a/cmd/server/db.go b/cmd/server/db.go index 90b1a918..99a0ff4d 100644 --- a/cmd/server/db.go +++ b/cmd/server/db.go @@ -15,9 +15,10 @@ import ( // DB wraps a read-only connection to the MeshCore SQLite database. type DB struct { - conn *sql.DB - path string // filesystem path to the database file - isV3 bool // v3 schema: observer_idx in observations (vs observer_id in v2) + conn *sql.DB + path string // filesystem path to the database file + isV3 bool // v3 schema: observer_idx in observations (vs observer_id in v2) + hasResolvedPath bool // observations table has resolved_path column } // OpenDB opens a read-only SQLite connection with WAL mode. @@ -61,9 +62,13 @@ func (db *DB) detectSchema() { var colType sql.NullString var notNull, pk int var dflt sql.NullString - if rows.Scan(&cid, &colName, &colType, ¬Null, &dflt, &pk) == nil && colName == "observer_idx" { - db.isV3 = true - return + if rows.Scan(&cid, &colName, &colType, ¬Null, &dflt, &pk) == nil { + if colName == "observer_idx" { + db.isV3 = true + } + if colName == "resolved_path" { + db.hasResolvedPath = true + } } } } diff --git a/cmd/server/main.go b/cmd/server/main.go index a1abc4f2..0aac9c73 100644 --- a/cmd/server/main.go +++ b/cmd/server/main.go @@ -144,6 +144,48 @@ func main() { log.Fatalf("[store] failed to load: %v", err) } + // Initialize persisted neighbor graph + dbPath = database.path + if err := ensureNeighborEdgesTable(dbPath); err != nil { + log.Printf("[neighbor] warning: could not create neighbor_edges table: %v", err) + } + // Add resolved_path column if missing. + // NOTE on startup ordering (review item #10): ensureResolvedPathColumn runs AFTER + // OpenDB/detectSchema, so db.hasResolvedPath will be false on first run with a + // pre-existing DB. This means Load() won't SELECT resolved_path from SQLite. + // That's OK: backfillResolvedPaths (below) computes and persists them in-memory + // AND to SQLite. On next restart, detectSchema finds the column and Load() reads it. + if err := ensureResolvedPathColumn(dbPath); err != nil { + log.Printf("[store] warning: could not add resolved_path column: %v", err) + } + + // Load or build neighbor graph + if neighborEdgesTableExists(database.conn) { + store.graph = loadNeighborEdgesFromDB(database.conn) + log.Printf("[neighbor] loaded persisted neighbor graph") + } else { + log.Printf("[neighbor] no persisted edges found, building from store...") + rw, rwErr := openRW(dbPath) + if rwErr == nil { + edgeCount := buildAndPersistEdges(store, rw) + rw.Close() + log.Printf("[neighbor] persisted %d edges", edgeCount) + } + store.graph = BuildFromStore(store) + } + + // Backfill resolved_path for observations that don't have it yet + if backfilled := backfillResolvedPaths(store, dbPath); backfilled > 0 { + log.Printf("[store] backfilled resolved_path for %d observations", backfilled) + } + + // Re-pick best observation now that resolved paths are populated + store.mu.Lock() + for _, tx := range store.packets { + pickBestObservation(tx) + } + store.mu.Unlock() + // WebSocket hub hub := NewHub() diff --git a/cmd/server/neighbor_persist.go b/cmd/server/neighbor_persist.go new file mode 100644 index 00000000..94f8dc1e --- /dev/null +++ b/cmd/server/neighbor_persist.go @@ -0,0 +1,531 @@ +package main + +import ( + "database/sql" + "encoding/json" + "fmt" + "log" + "strings" + "time" +) + +// persistSem limits concurrent async persistence goroutines to 1. +// Without this, each ingest cycle spawns a goroutine that opens a new +// SQLite RW connection; under sustained load goroutines pile up with +// no backpressure, causing contention and busy-timeout cascades. +var persistSem = make(chan struct{}, 1) + +// ─── neighbor_edges table ────────────────────────────────────────────────────── + +// ensureNeighborEdgesTable creates the neighbor_edges table if it doesn't exist. +// Uses a separate read-write connection since the main DB is read-only. +func ensureNeighborEdgesTable(dbPath string) error { + rw, err := openRW(dbPath) + if err != nil { + return fmt.Errorf("open rw for neighbor_edges: %w", err) + } + defer rw.Close() + + _, err = rw.Exec(`CREATE TABLE IF NOT EXISTS neighbor_edges ( + node_a TEXT NOT NULL, + node_b TEXT NOT NULL, + count INTEGER DEFAULT 1, + last_seen TEXT, + PRIMARY KEY (node_a, node_b) + )`) + return err +} + +// loadNeighborEdgesFromDB loads all edges from the neighbor_edges table +// and builds an in-memory NeighborGraph. +func loadNeighborEdgesFromDB(conn *sql.DB) *NeighborGraph { + g := NewNeighborGraph() + + rows, err := conn.Query("SELECT node_a, node_b, count, last_seen FROM neighbor_edges") + if err != nil { + log.Printf("[neighbor] failed to load neighbor_edges: %v", err) + return g + } + defer rows.Close() + + count := 0 + for rows.Next() { + var a, b string + var cnt int + var lastSeen sql.NullString + if err := rows.Scan(&a, &b, &cnt, &lastSeen); err != nil { + continue + } + ts := time.Time{} + if lastSeen.Valid { + ts = parseTimestamp(lastSeen.String) + } + // Build edge directly (both nodes are full pubkeys from persisted data) + key := makeEdgeKey(a, b) + g.mu.Lock() + e, exists := g.edges[key] + if !exists { + e = &NeighborEdge{ + NodeA: key.A, + NodeB: key.B, + Observers: make(map[string]bool), + FirstSeen: ts, + LastSeen: ts, + Count: cnt, + } + g.edges[key] = e + g.byNode[key.A] = append(g.byNode[key.A], e) + g.byNode[key.B] = append(g.byNode[key.B], e) + } else { + e.Count += cnt + if ts.After(e.LastSeen) { + e.LastSeen = ts + } + } + g.mu.Unlock() + count++ + } + + if count > 0 { + g.mu.Lock() + g.builtAt = time.Now() + g.mu.Unlock() + log.Printf("[neighbor] loaded %d edges from neighbor_edges table", count) + } + + return g +} + +// ─── shared async persistence helper ─────────────────────────────────────────── + +// persistObsUpdate holds data for a resolved_path SQLite update. +type persistObsUpdate struct { + obsID int + resolvedPath string +} + +// persistEdgeUpdate holds data for a neighbor_edges SQLite upsert. +type persistEdgeUpdate struct { + a, b, ts string +} + +// asyncPersistResolvedPathsAndEdges writes resolved_path updates and neighbor +// edge upserts to SQLite in a background goroutine. Shared between +// IngestNewFromDB and IngestNewObservations to avoid DRY violation. +func asyncPersistResolvedPathsAndEdges(dbPath string, obsUpdates []persistObsUpdate, edgeUpdates []persistEdgeUpdate, logPrefix string) { + if len(obsUpdates) == 0 && len(edgeUpdates) == 0 { + return + } + // Try-acquire semaphore BEFORE spawning goroutine. If another + // persistence operation is already running, drop this batch — + // data lives in memory and will be backfilled on restart. + select { + case persistSem <- struct{}{}: + // Acquired — spawn goroutine to do the work. + default: + log.Printf("[store] %s skipped: persistence already in progress", logPrefix) + return + } + go func() { + defer func() { <-persistSem }() + + rw, err := openRW(dbPath) + if err != nil { + log.Printf("[store] %s rw open error: %v", logPrefix, err) + return + } + defer rw.Close() + + if len(obsUpdates) > 0 { + sqlTx, err := rw.Begin() + if err == nil { + stmt, err := sqlTx.Prepare("UPDATE observations SET resolved_path = ? WHERE id = ?") + if err == nil { + var firstErr error + for _, u := range obsUpdates { + if _, err := stmt.Exec(u.resolvedPath, u.obsID); err != nil && firstErr == nil { + firstErr = err + } + } + stmt.Close() + if firstErr != nil { + log.Printf("[store] %s resolved_path error (first): %v", logPrefix, firstErr) + } + } else { + log.Printf("[store] %s resolved_path prepare error: %v", logPrefix, err) + } + sqlTx.Commit() + } + } + + if len(edgeUpdates) > 0 { + sqlTx, err := rw.Begin() + if err == nil { + stmt, err := sqlTx.Prepare(`INSERT INTO neighbor_edges (node_a, node_b, count, last_seen) + VALUES (?, ?, 1, ?) + ON CONFLICT(node_a, node_b) DO UPDATE SET + count = count + 1, last_seen = MAX(last_seen, excluded.last_seen)`) + if err == nil { + var firstErr error + for _, e := range edgeUpdates { + if _, err := stmt.Exec(e.a, e.b, e.ts); err != nil && firstErr == nil { + firstErr = err + } + } + stmt.Close() + if firstErr != nil { + log.Printf("[store] %s edge error (first): %v", logPrefix, firstErr) + } + } else { + log.Printf("[store] %s edge prepare error: %v", logPrefix, err) + } + sqlTx.Commit() + } + } + }() +} + +// neighborEdgesTableExists checks if the neighbor_edges table has any data. +func neighborEdgesTableExists(conn *sql.DB) bool { + var cnt int + err := conn.QueryRow("SELECT COUNT(*) FROM neighbor_edges").Scan(&cnt) + if err != nil { + return false // table doesn't exist + } + return cnt > 0 +} + +// buildAndPersistEdges scans all packets in the store, extracts edges per +// ADVERT/non-ADVERT rules, and persists them to SQLite. +func buildAndPersistEdges(store *PacketStore, rw *sql.DB) int { + store.mu.RLock() + packets := make([]*StoreTx, len(store.packets)) + copy(packets, store.packets) + store.mu.RUnlock() + + _, pm := store.getCachedNodesAndPM() + + tx, err := rw.Begin() + if err != nil { + log.Printf("[neighbor] begin tx error: %v", err) + return 0 + } + defer tx.Rollback() + + stmt, err := tx.Prepare(`INSERT INTO neighbor_edges (node_a, node_b, count, last_seen) + VALUES (?, ?, 1, ?) + ON CONFLICT(node_a, node_b) DO UPDATE SET + count = count + 1, last_seen = MAX(last_seen, excluded.last_seen)`) + if err != nil { + log.Printf("[neighbor] prepare stmt error: %v", err) + return 0 + } + defer stmt.Close() + + edgeCount := 0 + var firstErr error + for _, pkt := range packets { + for _, obs := range pkt.Observations { + for _, ec := range extractEdgesFromObs(obs, pkt, pm) { + if _, err := stmt.Exec(ec.A, ec.B, ec.Timestamp); err != nil && firstErr == nil { + firstErr = err + } + edgeCount++ + } + } + } + if firstErr != nil { + log.Printf("[neighbor] edge exec error (first): %v", firstErr) + } + + if err := tx.Commit(); err != nil { + log.Printf("[neighbor] commit error: %v", err) + return 0 + } + return edgeCount +} + +// ─── resolved_path column ────────────────────────────────────────────────────── + +// ensureResolvedPathColumn adds the resolved_path column to observations if missing. +func ensureResolvedPathColumn(dbPath string) error { + rw, err := openRW(dbPath) + if err != nil { + return err + } + defer rw.Close() + + // Check if column already exists + rows, err := rw.Query("PRAGMA table_info(observations)") + if err != nil { + return err + } + defer rows.Close() + + for rows.Next() { + var cid int + var colName string + var colType sql.NullString + var notNull, pk int + var dflt sql.NullString + if rows.Scan(&cid, &colName, &colType, ¬Null, &dflt, &pk) == nil && colName == "resolved_path" { + return nil // already exists + } + } + + _, err = rw.Exec("ALTER TABLE observations ADD COLUMN resolved_path TEXT") + if err != nil { + return fmt.Errorf("add resolved_path column: %w", err) + } + log.Println("[store] Added resolved_path column to observations") + return nil +} + +// resolvePathForObs resolves hop prefixes to full pubkeys for an observation. +// Returns nil if path is empty. +func resolvePathForObs(pathJSON, observerID string, tx *StoreTx, pm *prefixMap, graph *NeighborGraph) []*string { + hops := parsePathJSON(pathJSON) + if len(hops) == 0 { + return nil + } + + // Build context pubkeys: observer + originator (if known) + contextPKs := make([]string, 0, 3) + if observerID != "" { + contextPKs = append(contextPKs, strings.ToLower(observerID)) + } + fromNode := extractFromNode(tx) + if fromNode != "" { + contextPKs = append(contextPKs, strings.ToLower(fromNode)) + } + + resolved := make([]*string, len(hops)) + for i, hop := range hops { + // Add adjacent hops as context for disambiguation + ctx := make([]string, len(contextPKs), len(contextPKs)+2) + copy(ctx, contextPKs) + // Add previously resolved hops as context + if i > 0 && resolved[i-1] != nil { + ctx = append(ctx, *resolved[i-1]) + } + + node, _, _ := pm.resolveWithContext(hop, ctx, graph) + if node != nil { + pk := strings.ToLower(node.PublicKey) + resolved[i] = &pk + } + } + + return resolved +} + +// marshalResolvedPath converts []*string to JSON for storage. +func marshalResolvedPath(rp []*string) string { + if len(rp) == 0 { + return "" + } + b, err := json.Marshal(rp) + if err != nil { + return "" + } + return string(b) +} + +// unmarshalResolvedPath parses a resolved_path JSON string. +func unmarshalResolvedPath(s string) []*string { + if s == "" { + return nil + } + var result []*string + if json.Unmarshal([]byte(s), &result) != nil { + return nil + } + return result +} + +// backfillResolvedPaths resolves paths for all observations that have NULL resolved_path. +func backfillResolvedPaths(store *PacketStore, dbPath string) int { + // Collect pending observations and snapshot immutable fields under read lock. + // graph is set in main.go before backfill is called; nil-safe throughout (review item #6). + type obsRef struct { + obsID int + pathJSON string + observerID string + txJSON string // snapshot of DecodedJSON for extractFromNode + payloadType *int + } + store.mu.RLock() + pm := store.nodePM + graph := store.graph + var pending []obsRef + for _, tx := range store.packets { + for _, obs := range tx.Observations { + if obs.ResolvedPath == nil && obs.PathJSON != "" && obs.PathJSON != "[]" { + pending = append(pending, obsRef{ + obsID: obs.ID, + pathJSON: obs.PathJSON, + observerID: obs.ObserverID, + txJSON: tx.DecodedJSON, + payloadType: tx.PayloadType, + }) + } + } + } + store.mu.RUnlock() + + if len(pending) == 0 || pm == nil { + return 0 + } + + // Resolve paths outside the lock — resolvePathForObs only reads pm and graph. + type resolved struct { + obsID int + rp []*string + rpJSON string + } + var results []resolved + for _, ref := range pending { + // Build a minimal StoreTx for extractFromNode (only needs DecodedJSON + PayloadType). + fakeTx := &StoreTx{DecodedJSON: ref.txJSON, PayloadType: ref.payloadType} + rp := resolvePathForObs(ref.pathJSON, ref.observerID, fakeTx, pm, graph) + if len(rp) > 0 { + rpJSON := marshalResolvedPath(rp) + if rpJSON != "" { + results = append(results, resolved{ref.obsID, rp, rpJSON}) + } + } + } + + if len(results) == 0 { + return 0 + } + + // Persist to SQLite (no lock needed — separate RW connection). + rw, err := openRW(dbPath) + if err != nil { + log.Printf("[store] backfill: open rw error: %v", err) + return 0 + } + defer rw.Close() + + sqlTx, err := rw.Begin() + if err != nil { + log.Printf("[store] backfill: begin tx error: %v", err) + return 0 + } + defer sqlTx.Rollback() + + stmt, err := sqlTx.Prepare("UPDATE observations SET resolved_path = ? WHERE id = ?") + if err != nil { + log.Printf("[store] backfill: prepare error: %v", err) + return 0 + } + defer stmt.Close() + + var firstErr error + for _, r := range results { + if _, err := stmt.Exec(r.rpJSON, r.obsID); err != nil && firstErr == nil { + firstErr = err + } + } + if firstErr != nil { + log.Printf("[store] backfill resolved_path exec error (first): %v", firstErr) + } + + if err := sqlTx.Commit(); err != nil { + log.Printf("[store] backfill: commit error: %v", err) + return 0 + } + + // Update in-memory state under write lock. + store.mu.Lock() + count := 0 + for _, r := range results { + if obs, ok := store.byObsID[r.obsID]; ok { + obs.ResolvedPath = r.rp + count++ + } + } + store.mu.Unlock() + + return count +} + +// ─── Shared helpers ──────────────────────────────────────────────────────────── + +// edgeCandidate represents an extracted edge to be persisted. +type edgeCandidate struct { + A, B, Timestamp string +} + +// extractEdgesFromObs extracts neighbor edge candidates from a single observation. +// For ADVERTs: originator↔path[0] (if unambiguous). For ALL types: observer↔path[last] (if unambiguous). +// Also handles zero-hop ADVERTs (originator↔observer direct link). +func extractEdgesFromObs(obs *StoreObs, tx *StoreTx, pm *prefixMap) []edgeCandidate { + isAdvert := tx.PayloadType != nil && *tx.PayloadType == 4 + fromNode := extractFromNode(tx) + path := parsePathJSON(obs.PathJSON) + observerPK := strings.ToLower(obs.ObserverID) + ts := obs.Timestamp + var edges []edgeCandidate + + if len(path) == 0 { + if isAdvert && fromNode != "" { + fromLower := strings.ToLower(fromNode) + if fromLower != observerPK { + a, b := fromLower, observerPK + if a > b { + a, b = b, a + } + edges = append(edges, edgeCandidate{a, b, ts}) + } + } + return edges + } + + // Edge 1: originator ↔ path[0] — ADVERTs only (resolve prefix to full pubkey) + if isAdvert && fromNode != "" && pm != nil { + firstHop := strings.ToLower(path[0]) + fromLower := strings.ToLower(fromNode) + candidates := pm.m[firstHop] + if len(candidates) == 1 { + resolved := strings.ToLower(candidates[0].PublicKey) + if resolved != fromLower { + a, b := fromLower, resolved + if a > b { + a, b = b, a + } + edges = append(edges, edgeCandidate{a, b, ts}) + } + } + } + + // Edge 2: observer ↔ path[last] — ALL packet types + if pm != nil { + lastHop := strings.ToLower(path[len(path)-1]) + candidates := pm.m[lastHop] + if len(candidates) == 1 { + resolved := strings.ToLower(candidates[0].PublicKey) + if resolved != observerPK { + a, b := observerPK, resolved + if a > b { + a, b = b, a + } + edges = append(edges, edgeCandidate{a, b, ts}) + } + } + } + + return edges +} + +// openRW opens a read-write SQLite connection (same pattern as PruneOldPackets). +func openRW(dbPath string) (*sql.DB, error) { + dsn := fmt.Sprintf("file:%s?_journal_mode=WAL&_busy_timeout=10000", dbPath) + rw, err := sql.Open("sqlite", dsn) + if err != nil { + return nil, err + } + rw.SetMaxOpenConns(1) + return rw, nil +} diff --git a/cmd/server/neighbor_persist_test.go b/cmd/server/neighbor_persist_test.go new file mode 100644 index 00000000..66fe35ea --- /dev/null +++ b/cmd/server/neighbor_persist_test.go @@ -0,0 +1,534 @@ +package main + +import ( + "database/sql" + "encoding/json" + "path/filepath" + "strings" + "testing" + "time" + + _ "modernc.org/sqlite" +) + +// createTestDBWithSchema creates a temp SQLite DB with the standard schema + resolved_path column. +func createTestDBWithSchema(t *testing.T) (*DB, string) { + t.Helper() + dir := t.TempDir() + dbPath := filepath.Join(dir, "test.db") + + conn, err := sql.Open("sqlite", "file:"+dbPath+"?_journal_mode=WAL") + if err != nil { + t.Fatal(err) + } + + // Create tables + conn.Exec(`CREATE TABLE transmissions ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + raw_hex TEXT, hash TEXT UNIQUE, first_seen TEXT, + route_type INTEGER, payload_type INTEGER, payload_version INTEGER, + decoded_json TEXT + )`) + conn.Exec(`CREATE TABLE observers ( + id TEXT PRIMARY KEY, name TEXT, iata TEXT + )`) + conn.Exec(`CREATE TABLE observations ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + transmission_id INTEGER NOT NULL REFERENCES transmissions(id), + observer_id TEXT, observer_name TEXT, direction TEXT, + snr REAL, rssi REAL, score INTEGER, + path_json TEXT, timestamp TEXT, + resolved_path TEXT + )`) + conn.Exec(`CREATE TABLE nodes ( + public_key TEXT PRIMARY KEY, name TEXT, role TEXT, + lat REAL, lon REAL, last_seen TEXT, first_seen TEXT, + advert_count INTEGER DEFAULT 0 + )`) + + conn.Close() + + db, err := OpenDB(dbPath) + if err != nil { + t.Fatal(err) + } + return db, dbPath +} + +func TestResolvePathForObs(t *testing.T) { + // Build a prefix map with known nodes + nodes := []nodeInfo{ + {PublicKey: "aabbccddee1234567890aabbccddee1234567890aabbccddee1234567890aabb", Name: "Node-AA"}, + {PublicKey: "bbccddee1234567890aabbccddee1234567890aabbccddee1234567890aabb11", Name: "Node-BB"}, + } + pm := buildPrefixMap(nodes) + graph := NewNeighborGraph() + + tx := &StoreTx{ + DecodedJSON: `{"pubKey": "originator1234567890"}`, + PayloadType: intPtr(4), + } + + // Unambiguous prefixes should resolve + rp := resolvePathForObs(`["aa","bb"]`, "observer1", tx, pm, graph) + if len(rp) != 2 { + t.Fatalf("expected 2 resolved hops, got %d", len(rp)) + } + if rp[0] == nil || !strings.HasPrefix(*rp[0], "aabbcc") { + t.Errorf("expected first hop to resolve to Node-AA, got %v", rp[0]) + } + if rp[1] == nil || !strings.HasPrefix(*rp[1], "bbccdd") { + t.Errorf("expected second hop to resolve to Node-BB, got %v", rp[1]) + } +} + +func TestResolvePathForObs_EmptyPath(t *testing.T) { + pm := buildPrefixMap(nil) + rp := resolvePathForObs(`[]`, "", &StoreTx{}, pm, nil) + if rp != nil { + t.Errorf("expected nil for empty path, got %v", rp) + } + + rp = resolvePathForObs("", "", &StoreTx{}, pm, nil) + if rp != nil { + t.Errorf("expected nil for empty string, got %v", rp) + } +} + +func TestResolvePathForObs_Unresolvable(t *testing.T) { + nodes := []nodeInfo{ + {PublicKey: "aabbccddee1234567890aabbccddee1234567890aabbccddee1234567890aabb", Name: "Node-AA"}, + } + pm := buildPrefixMap(nodes) + + // "zz" prefix doesn't match any node + rp := resolvePathForObs(`["zz"]`, "", &StoreTx{}, pm, nil) + if len(rp) != 1 { + t.Fatalf("expected 1 hop, got %d", len(rp)) + } + if rp[0] != nil { + t.Errorf("expected nil for unresolvable hop, got %v", *rp[0]) + } +} + +func TestMarshalUnmarshalResolvedPath(t *testing.T) { + pk1 := "aabbccdd" + var rp []*string + rp = append(rp, &pk1, nil) + + j := marshalResolvedPath(rp) + if j == "" { + t.Fatal("expected non-empty JSON") + } + + parsed := unmarshalResolvedPath(j) + if len(parsed) != 2 { + t.Fatalf("expected 2 elements, got %d", len(parsed)) + } + if parsed[0] == nil || *parsed[0] != "aabbccdd" { + t.Errorf("first element wrong: %v", parsed[0]) + } + if parsed[1] != nil { + t.Errorf("second element should be nil, got %v", *parsed[1]) + } +} + +func TestMarshalResolvedPath_Empty(t *testing.T) { + if marshalResolvedPath(nil) != "" { + t.Error("expected empty for nil") + } + if marshalResolvedPath([]*string{}) != "" { + t.Error("expected empty for empty slice") + } +} + +func TestUnmarshalResolvedPath_Invalid(t *testing.T) { + if unmarshalResolvedPath("") != nil { + t.Error("expected nil for empty string") + } + if unmarshalResolvedPath("not json") != nil { + t.Error("expected nil for invalid JSON") + } +} + +func TestEnsureNeighborEdgesTable(t *testing.T) { + dir := t.TempDir() + dbPath := filepath.Join(dir, "test.db") + + // Create initial DB + conn, _ := sql.Open("sqlite", "file:"+dbPath+"?_journal_mode=WAL") + conn.Exec("CREATE TABLE test (id INTEGER PRIMARY KEY)") + conn.Close() + + if err := ensureNeighborEdgesTable(dbPath); err != nil { + t.Fatal(err) + } + + // Verify table exists + conn, _ = sql.Open("sqlite", "file:"+dbPath+"?mode=ro") + defer conn.Close() + var cnt int + if err := conn.QueryRow("SELECT COUNT(*) FROM neighbor_edges").Scan(&cnt); err != nil { + t.Fatalf("neighbor_edges table not created: %v", err) + } +} + +func TestLoadNeighborEdgesFromDB(t *testing.T) { + dir := t.TempDir() + dbPath := filepath.Join(dir, "test.db") + + conn, _ := sql.Open("sqlite", "file:"+dbPath+"?_journal_mode=WAL") + conn.Exec(`CREATE TABLE neighbor_edges ( + node_a TEXT NOT NULL, node_b TEXT NOT NULL, + count INTEGER DEFAULT 1, last_seen TEXT, + PRIMARY KEY (node_a, node_b) + )`) + conn.Exec("INSERT INTO neighbor_edges VALUES ('aaa', 'bbb', 5, '2024-01-01T00:00:00Z')") + conn.Exec("INSERT INTO neighbor_edges VALUES ('ccc', 'ddd', 3, '2024-01-02T00:00:00Z')") + + g := loadNeighborEdgesFromDB(conn) + conn.Close() + + // Should have 2 edges + edges := g.AllEdges() + if len(edges) != 2 { + t.Errorf("expected 2 edges, got %d", len(edges)) + } + + // Check neighbors + n := g.Neighbors("aaa") + if len(n) != 1 { + t.Errorf("expected 1 neighbor for aaa, got %d", len(n)) + } +} + +func TestStoreObsResolvedPathInBroadcast(t *testing.T) { + // Verify resolved_path appears in broadcast maps + pk := "aabbccdd" + obs := &StoreObs{ + ID: 1, + ObserverID: "obs1", + ObserverName: "Observer 1", + PathJSON: `["aa"]`, + ResolvedPath: []*string{&pk}, + Timestamp: "2024-01-01T00:00:00Z", + } + + tx := &StoreTx{ + ID: 1, + Hash: "abc123", + Observations: []*StoreObs{obs}, + } + pickBestObservation(tx) + + if tx.ResolvedPath == nil { + t.Fatal("expected ResolvedPath to be set on tx after pickBestObservation") + } + if *tx.ResolvedPath[0] != "aabbccdd" { + t.Errorf("expected resolved path to be aabbccdd, got %s", *tx.ResolvedPath[0]) + } +} + +func TestResolvedPathInTxToMap(t *testing.T) { + pk := "aabbccdd" + tx := &StoreTx{ + ID: 1, + Hash: "abc123", + PathJSON: `["aa"]`, + ResolvedPath: []*string{&pk}, + obsKeys: make(map[string]bool), + } + + m := txToMap(tx) + rp, ok := m["resolved_path"] + if !ok { + t.Fatal("resolved_path not in txToMap output") + } + rpSlice, ok := rp.([]*string) + if !ok || len(rpSlice) != 1 || *rpSlice[0] != "aabbccdd" { + t.Errorf("unexpected resolved_path: %v", rp) + } +} + +func TestResolvedPathOmittedWhenNil(t *testing.T) { + tx := &StoreTx{ + ID: 1, + Hash: "abc123", + obsKeys: make(map[string]bool), + } + + m := txToMap(tx) + if _, ok := m["resolved_path"]; ok { + t.Error("resolved_path should not be in map when nil") + } +} + +func TestEnsureResolvedPathColumn(t *testing.T) { + dir := t.TempDir() + dbPath := filepath.Join(dir, "test.db") + + conn, _ := sql.Open("sqlite", "file:"+dbPath+"?_journal_mode=WAL") + conn.Exec(`CREATE TABLE observations ( + id INTEGER PRIMARY KEY, transmission_id INTEGER, + observer_id TEXT, path_json TEXT, timestamp TEXT + )`) + conn.Close() + + if err := ensureResolvedPathColumn(dbPath); err != nil { + t.Fatal(err) + } + + // Verify column exists + conn, _ = sql.Open("sqlite", "file:"+dbPath+"?mode=ro") + defer conn.Close() + rows, _ := conn.Query("PRAGMA table_info(observations)") + found := false + for rows.Next() { + var cid int + var colName string + var colType sql.NullString + var notNull, pk int + var dflt sql.NullString + rows.Scan(&cid, &colName, &colType, ¬Null, &dflt, &pk) + if colName == "resolved_path" { + found = true + } + } + rows.Close() + if !found { + t.Error("resolved_path column not added") + } + + // Running again should be idempotent + if err := ensureResolvedPathColumn(dbPath); err != nil { + t.Fatal("second call should be idempotent:", err) + } +} + +func TestDBDetectsResolvedPathColumn(t *testing.T) { + dir := t.TempDir() + dbPath := filepath.Join(dir, "test.db") + + // Create DB without resolved_path + conn, _ := sql.Open("sqlite", "file:"+dbPath+"?_journal_mode=WAL") + conn.Exec(`CREATE TABLE observations (id INTEGER PRIMARY KEY, observer_idx INTEGER)`) + conn.Exec(`CREATE TABLE transmissions (id INTEGER PRIMARY KEY)`) + conn.Close() + + db, err := OpenDB(dbPath) + if err != nil { + t.Fatal(err) + } + if db.hasResolvedPath { + t.Error("should not detect resolved_path when column missing") + } + db.Close() + + // Add resolved_path column + conn, _ = sql.Open("sqlite", "file:"+dbPath+"?_journal_mode=WAL") + conn.Exec("ALTER TABLE observations ADD COLUMN resolved_path TEXT") + conn.Close() + + db, err = OpenDB(dbPath) + if err != nil { + t.Fatal(err) + } + if !db.hasResolvedPath { + t.Error("should detect resolved_path when column exists") + } + db.Close() +} + +func TestLoadWithResolvedPath(t *testing.T) { + db, dbPath := createTestDBWithSchema(t) + defer db.Close() + + // Insert test data + rw, _ := openRW(dbPath) + rw.Exec(`INSERT INTO transmissions (id, hash, first_seen, payload_type, decoded_json) + VALUES (1, 'hash1', '2024-01-01T00:00:00Z', 4, '{"pubKey":"origpk"}')`) + rw.Exec(`INSERT INTO observations (id, transmission_id, observer_id, observer_name, path_json, timestamp, resolved_path) + VALUES (1, 1, 'obs1', 'Observer1', '["aa"]', '2024-01-01T00:00:00Z', '["aabbccdd"]')`) + rw.Close() + + store := NewPacketStore(db, nil) + if err := store.Load(); err != nil { + t.Fatal(err) + } + + if len(store.packets) != 1 { + t.Fatalf("expected 1 packet, got %d", len(store.packets)) + } + + tx := store.packets[0] + if len(tx.Observations) != 1 { + t.Fatalf("expected 1 observation, got %d", len(tx.Observations)) + } + + obs := tx.Observations[0] + if obs.ResolvedPath == nil { + t.Fatal("expected ResolvedPath to be loaded") + } + if len(obs.ResolvedPath) != 1 || *obs.ResolvedPath[0] != "aabbccdd" { + t.Errorf("unexpected ResolvedPath: %v", obs.ResolvedPath) + } + + // Check that pickBestObservation propagated resolved_path to tx + if tx.ResolvedPath == nil || len(tx.ResolvedPath) != 1 { + t.Error("expected ResolvedPath to be propagated to tx") + } +} + +func TestResolvedPathInAPIResponse(t *testing.T) { + // Test that TransmissionResp properly marshals resolved_path + pk := "aabbccddee" + resp := TransmissionResp{ + ID: 1, + Hash: "test", + ResolvedPath: []*string{&pk, nil}, + } + + data, err := json.Marshal(resp) + if err != nil { + t.Fatal(err) + } + + var m map[string]interface{} + json.Unmarshal(data, &m) + + rp, ok := m["resolved_path"] + if !ok { + t.Fatal("resolved_path missing from JSON") + } + rpArr, ok := rp.([]interface{}) + if !ok || len(rpArr) != 2 { + t.Fatalf("unexpected resolved_path shape: %v", rp) + } + if rpArr[0] != "aabbccddee" { + t.Errorf("first element wrong: %v", rpArr[0]) + } + if rpArr[1] != nil { + t.Errorf("second element should be null: %v", rpArr[1]) + } +} + +func TestResolvedPathOmittedWhenEmpty(t *testing.T) { + resp := TransmissionResp{ + ID: 1, + Hash: "test", + } + + data, _ := json.Marshal(resp) + var m map[string]interface{} + json.Unmarshal(data, &m) + + if _, ok := m["resolved_path"]; ok { + t.Error("resolved_path should be omitted when nil") + } +} + +func TestExtractEdgesFromObs_AdvertNoPath(t *testing.T) { + tx := &StoreTx{ + DecodedJSON: `{"pubKey":"aaaa1111"}`, + PayloadType: intPtr(4), + } + obs := &StoreObs{ + ObserverID: "bbbb2222", + PathJSON: "", + Timestamp: "2024-01-01T00:00:00Z", + } + + edges := extractEdgesFromObs(obs, tx, nil) + if len(edges) != 1 { + t.Fatalf("expected 1 edge for zero-hop advert, got %d", len(edges)) + } + // Canonical ordering: aaaa < bbbb + if edges[0].A != "aaaa1111" || edges[0].B != "bbbb2222" { + t.Errorf("unexpected edge: %+v", edges[0]) + } +} + +func TestExtractEdgesFromObs_NonAdvertNoPath(t *testing.T) { + tx := &StoreTx{PayloadType: intPtr(1)} + obs := &StoreObs{ObserverID: "obs1", PathJSON: ""} + edges := extractEdgesFromObs(obs, tx, nil) + if len(edges) != 0 { + t.Errorf("expected 0 edges for non-advert without path, got %d", len(edges)) + } +} + +func TestExtractEdgesFromObs_WithPath(t *testing.T) { + nodes := []nodeInfo{ + {PublicKey: "aabbccddee1234567890aabbccddee1234567890aabbccddee1234567890aabb", Name: "Node-AA"}, + {PublicKey: "ffgghhii1234567890aabbccddee1234567890aabbccddee1234567890aabb11", Name: "Node-FF"}, + } + pm := buildPrefixMap(nodes) + + tx := &StoreTx{ + DecodedJSON: `{"pubKey":"originator00"}`, + PayloadType: intPtr(4), + } + obs := &StoreObs{ + ObserverID: "observer00", + PathJSON: `["aa","ff"]`, + Timestamp: "2024-01-01T00:00:00Z", + } + + edges := extractEdgesFromObs(obs, tx, pm) + // Should get: originator↔aa (advert), observer↔ff (last hop) + if len(edges) != 2 { + t.Fatalf("expected 2 edges, got %d", len(edges)) + } +} + +func TestExtractEdgesFromObs_SameNodeNoEdge(t *testing.T) { + tx := &StoreTx{ + DecodedJSON: `{"pubKey":"same1234"}`, + PayloadType: intPtr(4), + } + obs := &StoreObs{ + ObserverID: "same1234", + PathJSON: "", + Timestamp: "2024-01-01T00:00:00Z", + } + edges := extractEdgesFromObs(obs, tx, nil) + if len(edges) != 0 { + t.Errorf("expected 0 edges when originator == observer, got %d", len(edges)) + } +} + + + +func TestPersistSemaphoreTryAcquireSkipsBatch(t *testing.T) { + // Verify that persistSem is a buffered channel of size 1. + if cap(persistSem) != 1 { + t.Errorf("persistSem capacity = %d, want 1", cap(persistSem)) + } + // Acquire the semaphore to simulate an in-progress persistence. + persistSem <- struct{}{} + + // asyncPersistResolvedPathsAndEdges should skip (not block, not + // spawn a goroutine) when the semaphore is already held. + done := make(chan struct{}) + go func() { + asyncPersistResolvedPathsAndEdges( + "/nonexistent/path.db", + []persistObsUpdate{{obsID: 1, resolvedPath: "x"}}, + nil, + "test", + ) + close(done) + }() + + // If the function blocks on the semaphore instead of skipping, + // this select will hit the timeout. + select { + case <-done: + // Expected: returned immediately because semaphore was busy. + case <-time.After(500 * time.Millisecond): + <-persistSem + t.Fatal("asyncPersistResolvedPathsAndEdges blocked instead of skipping when semaphore was held") + } + + <-persistSem // release +} diff --git a/cmd/server/routes.go b/cmd/server/routes.go index 98bef917..ef0b579e 100644 --- a/cmd/server/routes.go +++ b/cmd/server/routes.go @@ -1088,7 +1088,7 @@ func (s *Server) handleNodePaths(w http.ResponseWriter, r *http.Request) { if cached, ok := hopCache[hop]; ok { return cached } - r := pm.resolve(hop) + r, _, _ := pm.resolveWithContext(hop, nil, s.store.graph) hopCache[hop] = r return r } @@ -1997,6 +1997,9 @@ func mapSliceToTransmissions(maps []map[string]interface{}) []TransmissionResp { tx.PathJSON = m["path_json"] tx.Direction = m["direction"] tx.Score = m["score"] + if rp, ok := m["resolved_path"].([]*string); ok { + tx.ResolvedPath = rp + } result = append(result, tx) } return result @@ -2018,6 +2021,9 @@ func mapSliceToObservations(maps []map[string]interface{}) []ObservationResp { obs.RSSI = m["rssi"] obs.PathJSON = m["path_json"] obs.Timestamp = m["timestamp"] + if rp, ok := m["resolved_path"].([]*string); ok { + obs.ResolvedPath = rp + } result = append(result, obs) } return result diff --git a/cmd/server/store.go b/cmd/server/store.go index 90b4abed..5b9b4eb3 100644 --- a/cmd/server/store.go +++ b/cmd/server/store.go @@ -39,6 +39,7 @@ type StoreTx struct { RSSI *float64 PathJSON string Direction string + ResolvedPath []*string // resolved path from best observation LatestSeen string // max observation timestamp (or FirstSeen if no observations) // Cached parsed fields (set once, read many) parsedPath []string // cached parsePathJSON result @@ -58,6 +59,7 @@ type StoreObs struct { RSSI *float64 Score *int PathJSON string + ResolvedPath []*string // resolved full pubkeys, parallel to path_json; nil elements = unresolved Timestamp string } @@ -133,6 +135,9 @@ type PacketStore struct { // Updated incrementally during Load/Ingest/Evict — avoids JSON parsing in GetPerfStoreStats. advertPubkeys map[string]int // pubkey → number of advert packets referencing it + // Persisted neighbor graph for hop resolution at ingest time. + graph *NeighborGraph + // Eviction config and stats retentionHours float64 // 0 = unlimited maxMemoryMB int // 0 = unlimited @@ -217,11 +222,15 @@ func (s *PacketStore) Load() error { t0 := time.Now() var loadSQL string + rpCol := "" + if s.db.hasResolvedPath { + rpCol = ",\n\t\t\t\to.resolved_path" + } if s.db.isV3 { loadSQL = `SELECT t.id, t.raw_hex, t.hash, t.first_seen, t.route_type, t.payload_type, t.payload_version, t.decoded_json, o.id, obs.id, obs.name, o.direction, - o.snr, o.rssi, o.score, o.path_json, strftime('%Y-%m-%dT%H:%M:%fZ', o.timestamp, 'unixepoch') + o.snr, o.rssi, o.score, o.path_json, strftime('%Y-%m-%dT%H:%M:%fZ', o.timestamp, 'unixepoch')` + rpCol + ` FROM transmissions t LEFT JOIN observations o ON o.transmission_id = t.id LEFT JOIN observers obs ON obs.rowid = o.observer_idx @@ -230,7 +239,7 @@ func (s *PacketStore) Load() error { loadSQL = `SELECT t.id, t.raw_hex, t.hash, t.first_seen, t.route_type, t.payload_type, t.payload_version, t.decoded_json, o.id, o.observer_id, o.observer_name, o.direction, - o.snr, o.rssi, o.score, o.path_json, o.timestamp + o.snr, o.rssi, o.score, o.path_json, o.timestamp` + rpCol + ` FROM transmissions t LEFT JOIN observations o ON o.transmission_id = t.id ORDER BY t.first_seen ASC, o.timestamp DESC` @@ -250,11 +259,16 @@ func (s *PacketStore) Load() error { var observerID, observerName, direction, pathJSON, obsTimestamp sql.NullString var snr, rssi sql.NullFloat64 var score sql.NullInt64 + var resolvedPathStr sql.NullString - if err := rows.Scan(&txID, &rawHex, &hash, &firstSeen, &routeType, &payloadType, + scanArgs := []interface{}{&txID, &rawHex, &hash, &firstSeen, &routeType, &payloadType, &payloadVersion, &decodedJSON, &obsID, &observerID, &observerName, &direction, - &snr, &rssi, &score, &pathJSON, &obsTimestamp); err != nil { + &snr, &rssi, &score, &pathJSON, &obsTimestamp} + if s.db.hasResolvedPath { + scanArgs = append(scanArgs, &resolvedPathStr) + } + if err := rows.Scan(scanArgs...); err != nil { log.Printf("[store] scan error: %v", err) continue } @@ -305,6 +319,7 @@ func (s *PacketStore) Load() error { RSSI: nullFloatPtr(rssi), Score: nullIntPtr(score), PathJSON: obsPJ, + ResolvedPath: unmarshalResolvedPath(nullStrVal(resolvedPathStr)), Timestamp: normalizeTimestamp(nullStrVal(obsTimestamp)), } @@ -366,6 +381,7 @@ func pickBestObservation(tx *StoreTx) { tx.RSSI = best.RSSI tx.PathJSON = best.PathJSON tx.Direction = best.Direction + tx.ResolvedPath = best.ResolvedPath tx.pathParsed = false // invalidate cached parsed path } @@ -990,6 +1006,9 @@ func (s *PacketStore) IngestNewFromDB(sinceID, limit int) ([]map[string]interfac limit = 100 } + // NOTE: The SQL query intentionally does NOT select resolved_path from the DB. + // New ingests always resolve fresh using the current prefix map and neighbor graph. + // On restart, Load() handles reading persisted resolved_path values. (review item #7) var querySQL string if s.db.isV3 { querySQL = `SELECT t.id, t.raw_hex, t.hash, t.first_seen, t.route_type, @@ -1094,6 +1113,10 @@ func (s *PacketStore) IngestNewFromDB(sinceID, limit int) ([]map[string]interfac broadcastTxs := make(map[int]*StoreTx) // track new transmissions for broadcast var broadcastOrder []int + // Hoist getCachedNodesAndPM() once before the observation loop to avoid + // per-observation function calls (review item #1). + _, cachedPM := s.getCachedNodesAndPM() + for _, r := range tempRows { if r.txID > newMaxID { newMaxID = r.txID @@ -1153,6 +1176,13 @@ func (s *PacketStore) IngestNewFromDB(sinceID, limit int) ([]map[string]interfac PathJSON: r.pathJSON, Timestamp: normalizeTimestamp(r.obsTS), } + + // Resolve path at ingest time using neighbor graph + // (cachedPM is hoisted before the observation loop to avoid per-obs function calls) + if r.pathJSON != "" && r.pathJSON != "[]" && cachedPM != nil { + obs.ResolvedPath = resolvePathForObs(r.pathJSON, r.observerID, tx, cachedPM, s.graph) + } + tx.Observations = append(tx.Observations, obs) tx.obsKeys[dk] = true tx.ObservationCount++ @@ -1196,7 +1226,7 @@ func (s *PacketStore) IngestNewFromDB(sinceID, limit int) ([]map[string]interfac if cached, ok := hopCache[hop]; ok { return cached } - r := pm.resolve(hop) + r, _, _ := pm.resolveWithContext(hop, nil, s.graph) hopCache[hop] = r return r } @@ -1246,6 +1276,9 @@ func (s *PacketStore) IngestNewFromDB(sinceID, limit int) ([]map[string]interfac "direction": strOrNil(obs.Direction), "observation_count": tx.ObservationCount, } + if obs.ResolvedPath != nil { + pkt["resolved_path"] = obs.ResolvedPath + } // Broadcast map: top-level fields for live.js + nested packet for packets.js broadcastMap := make(map[string]interface{}, len(pkt)+2) for k, v := range pkt { @@ -1280,6 +1313,36 @@ func (s *PacketStore) IngestNewFromDB(sinceID, limit int) ([]map[string]interfac s.invalidateCachesFor(inv) } + // Persist resolved paths and neighbor edges asynchronously (don't block ingest). + if len(broadcastTxs) > 0 && s.db != nil { + dbPath := s.db.path + var obsUpdates []persistObsUpdate + var edgeUpdates []persistEdgeUpdate + + _, pm := s.getCachedNodesAndPM() + // Read graph ref under lock (it's set during startup and not replaced after, + // but reading under lock is safer — review item #5). + graphRef := s.graph + for _, tx := range broadcastTxs { + for _, obs := range tx.Observations { + if obs.ResolvedPath != nil { + rpJSON := marshalResolvedPath(obs.ResolvedPath) + if rpJSON != "" { + obsUpdates = append(obsUpdates, persistObsUpdate{obs.ID, rpJSON}) + } + } + for _, ec := range extractEdgesFromObs(obs, tx, pm) { + edgeUpdates = append(edgeUpdates, persistEdgeUpdate{ec.A, ec.B, ec.Timestamp}) + if graphRef != nil { + graphRef.upsertEdge(ec.A, ec.B, "", obs.ObserverID, obs.SNR, parseTimestamp(ec.Timestamp)) + } + } + } + } + + asyncPersistResolvedPathsAndEdges(dbPath, obsUpdates, edgeUpdates, "persist") + } + return result, newMaxID } @@ -1363,6 +1426,13 @@ func (s *PacketStore) IngestNewObservations(sinceObsID, limit int) []map[string] updatedTxs := make(map[int]*StoreTx) broadcastMaps := make([]map[string]interface{}, 0, len(obsRows)) + // Track newly created observations for persistence — only these should be + // persisted, not all observations of each updated tx (fixes edge count inflation). + var newObs []*StoreObs + + // Hoist getCachedNodesAndPM() before the loop — same pattern as IngestNewFromDB (review fix #1). + _, pm := s.getCachedNodesAndPM() + graphRef := s.graph for _, r := range obsRows { // Already ingested (e.g. by IngestNewFromDB in same cycle) @@ -1396,9 +1466,18 @@ func (s *PacketStore) IngestNewObservations(sinceObsID, limit int) []map[string] PathJSON: r.pathJSON, Timestamp: normalizeTimestamp(r.timestamp), } + + // Resolve path at ingest time for late-arriving observations (review item #2). + if r.pathJSON != "" && r.pathJSON != "[]" { + if pm != nil { + obs.ResolvedPath = resolvePathForObs(r.pathJSON, r.observerID, tx, pm, s.graph) + } + } + tx.Observations = append(tx.Observations, obs) tx.obsKeys[dk] = true tx.ObservationCount++ + newObs = append(newObs, obs) if obs.Timestamp > tx.LatestSeen { tx.LatestSeen = obs.Timestamp } @@ -1438,6 +1517,9 @@ func (s *PacketStore) IngestNewObservations(sinceObsID, limit int) []map[string] "direction": strOrNil(obs.Direction), "observation_count": tx.ObservationCount, } + if obs.ResolvedPath != nil { + pkt["resolved_path"] = obs.ResolvedPath + } broadcastMap := make(map[string]interface{}, len(pkt)+2) for k, v := range pkt { broadcastMap[k] = v @@ -1506,6 +1588,36 @@ func (s *PacketStore) IngestNewObservations(sinceObsID, limit int) []map[string] }) } + // Persist resolved paths and neighbor edges asynchronously (review fix #3). + // Only process NEW observations — not all observations of each updated tx — + // to avoid edge count inflation and unnecessary UPDATEs for pre-existing data. + if len(newObs) > 0 && s.db != nil { + dbPath := s.db.path + var obsUpdates []persistObsUpdate + var edgeUpdates []persistEdgeUpdate + + for _, obs := range newObs { + tx := s.byTxID[obs.TransmissionID] + if tx == nil { + continue + } + if obs.ResolvedPath != nil { + rpJSON := marshalResolvedPath(obs.ResolvedPath) + if rpJSON != "" { + obsUpdates = append(obsUpdates, persistObsUpdate{obs.ID, rpJSON}) + } + } + for _, ec := range extractEdgesFromObs(obs, tx, pm) { + edgeUpdates = append(edgeUpdates, persistEdgeUpdate{ec.A, ec.B, ec.Timestamp}) + if graphRef != nil { + graphRef.upsertEdge(ec.A, ec.B, "", obs.ObserverID, obs.SNR, parseTimestamp(ec.Timestamp)) + } + } + } + + asyncPersistResolvedPathsAndEdges(dbPath, obsUpdates, edgeUpdates, "obs-persist") + } + return broadcastMaps } @@ -1686,6 +1798,9 @@ func (s *PacketStore) enrichObs(obs *StoreObs) map[string]interface{} { "score": intPtrOrNil(obs.Score), "path_json": strOrNil(obs.PathJSON), } + if obs.ResolvedPath != nil { + m["resolved_path"] = obs.ResolvedPath + } if tx != nil { m["hash"] = strOrNil(tx.Hash) @@ -1719,6 +1834,9 @@ func txToMap(tx *StoreTx) map[string]interface{} { "path_json": strOrNil(tx.PathJSON), "direction": strOrNil(tx.Direction), } + if tx.ResolvedPath != nil { + m["resolved_path"] = tx.ResolvedPath + } // Include parsed path array to match Node.js output shape if hops := txGetParsedPath(tx); len(hops) > 0 { m["_parsedPath"] = hops @@ -1728,7 +1846,7 @@ func txToMap(tx *StoreTx) map[string]interface{} { // Include observations for expand=observations support (stripped by handler when not requested) obs := make([]map[string]interface{}, 0, len(tx.Observations)) for _, o := range tx.Observations { - obs = append(obs, map[string]interface{}{ + om := map[string]interface{}{ "id": o.ID, "observer_id": strOrNil(o.ObserverID), "observer_name": strOrNil(o.ObserverName), @@ -1737,7 +1855,11 @@ func txToMap(tx *StoreTx) map[string]interface{} { "path_json": strOrNil(o.PathJSON), "timestamp": strOrNil(o.Timestamp), "direction": strOrNil(o.Direction), - }) + } + if o.ResolvedPath != nil { + om["resolved_path"] = o.ResolvedPath + } + obs = append(obs, om) } m["observations"] = obs return m @@ -1884,7 +2006,7 @@ func (s *PacketStore) buildDistanceIndex() { if cached, ok := hopCache[hop]; ok { return cached } - r := pm.resolve(hop) + r, _, _ := pm.resolveWithContext(hop, nil, s.graph) hopCache[hop] = r return r } @@ -3545,7 +3667,7 @@ func (s *PacketStore) computeAnalyticsTopology(region string) map[string]interfa if cached, ok := hopCache[hop]; ok { return cached } - r := pm.resolve(hop) + r, _, _ := pm.resolveWithContext(hop, nil, s.graph) hopCache[hop] = r return r } @@ -5536,7 +5658,7 @@ func (s *PacketStore) computeAnalyticsSubpaths(region string, minLen, maxLen, li } return hop } - r := pm.resolve(hop) + r, _, _ := pm.resolveWithContext(hop, nil, s.graph) hopCache[hop] = r if r != nil { return r.Name @@ -5673,7 +5795,7 @@ func (s *PacketStore) GetSubpathDetail(rawHops []string) map[string]interface{} // Resolve the requested hops nodes := make([]map[string]interface{}, len(rawHops)) for i, hop := range rawHops { - r := pm.resolve(hop) + r, _, _ := pm.resolveWithContext(hop, nil, s.graph) entry := map[string]interface{}{"hop": hop, "name": hop, "lat": nil, "lon": nil, "pubkey": nil} if r != nil { entry["name"] = r.Name @@ -5752,7 +5874,7 @@ func (s *PacketStore) GetSubpathDetail(rawHops []string) map[string]interface{} // Full parent path (resolved) resolved := make([]string, len(hops)) for i, h := range hops { - r := pm.resolve(h) + r, _, _ := pm.resolveWithContext(h, nil, s.graph) if r != nil { resolved[i] = r.Name } else { diff --git a/cmd/server/types.go b/cmd/server/types.go index 1a7812e1..6b6fbceb 100644 --- a/cmd/server/types.go +++ b/cmd/server/types.go @@ -240,6 +240,7 @@ type TransmissionResp struct { SNR interface{} `json:"snr"` RSSI interface{} `json:"rssi"` PathJSON interface{} `json:"path_json"` + ResolvedPath []*string `json:"resolved_path,omitempty"` Direction interface{} `json:"direction"` Score interface{} `json:"score,omitempty"` Observations []ObservationResp `json:"observations,omitempty"` @@ -254,6 +255,7 @@ type ObservationResp struct { SNR interface{} `json:"snr"` RSSI interface{} `json:"rssi"` PathJSON interface{} `json:"path_json"` + ResolvedPath []*string `json:"resolved_path,omitempty"` Timestamp interface{} `json:"timestamp"` }