From e1c4fa2f70ae6b5efc0054b91433460d07f06ce1 Mon Sep 17 00:00:00 2001 From: n30nex Date: Sat, 26 Sep 2026 03:12:13 -0400 Subject: [PATCH 1/2] fix(ingest): resolve advert endpoints after node updates --- .../ingest/advert_endpoint_update_test.go | 144 ++++++++++++++++++ internal/ingest/endpoint_matching_test.go | 4 + internal/ingest/packet.go | 14 +- 3 files changed, 155 insertions(+), 7 deletions(-) create mode 100644 internal/ingest/advert_endpoint_update_test.go diff --git a/internal/ingest/advert_endpoint_update_test.go b/internal/ingest/advert_endpoint_update_test.go new file mode 100644 index 00000000..aea19c4d --- /dev/null +++ b/internal/ingest/advert_endpoint_update_test.go @@ -0,0 +1,144 @@ +// Copyright 2026 Beacon Contributors +// SPDX-License-Identifier: AGPL-3.0-or-later + +package ingest + +import ( + "context" + "encoding/json" + "errors" + "testing" + "time" + + "github.com/MeshCore-Beacon/beacon-server/internal/api" + "github.com/MeshCore-Beacon/beacon-server/internal/hub" + "github.com/google/uuid" + "github.com/meshcore-go/meshcore-go" +) + +type advertEndpointDB struct { + *endpointCaptureDB + found, duplicate bool + lookups, upserts int +} + +func (s *advertEndpointDB) GetNodeByPubkey(context.Context, []byte) (uuid.UUID, error) { + if !s.found { + return uuid.Nil, errors.New("node not found") + } + return s.node.ID, nil +} + +func (s *advertEndpointDB) GetNodesByIDs(ctx context.Context, ids []uuid.UUID) (map[uuid.UUID]*api.ResolvedNode, error) { + s.lookups++ + return s.endpointCaptureDB.GetNodesByIDs(ctx, ids) +} + +func (s *advertEndpointDB) UpsertNode(_ context.Context, n UpsertNodeParams, _ RadioSettings) (uuid.UUID, error) { + s.node.Name = &n.Name + s.found = true + s.upserts++ + return s.node.ID, nil +} + +func (s *advertEndpointDB) InsertObservation(context.Context, InsertObservationParams) (bool, error) { + return !s.duplicate, nil +} + +func TestAdvertEndpointUsesUpdatedName(t *testing.T) { + for _, existing := range []bool{false, true} { + t.Run(map[bool]string{false: "first advert", true: "renamed node"}[existing], func(t *testing.T) { + w, base := newTestWorker() + old := "Previous repeater" + store := &advertEndpointDB{endpointCaptureDB: &endpointCaptureDB{stubDB: base, node: api.ResolvedNode{ID: uuid.New(), Name: &old, PublicKey: "aa"}}, found: existing} + w.db = store + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + go w.hub.Run() + client := w.hub.NewClient() + w.hub.AddScope(client, "rename", hub.Scope{Events: []hub.EventType{hub.EventPacketObservation, hub.EventObserverStatus}}) + defer w.hub.Remove(client) + waitForSummarySubscriber(t, ctx, w.hub, client) + packet := buildAdvertPacketWithData(t, append([]byte{meshcore.AdvertTypeRepeater | meshcore.AdvertNameMask}, []byte("Renamed repeater")...), false) + w.handlePacket(ctx, "YYZ", "0102", packetEnvelope(t, packet)) + for { + select { + case event := <-client.Send: + if event.Type != hub.EventPacketObservation { + continue + } + var got packetObservationEvent + if err := json.Unmarshal(event.Payload, &got); err != nil { + t.Fatal(err) + } + hop := got.Observation.ResolvedSource + if hop == nil || len(hop.Nodes) != 1 || hop.Nodes[0].Name == nil || *hop.Nodes[0].Name != "Renamed repeater" { + t.Fatalf("live source did not use the newly saved name: %+v", hop) + } + if store.upserts != 1 || hop.Nodes[0].ID != store.node.ID || hop.Confidence != "high" { + t.Fatal("advert identity or confidence changed") + } + return + case <-ctx.Done(): + t.Fatal("packet event missing") + } + } + }) + } +} + +func TestDuplicateAdvertSkipsLiveEndpointLookup(t *testing.T) { + w, base := newTestWorker() + store := &advertEndpointDB{endpointCaptureDB: &endpointCaptureDB{stubDB: base, node: api.ResolvedNode{ID: uuid.New(), PublicKey: "aa"}}, found: true, duplicate: true} + w.db = store + w.handlePacket(context.Background(), "YYZ", "0102", packetEnvelope(t, buildAdvertPacket(t, false))) + if store.lookups != 0 || store.upserts != 0 { + t.Fatalf("duplicate observation performed live-only work: lookups=%d upserts=%d", store.lookups, store.upserts) + } +} + +func TestRepeatAdvertUsesCurrentIdentityWithoutReapplyingAdvert(t *testing.T) { + r := newRepeatHarness(t, true) + store := &advertEndpointDB{endpointCaptureDB: &endpointCaptureDB{stubDB: r.db.stubDB, node: api.ResolvedNode{ID: uuid.New(), PublicKey: "aa"}}} + r.w.db = store + packet := buildAdvertPacketWithData(t, append([]byte{meshcore.AdvertTypeRepeater | meshcore.AdvertNameMask}, []byte("Advert name")...), false) + hear := func(path byte) []packetObservationEvent { + packet.Path, packet.PathLength = []byte{path}, 1 + r.w.handlePacket(r.ctx, "YOW", "0102", packetEnvelope(t, packet)) + r.w.hub.Broadcast(hub.Event{Type: hub.EventObserverStatus}) + var events []packetObservationEvent + for { + select { + case event := <-r.client.Send: + if event.Type == hub.EventObserverStatus { + return events + } + var payload packetObservationEvent + if err := json.Unmarshal(event.Payload, &payload); err != nil { + t.Fatal(err) + } + events = append(events, payload) + case <-r.ctx.Done(): + t.Fatal("repeat marker missing") + } + } + } + first := hear(0x11) + if len(first) != 1 || first[0].Packet.IsRepeat || store.upserts != 1 { + t.Fatal("first advert did not produce one stored hearing") + } + // A later advert has updated the identity before this older packet is heard again. + current := "Current name" + store.node.Name, store.duplicate = ¤t, true + repeated := hear(0x22) + if len(repeated) != 1 || !repeated[0].Packet.IsRepeat || repeated[0].Observation.ResolvedSource == nil { + t.Fatal("new-path repeat lost its resolved source") + } + if got := repeated[0].Observation.ResolvedSource.Nodes[0].Name; got == nil || *got != current || store.upserts != 1 { + t.Fatal("repeat reapplied the old advert or returned a stale identity") + } + lookups := store.lookups + if got := hear(0x22); len(got) != 0 || store.lookups != lookups || store.upserts != 1 { + t.Fatal("suppressed repeat performed live-only work") + } +} diff --git a/internal/ingest/endpoint_matching_test.go b/internal/ingest/endpoint_matching_test.go index 52d12195..24596fb7 100644 --- a/internal/ingest/endpoint_matching_test.go +++ b/internal/ingest/endpoint_matching_test.go @@ -19,6 +19,10 @@ type endpointRoutingDB struct { paths []string } +func (s *endpointRoutingDB) InsertObservation(context.Context, InsertObservationParams) (bool, error) { + return true, nil +} + func (s *endpointRoutingDB) ResolveEndpointHashes(_ context.Context, iata string, hashes [][]byte) (map[string][]api.ResolvedPathEntry, error) { for _, hash := range hashes { s.endpoints = append(s.endpoints, iata+":"+hex.EncodeToString(hash)) diff --git a/internal/ingest/packet.go b/internal/ingest/packet.go index 92f3da78..630dcd12 100644 --- a/internal/ingest/packet.go +++ b/internal/ingest/packet.go @@ -862,7 +862,13 @@ func (w *Worker) handlePacket(ctx context.Context, iata, pubkeyHex string, raw [ // A duplicate is streamed, never stored, to includeRepeats clients when its path is new. repeat := !inserted && w.hub.RepeatsWanted() && w.hub.MarkSent(packetHash[:], id[:], packet.Path) if inserted || repeat { - // Endpoints for the live event only; stored rows resolve them at read time. + if inserted { + w.handlePayloadTypeSideEffects(ctx, packet, iata, packetHash[:], radio, scopeID, matchedScope, pubkeyBytes, float32(parseNumber(envelope.SNR))) + if w.hub.RepeatsWanted() { + w.hub.MarkSent(packetHash[:], id[:], packet.Path) // so broker copies of it aren't repeats + } + } + // Resolve after advert updates; suppressed copies need no endpoint lookup. var resolvedSource, resolvedDestination *api.ResolvedHop if packet.PayloadType() == meshcore.PayloadTypeAdvert && originPubkey != nil { // Exact match: ADVERT carries the sender's real identity pubkey, not a @@ -885,12 +891,6 @@ func (w *Worker) handlePacket(ctx context.Context, iata, pubkeyHex string, raw [ resolvedDestination = &hop } } - if inserted { - w.handlePayloadTypeSideEffects(ctx, packet, iata, packetHash[:], radio, scopeID, matchedScope, pubkeyBytes, float32(parseNumber(envelope.SNR))) - if w.hub.RepeatsWanted() { - w.hub.MarkSent(packetHash[:], id[:], packet.Path) // so broker copies of it aren't repeats - } - } evt := packetObservationEvent{} evt.PacketHash = hex.EncodeToString(packetHash[:]) evt.Packet.PayloadType = packet.PayloadType() From 319df7da5246a71f9693f923d2124bf831eabb82 Mon Sep 17 00:00:00 2001 From: n30nex Date: Tue, 29 Sep 2026 18:27:42 -0400 Subject: [PATCH 2/2] fix: address maintainer review for #166 --- internal/ingest/endpoint_matching_test.go | 4 ---- 1 file changed, 4 deletions(-) diff --git a/internal/ingest/endpoint_matching_test.go b/internal/ingest/endpoint_matching_test.go index 24596fb7..52d12195 100644 --- a/internal/ingest/endpoint_matching_test.go +++ b/internal/ingest/endpoint_matching_test.go @@ -19,10 +19,6 @@ type endpointRoutingDB struct { paths []string } -func (s *endpointRoutingDB) InsertObservation(context.Context, InsertObservationParams) (bool, error) { - return true, nil -} - func (s *endpointRoutingDB) ResolveEndpointHashes(_ context.Context, iata string, hashes [][]byte) (map[string][]api.ResolvedPathEntry, error) { for _, hash := range hashes { s.endpoints = append(s.endpoints, iata+":"+hex.EncodeToString(hash))