Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
144 changes: 144 additions & 0 deletions internal/ingest/advert_endpoint_update_test.go
Original file line number Diff line number Diff line change
@@ -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 = &current, 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")
}
}
14 changes: 7 additions & 7 deletions internal/ingest/packet.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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()
Expand Down
Loading