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
4 changes: 2 additions & 2 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -52,10 +52,10 @@ jobs:
BEACON_TEST_POSTGRES_DSN: postgres://postgres:backup-ci-only@127.0.0.1:5432/postgres?sslmode=disable
run: go test ./cmd/beacon -run '^TestBackupConfigurationPostgres$' -count=1 -v

- name: Verify signal and path aggregates against PostgreSQL 16
- name: Verify stats and endpoint queries against PostgreSQL 16
env:
BEACON_TEST_POSTGRES_DSN: postgres://postgres:backup-ci-only@127.0.0.1:5432/postgres?sslmode=disable
run: go test ./db -run '^Test(Signal|Paths)Postgres$' -count=1 -v
run: go test ./db -run '^Test(Signal|Paths|PacketEndpointsResolveLive)Postgres$' -count=1 -v

- name: Verify backup command against PostgreSQL 16
run: |
Expand Down
151 changes: 151 additions & 0 deletions db/endpoints.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,151 @@
// Copyright 2026 Beacon Contributors
// SPDX-License-Identifier: AGPL-3.0-or-later

package db

import (
"context"
"encoding/hex"
"log/slog"

sqlc "github.com/MeshCore-Beacon/beacon-server/db/sqlc"
"github.com/MeshCore-Beacon/beacon-server/internal/api"
"github.com/meshcore-go/meshcore-go"
)

// packetEndpoints is what a packet's payload says about its sender and recipient.
type packetEndpoints struct {
advert bool // source is the advert's own key, not a hash
advertKey []byte // nil for legacy adverts stored without one
source, destination []byte // one-byte hashes
}

func parsePacketEndpoints(payloadType int16, raw, originPubkey []byte) packetEndpoints {
var ep packetEndpoints
switch payloadType {
case int16(meshcore.PayloadTypeAdvert):
ep.advert, ep.advertKey = true, originPubkey
case int16(meshcore.PayloadTypeAnonReq):
if p, err := meshcore.AnonReqFromBytes(raw); err == nil {
ep.destination = []byte{p.Destination}
}
case int16(meshcore.PayloadTypeReq):
if p, err := meshcore.RequestFromBytes(raw); err == nil {
ep.source, ep.destination = []byte{p.Source}, []byte{p.Destination}
}
case int16(meshcore.PayloadTypeResponse):
if p, err := meshcore.ResponseFromBytes(raw); err == nil {
ep.source, ep.destination = []byte{p.Source}, []byte{p.Destination}
}
case int16(meshcore.PayloadTypeTxtMsg):
if p, err := meshcore.TextMessageFromBytes(raw); err == nil {
ep.source, ep.destination = []byte{p.Source}, []byte{p.Destination}
}
case int16(meshcore.PayloadTypePath):
if p, err := meshcore.PathFromBytes(raw); err == nil {
ep.source, ep.destination = []byte{p.Source}, []byte{p.Destination}
}
}
return ep
}

type endpointLookup struct {
iata string
ep packetEndpoints
}

