Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
18 commits
Select commit Hold shift + click to select a range
bbba9a7
feat(ingestor): resolve the last hop of a flood path from the observe…
dborup Oct 3, 2026
dc08e63
feat(ingestor): build the hop prefix index from relay roles only (#188)
dborup Oct 3, 2026
f1d408c
feat(ingestor): backfill NULL resolved_path rows after the index is p…
dborup Oct 3, 2026
1cfb42d
fix(ingestor): resolved_path backfill refuses to run before the index…
dborup Oct 3, 2026
f9bae20
test(ingestor): reproduce the backfill running on a pre-build or empt…
dborup Oct 3, 2026
cbaa9ca
fix(ingestor): backfill waits for a post-build, non-empty neighbour g…
dborup Oct 3, 2026
d1db161
docs(config): document resolvedPathBackfill and its zero-means-defaul…
dborup Oct 3, 2026
0c623db
test(ingestor): kill the resolver and backfill mutants that survived …
dborup Oct 3, 2026
f5add67
test(ingestor): a tick publishes the post-build graph after a failed …
dborup Oct 3, 2026
e610dbb
test(ingestor): reproduce the observer resolved as its own last hop (…
dborup Oct 3, 2026
503ea75
fix(ingestor): the observer is never its own last hop in a flood path…
dborup Oct 3, 2026
2e2886a
docs(ingestor): correct the TRACE note and firmware references (#188)
dborup Oct 3, 2026
486aa9e
Merge origin/master into codex/issue-188-observer-anchor
dborup Oct 3, 2026
e369ae6
test(ingestor): reproduce the warm-up stopping after a full batch wit…
dborup Oct 3, 2026
da0494c
fix(ingestor): end the edge-build warm-up on rows scanned, not edges …
dborup Oct 3, 2026
addccb8
test(ingestor): reproduce observer edges built from DIRECT paths (#188)
dborup Oct 3, 2026
194512d
fix(ingestor): build observer<->last-hop edges from flood routes only…
dborup Oct 3, 2026
0f4c167
test(ingestor): pin the TRANSPORT_FLOOD observer edge separately (#188)
dborup Oct 3, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions cmd/ingestor/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -107,6 +107,10 @@ type Config struct {
// (internal/channelregistry). Approved channels are always loaded into
// the channel keys, whether or not new submissions are enabled.
ChannelProposals *channelregistry.Config `json:"channelProposals,omitempty"`

// ResolvedPathBackfill tunes the start-up re-resolution of NULL
// observations.resolved_path rows (#188, resolved_path_backfill.go).
ResolvedPathBackfill *ResolvedPathBackfillConfig `json:"resolvedPathBackfill,omitempty"`
}

// NeighborEdgesDaysOrDefault returns the configured pruning window or 5.
Expand Down
8 changes: 6 additions & 2 deletions cmd/ingestor/db.go
Original file line number Diff line number Diff line change
Expand Up @@ -1215,10 +1215,14 @@ func (s *Store) InsertTransmission(data *PacketData) (bool, error) {
// the ingestor now. Per #1560: use the context-aware resolver so
// 1-byte prefix collisions are disambiguated via NeighborGraph
// adjacency (anchored on from_pubkey for ADVERTs, previous hop
// otherwise). Empty resolved JSON → NULL via nilIfEmpty.
resolved := resolvePathWithContext(
// otherwise). Per #188: flood paths are also walked backwards from the
// observer, whose neighbour the last hop is. Empty resolved JSON → NULL
// via nilIfEmpty.
resolved := resolveObservationPath(
parsePathArray(data.PathJSON),
strings.ToLower(data.FromPubkey),
data.ObserverID,
data.RouteType,
s.neighborGraph.load(),
s.prefixIdx.load(),
)
Expand Down
7 changes: 7 additions & 0 deletions cmd/ingestor/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -440,6 +440,13 @@ func main() {
defer stopNeighborBuilder()
log.Printf("[neighbor-build] enabled (interval=%s)", NeighborEdgesBuilderInterval)

// #188: re-resolve NULL resolved_path rows now that the prefix index
// and neighbour graph are primed (StartNeighborEdgesBuilder above).
if enabled, batch, pause := cfg.ResolvedPathBackfillSettings(); enabled {
stopResolvedPathBackfill := store.StartResolvedPathBackfill(batch, pause)
defer stopResolvedPathBackfill()
}

// #1212: per-source stall watchdog. Detects "silently dead" sources
// where the client reports connected but no messages have flowed. Logs
// a WARN line every minute for any source silent for >5m. Scan every
Expand Down
128 changes: 104 additions & 24 deletions cmd/ingestor/neighbor_builder.go
Original file line number Diff line number Diff line change
Expand Up @@ -79,18 +79,39 @@ func (s *Store) StartNeighborEdgesBuilder(interval time.Duration) func() {
if err := s.RefreshNeighborGraph(); err != nil {
log.Printf("[neighbor-build] initial neighbor-graph refresh error: %v", err)
}
//
// Each call reads at most neighborBuilderMaxBatch observations, so the
// loop ends when a call read fewer than that (caught up), not when it
// produced few edges (PR #190 review). Every row read is newer than the
// watermark, so a call that persists an edge moves the watermark; a full
// batch with no edge cannot, and the next call would read the same rows.
caughtUp := false
for {
n, err := s.buildAndPersistNeighborEdges()
b, err := s.buildNeighborEdges()
if err != nil {
log.Printf("[neighbor-build] initial build error: %v", err)
break
}
wuTotal += n
if n < neighborBuilderMaxBatch {
wuTotal += b.edges
if b.caughtUp() {
caughtUp = true
break
}
if b.edges == 0 {
log.Printf("[neighbor-build] initial build cannot move past its watermark: %d observations yielded no edge; neighbor_edges stays incomplete and the resolved_path backfill waits", b.scanned)
break
}
}
log.Printf("[neighbor-build] initial build: %d edges upserted in %s", wuTotal, time.Since(wuStart))
// Publish the graph with the edges the warm-up just persisted (#188):
// the snapshot primed above predates them, and on a fresh or restored DB
// it is empty. Only a build that caught up publishes a post-build graph;
// otherwise the first tick that catches up does it.
if caughtUp {
if err := s.refreshBuiltNeighborGraph(); err != nil {
log.Printf("[neighbor-build] post-build neighbor-graph refresh error: %v", err)
}
}

var stopOnce sync.Once
go func() {
Expand All @@ -106,11 +127,18 @@ func (s *Store) StartNeighborEdgesBuilder(interval time.Duration) func() {
if err := s.RefreshPrefixIndex(); err != nil {
log.Printf("[neighbor-build] prefix-index refresh error: %v", err)
}
n, err := s.buildAndPersistNeighborEdges()
b, err := s.buildNeighborEdges()
n := b.edges
// Refresh the neighbor-graph snapshot after the edges
// build (#1560) so the context-aware resolver picks up
// newly persisted adjacencies on the next ingest.
if grErr := s.RefreshNeighborGraph(); grErr != nil {
// newly persisted adjacencies on the next ingest. After a
// successful build that caught up, it is a post-build
// snapshot (#188).
refresh := s.RefreshNeighborGraph
if err == nil && b.caughtUp() {
refresh = s.refreshBuiltNeighborGraph
}
if grErr := refresh(); grErr != nil {
log.Printf("[neighbor-build] neighbor-graph refresh error: %v", grErr)
}
dur := time.Since(start)
Expand All @@ -137,10 +165,30 @@ func (s *Store) StartNeighborEdgesBuilder(interval time.Duration) func() {
}
}

// buildAndPersistNeighborEdges scans transmissions + observations,
// neighborEdgesBuild reports one buildNeighborEdges call.
type neighborEdgesBuild struct {
edges int // edge rows upserted
scanned int // observation rows read, at most neighborBuilderMaxBatch
}

// caughtUp reports whether the scan read every observation newer than the
// watermark it started from: it read fewer rows than the cap.
func (b neighborEdgesBuild) caughtUp() bool {
return b.scanned < neighborBuilderMaxBatch
}

// buildAndPersistNeighborEdges is buildNeighborEdges reporting only the
// number of edge upserts.
func (s *Store) buildAndPersistNeighborEdges() (int, error) {
b, err := s.buildNeighborEdges()
return b.edges, err
}

// buildNeighborEdges scans transmissions + observations,
// extracts edge candidates (originator↔first-hop on ADVERTs;
// observer↔last-hop on all packet types) and upserts them into
// neighbor_edges. Returns count of attempted upserts.
// observer↔last-hop on flood routes) and upserts them into
// neighbor_edges. It reports the edge upserts and the observation rows
// read; fewer rows than neighborBuilderMaxBatch means it caught up.
//
// Watermark / delta semantics (#1339): the builder derives a watermark
// from MAX(neighbor_edges.last_seen). On an empty edges table (fresh
Expand All @@ -161,10 +209,11 @@ func (s *Store) StartNeighborEdgesBuilder(interval time.Duration) func() {
// SELECT of (lowered) pubkey prefixes from nodes. Prefixes with
// multiple candidates are skipped (matches the conservative
// resolution rule in cmd/server/extractEdgesFromObs).
func (s *Store) buildAndPersistNeighborEdges() (int, error) {
func (s *Store) buildNeighborEdges() (neighborEdgesBuild, error) {
var b neighborEdgesBuild
prefixIdx, err := buildPrefixIndex(s.db)
if err != nil {
return 0, fmt.Errorf("build prefix index: %w", err)
return b, fmt.Errorf("build prefix index: %w", err)
}

// Derive the watermark from the existing edges table. RFC3339
Expand All @@ -173,7 +222,7 @@ func (s *Store) buildAndPersistNeighborEdges() (int, error) {
// query and the parse return zero → full warm-up scan.
var watermarkRFC sql.NullString
if err := s.db.QueryRow(`SELECT MAX(last_seen) FROM neighbor_edges`).Scan(&watermarkRFC); err != nil {
return 0, fmt.Errorf("read watermark: %w", err)
return b, fmt.Errorf("read watermark: %w", err)
}
var watermarkEpoch int64
if watermarkRFC.Valid && watermarkRFC.String != "" {
Expand All @@ -184,6 +233,7 @@ func (s *Store) buildAndPersistNeighborEdges() (int, error) {

rows, err := s.db.Query(`SELECT
t.payload_type,
COALESCE(t.route_type, -1),
t.decoded_json,
COALESCE(t.from_pubkey, ''),
COALESCE(o.path_json, ''),
Expand All @@ -196,16 +246,18 @@ func (s *Store) buildAndPersistNeighborEdges() (int, error) {
ORDER BY o.timestamp
LIMIT ?`, watermarkEpoch, neighborBuilderMaxBatch)
if err != nil {
return 0, fmt.Errorf("scan observations: %w", err)
return b, fmt.Errorf("scan observations: %w", err)
}
defer rows.Close()

var edges []edgeRow
for rows.Next() {
b.scanned++
var payloadType sql.NullInt64
var routeType int
var decodedJSON, fromPubkey, pathJSON, observerID string
var epochTs int64
if err := rows.Scan(&payloadType, &decodedJSON, &fromPubkey, &pathJSON, &observerID, &epochTs); err != nil {
if err := rows.Scan(&payloadType, &routeType, &decodedJSON, &fromPubkey, &pathJSON, &observerID, &epochTs); err != nil {
continue
}
fromNode := strings.ToLower(fromPubkey)
Expand All @@ -228,7 +280,14 @@ func (s *Store) buildAndPersistNeighborEdges() (int, error) {
edges = append(edges, canonEdge(fromNode, resolved, ts))
}
}
if observerPK != "" {
// The last hop is the node the observer heard only on a flood
// route: flood forwarders append their hash (MeshCore Mesh.cpp:346-350
// routeRecvPacket, at a366955). A DIRECT path is the remaining
// planned route, from which each forwarder strips itself at the
// front (Mesh.cpp:89, removeSelfFromPath :334-342), so its last hop
// is the route's far end (PR #190 review). Unknown route types are
// skipped too.
if observerPK != "" && isFloodRoute(routeType) {
last := path[len(path)-1]
if resolved, ok := resolvePrefix(prefixIdx, last); ok && resolved != observerPK {
edges = append(edges, canonEdge(observerPK, resolved, ts))
Expand Down Expand Up @@ -274,7 +333,7 @@ func (s *Store) buildAndPersistNeighborEdges() (int, error) {
}

if len(edges) == 0 {
return 0, nil
return b, nil
}

// Wrap the whole edge-persist tx under writer-perf instrumentation
Expand Down Expand Up @@ -304,9 +363,10 @@ func (s *Store) buildAndPersistNeighborEdges() (int, error) {
return nil
})
if err != nil {
return 0, err
return b, err
}
return inserted, nil
b.edges = inserted
return b, nil
}

// canonEdge orders the pair so node_a <= node_b (matches the existing
Expand Down Expand Up @@ -336,20 +396,23 @@ func parsePathArray(s string) []string {
// considered ambiguous and skipped during resolution.
type prefixIndex map[string][]string

// buildPrefixIndex reads nodes.public_key and builds the prefix → pubkey
// map. We index every 1-byte (2 hex char) prefix length the firmware
// uses (1, 2, 3, 4, 6, 8). Memory cost is O(nodes × len(prefixLens)).
// buildPrefixIndex reads the relay nodes (isRelayRole) and builds the
// prefix → pubkey map. We index every 1-byte (2 hex char) prefix length the
// firmware uses (1, 2, 3, 4, 6, 8). Memory cost is O(nodes × len(prefixLens)).
func buildPrefixIndex(db *sql.DB) (prefixIndex, error) {
rows, err := db.Query(`SELECT public_key FROM nodes`)
rows, err := db.Query(`SELECT public_key, COALESCE(role, '') FROM nodes`)
if err != nil {
return nil, err
}
defer rows.Close()
idx := make(prefixIndex, 1024)
var prefixLens = []int{1 * 2, 2 * 2, 3 * 2, 4 * 2, 6 * 2, 8 * 2}
for rows.Next() {
var pk string
if err := rows.Scan(&pk); err != nil {
var pk, role string
if err := rows.Scan(&pk, &role); err != nil {
continue
}
if !isRelayRole(role) {
continue
}
pkLower := strings.ToLower(pk)
Expand All @@ -364,6 +427,23 @@ func buildPrefixIndex(db *sql.DB) (prefixIndex, error) {
return idx, nil
}

// isRelayRole reports whether a node of this role can appear as a hop in a
// path (#188). It is the server's canAppearInPath (cmd/server/store.go),
// duplicated because the two binaries share no package for it; the test
// cases mirror the server's TestCanAppearInPath.
//
// Firmware: repeaters and room servers forward unless disable_fwd is set
// (examples/simple_repeater and simple_room_server MyMesh::allowPacketForward).
// Companions and sensors ship with forwarding off (companion_radio
// NodePrefs.h repeat.disable_fwd = 1; simple_sensor SensorMesh.cpp
// disable_fwd = true) and can only opt in, so a hop through an opted-in
// companion stays unresolved, or resolves to the one relay that shares its
// prefix. That is the server's trade-off too.
func isRelayRole(role string) bool {
r := strings.ToLower(role)
return strings.Contains(r, "repeater") || strings.Contains(r, "room_server") || r == "room"
}

// resolvePrefix returns the single resolved pubkey if exactly one
// candidate matches, otherwise (zero || multiple), it returns ok=false
// (matches the conservative server-side resolver in
Expand Down
3 changes: 2 additions & 1 deletion cmd/ingestor/neighbor_builder_delta_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -35,8 +35,9 @@ func TestNeighborEdgesBuilderDeltaScan(t *testing.T) {
}
defer store.Close()

// Hop nodes need a relay role: the prefix index holds relays only (#188).
if _, err := store.db.Exec(
`INSERT INTO nodes (public_key, name) VALUES (?, ?), (?, ?)`,
`INSERT INTO nodes (public_key, name, role) VALUES (?, ?, 'repeater'), (?, ?, 'repeater')`,
"aaaaaaaaaa", "from-node",
"bbbbbbbbbb", "first-hop",
); err != nil {
Expand Down
Loading
Loading