// resolveEndpoints resolves a whole page of endpoints in at most two queries.
// Results line up with lookups; nil means no such endpoint or a failed lookup.
func (s *Store) resolveEndpoints(ctx context.Context, lookups []endpointLookup) (sources, destinations []*api.ResolvedHop) {
iatas, hashes, keys := map[string]bool{}, map[byte]bool{}, map[string][]byte{}
for _, l := range lookups {
for _, h := range [][]byte{l.ep.source, l.ep.destination} {
if len(h) == 1 {
iatas[l.iata], hashes[h[0]] = true, true
}
}
if len(l.ep.advertKey) > 0 {
keys[string(l.ep.advertKey)] = l.ep.advertKey
}
}

candidates := map[string]map[string][]api.ResolvedPathEntry{} // iata -> hash hex
hashesOK := len(hashes) > 0
if hashesOK {
params := sqlc.ResolveEndpointHashPairsParams{}
for iata := range iatas {
params.Iatas = append(params.Iatas, iata)
}
for h := range hashes {
params.Hashes = append(params.Hashes, []byte{h})
}
rows, err := s.q.ResolveEndpointHashPairs(ctx, params)
if err != nil {
slog.Error("store: endpoint resolution failed", "component", "db", "error", err)
hashesOK = false
}
for _, r := range rows {
if candidates[r.Iata] == nil {
candidates[r.Iata] = map[string][]api.ResolvedPathEntry{}
}
key := hex.EncodeToString(r.Hash)
candidates[r.Iata][key] = append(candidates[r.Iata][key], api.ResolvedPathEntry{
NodeID: r.NodeID, Name: r.Name, Latitude: r.Latitude,
Longitude: r.Longitude, PublicKey: r.PublicKey,
})
}
}

adverts := map[string]*api.ResolvedNode{}
advertsOK := true
if len(keys) > 0 {
pubkeys := make([][]byte, 0, len(keys))
for _, k := range keys {
pubkeys = append(pubkeys, k)
}
rows, err := s.q.GetNodesByPubkeys(ctx, pubkeys)
if err != nil {
slog.Error("store: advert source lookup failed", "component", "db", "error", err)
advertsOK = false
}
for _, r := range rows {
adverts[string(r.PublicKey)] = &api.ResolvedNode{
ID: r.ID, Name: r.Name, PublicKey: hex.EncodeToString(r.PublicKey),
Latitude: r.Latitude, Longitude: r.Longitude,
}
}
}

hop := func(iata string, h []byte) *api.ResolvedHop {
if len(h) != 1 || !hashesOK {
return nil
}
r := api.BuildResolvedPath([][]byte{h}, candidates[iata])[0]
return &r
}
sources = make([]*api.ResolvedHop, len(lookups))
destinations = make([]*api.ResolvedHop, len(lookups))
for i, l := range lookups {
if l.ep.advert {
if advertsOK {
r := api.ResolveExactNode(adverts[string(l.ep.advertKey)])
sources[i] = &r
}
} else {
sources[i] = hop(l.iata, l.ep.source)
}
destinations[i] = hop(l.iata, l.ep.destination)
}
return sources, destinations
}

// fillLatestEndpoints resolves the latest observer's endpoints for a page of packets.
// lookups line up with items; rows without a latest observer are skipped.
func (s *Store) fillLatestEndpoints(ctx context.Context, items []api.PacketSummary, lookups []endpointLookup) {
sources, destinations := s.resolveEndpoints(ctx, lookups)
for i := range items {
if lo := items[i].LatestObserver; lo != nil {
lo.ResolvedSource, lo.ResolvedDestination = sources[i], destinations[i]
}
}
}
140 changes: 140 additions & 0 deletions db/endpoints_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,140 @@
// Copyright 2026 Beacon Contributors
// SPDX-License-Identifier: AGPL-3.0-or-later

package db

import (
"context"
"errors"
"sort"
"testing"

sqlc "github.com/MeshCore-Beacon/beacon-server/db/sqlc"
mockdb "github.com/MeshCore-Beacon/beacon-server/db/sqlc/mock"
"github.com/MeshCore-Beacon/beacon-server/internal/api"
"github.com/google/uuid"
"github.com/meshcore-go/meshcore-go"
"go.uber.org/mock/gomock"
)

func TestParsePacketEndpoints(t *testing.T) {
msg := append([]byte{0xbb, 0xaa, 0, 0}, make([]byte, 16)...)
key := []byte{1, 2, 3}
for _, tc := range []struct {
name string
kind uint8
raw, key []byte
want packetEndpoints
}{
{"direct message", meshcore.PayloadTypeTxtMsg, msg, nil, packetEndpoints{source: []byte{0xaa}, destination: []byte{0xbb}}},
{"request", meshcore.PayloadTypeReq, msg, nil, packetEndpoints{source: []byte{0xaa}, destination: []byte{0xbb}}},
{"advert", meshcore.PayloadTypeAdvert, nil, key, packetEndpoints{advert: true, advertKey: key}},
{"channel message", meshcore.PayloadTypeGrpTxt, msg, nil, packetEndpoints{}},
{"truncated", meshcore.PayloadTypeTxtMsg, []byte{0xbb}, nil, packetEndpoints{}},
} {
t.Run(tc.name, func(t *testing.T) {
got := parsePacketEndpoints(int16(tc.kind), tc.raw, tc.key)
if got.advert != tc.want.advert || string(got.advertKey) != string(tc.want.advertKey) ||
string(got.source) != string(tc.want.source) || string(got.destination) != string(tc.want.destination) {
t.Fatalf("got %+v, want %+v", got, tc.want)
}
})
}
}

func hopNames(h *api.ResolvedHop) []string {
if h == nil {
return nil
}
names := []string{h.Confidence}
for _, n := range h.Nodes {
names = append(names, *n.Name)
}
sort.Strings(names[1:])
return names
}

func TestResolveEndpointsBatchesAPage(t *testing.T) {
ctrl := gomock.NewController(t)
mock := mockdb.NewMockQuerier(ctrl)
store := &Store{q: mock}
name := func(s string) *string { return &s }
alice, bob, carol := uuid.New(), uuid.New(), uuid.New()
advertKey := []byte{0xaa, 0x11}

mock.EXPECT().ResolveEndpointHashPairs(gomock.Any(), gomock.Any()).DoAndReturn(
func(_ context.Context, p sqlc.ResolveEndpointHashPairsParams) ([]sqlc.ResolveEndpointHashPairsRow, error) {
if len(p.Iatas) != 2 || len(p.Hashes) != 2 {
t.Errorf("want distinct IATAs and hashes, got %v %v", p.Iatas, p.Hashes)
}
return []sqlc.ResolveEndpointHashPairsRow{
{Iata: "YYZ", Hash: []byte{0xaa}, NodeID: alice, Name: name("Alice")},
{Iata: "YYZ", Hash: []byte{0xaa}, NodeID: bob, Name: name("Bob")},
{Iata: "YYZ", Hash: []byte{0xbb}, NodeID: carol, Name: name("Carol")},
{Iata: "YVR", Hash: []byte{0xaa}, NodeID: alice, Name: name("Alice")},
{Iata: "YVR", Hash: []byte{0xbb}, NodeID: carol, Name: name("Carol")}, // cross-product extra
}, nil
})
mock.EXPECT().GetNodesByPubkeys(gomock.Any(), [][]byte{advertKey}).Return(
[]sqlc.GetNodesByPubkeysRow{{ID: alice, PublicKey: advertKey, Name: name("Alice")}}, nil)

sources, destinations := store.resolveEndpoints(context.Background(), []endpointLookup{
{iata: "YYZ", ep: packetEndpoints{source: []byte{0xaa}, destination: []byte{0xbb}}},
{iata: "YVR", ep: packetEndpoints{source: []byte{0xaa}}},
{iata: "YYZ", ep: packetEndpoints{advert: true, advertKey: advertKey}},
{iata: "YYZ", ep: packetEndpoints{advert: true, advertKey: advertKey}},
{iata: "YYZ", ep: packetEndpoints{advert: true}},
{},
})

for i, tc := range []struct{ source, destination []string }{
{[]string{"ambiguous", "Alice", "Bob"}, []string{"high", "Carol"}},
{[]string{"high", "Alice"}, nil},
{[]string{"high", "Alice"}, nil},
{[]string{"high", "Alice"}, nil},
{[]string{"none"}, nil},
{nil, nil},
} {
if got := hopNames(sources[i]); !equalStrings(got, tc.source) {
t.Errorf("row %d source = %v, want %v", i, got, tc.source)
}
if got := hopNames(destinations[i]); !equalStrings(got, tc.destination) {
t.Errorf("row %d destination = %v, want %v", i, got, tc.destination)
}
}
}

func TestResolveEndpointsFailedLookupLeavesHopsEmpty(t *testing.T) {
ctrl := gomock.NewController(t)
mock := mockdb.NewMockQuerier(ctrl)
store := &Store{q: mock}
mock.EXPECT().ResolveEndpointHashPairs(gomock.Any(), gomock.Any()).Return(nil, errors.New("boom"))

sources, destinations := store.resolveEndpoints(context.Background(), []endpointLookup{
{iata: "YYZ", ep: packetEndpoints{source: []byte{0xaa}, destination: []byte{0xbb}}},
})
if sources[0] != nil || destinations[0] != nil {
t.Fatal("a failed lookup must not report endpoints as unresolved")
}
}

func TestResolveEndpointsSkipsQueriesWhenNothingToResolve(t *testing.T) {
ctrl := gomock.NewController(t)
store := &Store{q: mockdb.NewMockQuerier(ctrl)} // no EXPECT: any query fails the test
sources, _ := store.resolveEndpoints(context.Background(), []endpointLookup{{}, {iata: "YYZ"}})
if len(sources) != 2 || sources[0] != nil || sources[1] != nil {
t.Fatal("rows without endpoints must stay empty")
}
}

func equalStrings(a, b []string) bool {
if len(a) != len(b) {
return false
}
for i := range a {
if a[i] != b[i] {
return false
}
}
return true
}
6 changes: 6 additions & 0 deletions db/migrations/038_drop_resolved_endpoints.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
-- Copyright 2026 Beacon Contributors
-- SPDX-License-Identifier: AGPL-3.0-or-later

-- Endpoints now resolve at read time like path hops. The per-observation
-- snapshots were ~60% of packet_observations on prod.
ALTER TABLE packet_observations DROP COLUMN IF EXISTS resolved_endpoints;
Loading
Loading