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
2 changes: 1 addition & 1 deletion .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -58,7 +58,7 @@ jobs:
- 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|PacketEndpointsResolveLive|ObserverMetrics|RouteEvidence|RouteEvidenceIndex|AnalyticsRetention|AnalyticsRetentionConcurrent|DeleteOldPacketsBatches|MeshMapperCatalogue|ChannelMessageScopes)Postgres$' -count=1 -v
run: go test ./db -run '^Test(Signal|Paths|PacketSummaries|PacketEndpointsResolveLive|ObserverMetrics|RouteEvidence|RouteEvidenceIndex|AnalyticsRetention|AnalyticsRetentionConcurrent|DeleteOldPacketsBatches|MeshMapperCatalogue|ChannelMessageScopes|Reconfirm[A-Za-z]*)Postgres$' -count=1 -v

- name: Verify backup command against PostgreSQL 16
run: |
Expand Down
9 changes: 6 additions & 3 deletions PROFILING.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,17 +31,20 @@ can enable a longer session. Profiling adds overhead while a capture is active.

## Capture behavior

- One 30-second sample immediately, then every 30 minutes.
- One 30-second sample immediately, then every 30 minutes, offset five minutes past
the half hour so a maintenance trigger due on the hour is not lost to the cooldown.
- Route reconfirmation requests an additional sample when the maintenance task starts.
This includes the retention step before route validation. A five-minute cooldown
between capture starts prevents overlap and repeated triggers from increasing load;
a periodic sample due during the cooldown runs when it ends.
a periodic sample due during the cooldown runs when it ends. A triggered sample
during the cooldown is skipped, and any sample cancels a periodic one waiting on
the cooldown.
- Background task stacks carry a `task` label while profiling is enabled.
- Shutdown or expiry stops the active sample and saves the shorter profile.
- Each profile is limited to 8 MiB. The dedicated directory is limited to 256 MiB
and 512 files, including metadata and files left by interrupted runs. The recorder
reserves space for a full capture before starting and stops when a limit is reached.
Files are never automatically deleted. Existing unrelated files consume the budget.
Files are never automatically deleted. Existing unrelated regular files consume the budget; subdirectories are ignored.
- Profiles and metadata are written with mode `0600`. A `.partial` file indicates
an interrupted capture and is not a completed profile.

Expand Down
2 changes: 1 addition & 1 deletion cmd/beacon/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -337,7 +337,7 @@ func main() {
CORS: cfg.CORS, Server: cfg.Server, Auth: cfg.Auth, RateLimit: resolved.RateLimit,
AdminRoutes: map[string]http.Handler{
"/accounts": handlers.AccountsRouter(store),
"/backup": handlers.BackupRouter(backupOpts),
"/backup": handlers.BackupRouter(backupOpts, ctx),
},
})

Expand Down
18 changes: 18 additions & 0 deletions db/analytics_retention_integration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -167,6 +167,24 @@ func TestAnalyticsRetentionPostgres(t *testing.T) {
t.Fatalf("missing unfiltered time index %s: %s %v", name, definition, err)
}
}
// The signal/path stores must read the 039 view definitions, not just 035/036.
store := &Store{q: sqlc.New(tx)}
hour := time.Now().UTC().Truncate(time.Hour)
signal, err := store.GetSignalStats(ctx, hour.Add(-6*24*time.Hour), hour, nil)
if err != nil {
t.Fatal(err)
}
// 12 observations; only the n=1 rows carry a signal (snr 0, rssi -100).
if signal.Receptions != 12 || signal.SNR.Samples != 4 || signal.RSSI.Samples != 4 {
t.Fatalf("signal stats on 039 views: %+v", signal)
}
paths, err := store.GetPathStats(ctx, hour.Add(-6*24*time.Hour), hour, nil)
if err != nil {
t.Fatal(err)
}
if paths.Receptions != 12 {
t.Fatalf("path stats on 039 views: %+v", paths)
}
before := analyticsSnapshot(t, ctx, tx)
q := sqlc.New(tx)
cutoff := time.Now().Add(-72 * time.Hour)
Expand Down
3 changes: 2 additions & 1 deletion db/observers.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ package db
import (
"context"
"encoding/hex"
"encoding/json"
"fmt"
"log/slog"
"math"
Expand Down Expand Up @@ -116,7 +117,7 @@ func (s *Store) GetObserver(ctx context.Context, observerID uuid.UUID) (*api.Obs
RadioCR: obs.RadioCr,
BatteryLevel: obs.BatteryLevel,
UptimeSeconds: obs.UptimeSeconds,
StatusMetadata: obs.StatusMetadata,
StatusMetadata: json.RawMessage(obs.StatusMetadata),
FirstSeen: obs.FirstSeen.Time.UnixMilli(),
LastSeen: obs.LastSeen.Time.UnixMilli(),
ObservationCount: *obs.ObservationCount,
Expand Down
12 changes: 12 additions & 0 deletions db/observers_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,9 @@ package db

import (
"context"
"encoding/json"
"errors"
"strings"
"testing"
"time"

Expand Down Expand Up @@ -178,6 +180,7 @@ func TestGetObserver_OnlineStatus(t *testing.T) {

observerID := uuid.MustParse("00000000-0000-0000-0000-000000000001")
obsCount := int64(10)
statusMetadata := []byte(`{"stats":{"noise_floor":-97}}`)

mock.EXPECT().
GetObserverByID(gomock.Any(), observerID).
Expand All @@ -188,6 +191,7 @@ func TestGetObserver_OnlineStatus(t *testing.T) {
FirstSeen: pgtype.Timestamptz{Time: time.Now().Add(-time.Hour), Valid: true},
LastSeen: pgtype.Timestamptz{Time: time.Now().Add(-time.Minute), Valid: true},
LastStatusAt: pgtype.Timestamptz{Time: time.Now().Add(-time.Minute), Valid: true},
StatusMetadata: statusMetadata,
}, nil)

mock.EXPECT().
Expand All @@ -213,6 +217,14 @@ func TestGetObserver_OnlineStatus(t *testing.T) {
if observer.IATA != "YVR" {
t.Errorf("expected IATA YVR, got %s", observer.IATA)
}

out, err := json.Marshal(observer)
if err != nil {
t.Fatal(err)
}
if !strings.Contains(string(out), `"statusMetadata":{"stats":{"noise_floor":-97}}`) {
t.Fatalf("statusMetadata not an object: %s", out)
}
}

func TestGetObserver_OfflineStatus(t *testing.T) {
Expand Down
24 changes: 15 additions & 9 deletions db/queries/queries.sql
Original file line number Diff line number Diff line change
Expand Up @@ -1482,10 +1482,21 @@ REFRESH MATERIALIZED VIEW CONCURRENTLY mv_radio_presets;
-- name: RefreshObserverActivity :exec
REFRESH MATERIALIZED VIEW CONCURRENTLY mv_observer_activity_hourly;

-- name: AmbiguousPrefixes :many
-- Hop prefixes that match >1 node in an IATA, per width. Computed once per reconfirm run.
SELECT iata::text AS iata, 1::int AS len, prefix_1 AS prefix FROM node_short_ids GROUP BY iata, prefix_1 HAVING COUNT(*) > 1
UNION ALL
SELECT iata::text, 2::int, prefix_2 FROM node_short_ids GROUP BY iata, prefix_2 HAVING COUNT(*) > 1
UNION ALL
SELECT iata::text, 3::int, prefix_3 FROM node_short_ids GROUP BY iata, prefix_3 HAVING COUNT(*) > 1
UNION ALL
SELECT iata::text, 4::int, prefix_4 FROM node_short_ids GROUP BY iata, prefix_4 HAVING COUNT(*) > 1;

-- name: ReconfirmRoutes :one
-- Checks one batch of least-recently-reconfirmed routes: deletes those with a departed
-- hop node or a hop prefix now matching >1 node in that IATA (length-aware:
-- 1/2/3/4-byte hop prefixes check prefix_1/2/3/4), and stamps the survivors.
-- 1/2/3/4-byte hop prefixes check prefix_1/2/3/4; ambiguity set supplied by AmbiguousPrefixes),
-- and stamps the survivors.
WITH batch AS MATERIALIZED (
SELECT r.iata, r.path_key, r.node_ids, r.hash_prefix
FROM known_routes r
Expand All @@ -1494,14 +1505,9 @@ WITH batch AS MATERIALIZED (
LIMIT @batch_size
FOR UPDATE OF r SKIP LOCKED
),
amb AS MATERIALIZED (
SELECT iata, 1 AS len, prefix_1 AS p FROM node_short_ids GROUP BY iata, prefix_1 HAVING COUNT(*) > 1
UNION ALL
SELECT iata, 2, prefix_2 FROM node_short_ids GROUP BY iata, prefix_2 HAVING COUNT(*) > 1
UNION ALL
SELECT iata, 3, prefix_3 FROM node_short_ids GROUP BY iata, prefix_3 HAVING COUNT(*) > 1
UNION ALL
SELECT iata, 4, prefix_4 FROM node_short_ids GROUP BY iata, prefix_4 HAVING COUNT(*) > 1
amb AS (
SELECT a.iata::char(3) AS iata, a.len, a.p
FROM ROWS FROM (unnest(@amb_iata::text[]), unnest(@amb_len::int[]), unnest(@amb_prefix::bytea[])) AS a(iata, len, p)
),
dead AS (
SELECT b.iata, b.path_key
Expand Down
34 changes: 22 additions & 12 deletions db/reconfirm_integration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,10 +17,19 @@ import (
"github.com/jackc/pgx/v5/pgxpool"
)

func ambiguity(t *testing.T, ctx context.Context, store *Store) AmbiguousPrefixes {
t.Helper()
amb, err := store.AmbiguousPrefixes(ctx)
if err != nil {
t.Fatal(err)
}
return amb
}

func TestReconfirmReleasesEarlierBatchPostgres(t *testing.T) {
ctx, pool, store := reconfirmPool(t)
before := time.Now()
if n, err := store.ReconfirmRoutes(ctx, 1, before); err != nil || n != 1 {
if n, err := store.ReconfirmRoutes(ctx, 1, before, ambiguity(t, ctx, store)); err != nil || n != 1 {
t.Fatalf("first batch: %d, %v", n, err)
}
gate, err := pool.Begin(ctx)
Expand All @@ -40,8 +49,9 @@ CREATE TRIGGER pause_second_route BEFORE UPDATE ON known_routes FOR EACH ROW EXE
t.Fatal(err)
}
finished := make(chan error, 1)
amb := ambiguity(t, ctx, store)
go func() {
_, err := store.ReconfirmRoutes(ctx, 1, before)
_, err := store.ReconfirmRoutes(ctx, 1, before, amb)
finished <- err
}()
deadline := time.Now().Add(3 * time.Second)
Expand Down Expand Up @@ -131,7 +141,7 @@ UPDATE known_routes SET node_ids=node_ids || md5('missing')::uuid, hash_prefix=h
original := snapshot()
var checked int64
for range 10 {
n, err := store.ReconfirmRoutes(ctx, 3, before)
n, err := store.ReconfirmRoutes(ctx, 3, before, ambiguity(t, ctx, store))
if err != nil {
t.Fatal(err)
}
Expand Down Expand Up @@ -165,7 +175,7 @@ UPDATE known_routes SET node_ids=node_ids || md5('missing')::uuid, hash_prefix=h
func TestReconfirmCancelledBatchRollsBackPostgres(t *testing.T) {
ctx, pool, store := reconfirmPool(t)
before := time.Now()
if _, err := store.ReconfirmRoutes(ctx, 1, before); err != nil {
if _, err := store.ReconfirmRoutes(ctx, 1, before, ambiguity(t, ctx, store)); err != nil {
t.Fatal(err)
}
_, err := pool.Exec(ctx, `
Expand All @@ -177,7 +187,7 @@ CREATE TRIGGER pause_reconfirm BEFORE UPDATE ON known_routes FOR EACH ROW EXECUT
}
batchCtx, cancel := context.WithTimeout(ctx, 100*time.Millisecond)
defer cancel()
if _, err := store.ReconfirmRoutes(batchCtx, 1, before); err == nil {
if _, err := store.ReconfirmRoutes(batchCtx, 1, before, ambiguity(t, ctx, store)); err == nil {
t.Fatal("expected batch cancellation")
}
// Waiting for this DDL also waits for cancellation to release the row lock.
Expand All @@ -191,18 +201,18 @@ CREATE TRIGGER pause_reconfirm BEFORE UPDATE ON known_routes FOR EACH ROW EXECUT
if total != 2 || stamped != 1 {
t.Fatalf("cancellation changed committed work: total %d, checked %d", total, stamped)
}
if n, err := store.ReconfirmRoutes(ctx, 1, before); err != nil || n != 1 {
if n, err := store.ReconfirmRoutes(ctx, 1, before, ambiguity(t, ctx, store)); err != nil || n != 1 {
t.Fatalf("cancelled route was not available to retry: %d, %v", n, err)
}
}

func TestReconfirmFutureCutoffDoesNotRepeatPostgres(t *testing.T) {
ctx, _, store := reconfirmPool(t)
before := time.Now().Add(time.Hour)
if n, err := store.ReconfirmRoutes(ctx, 2, before); err != nil || n != 2 {
if n, err := store.ReconfirmRoutes(ctx, 2, before, ambiguity(t, ctx, store)); err != nil || n != 2 {
t.Fatalf("first batch: %d, %v", n, err)
}
if n, err := store.ReconfirmRoutes(ctx, 2, before); err != nil || n != 0 {
if n, err := store.ReconfirmRoutes(ctx, 2, before, ambiguity(t, ctx, store)); err != nil || n != 0 {
t.Fatalf("completed routes consumed the run budget twice: %d, %v", n, err)
}
}
Expand All @@ -220,7 +230,7 @@ CREATE TRIGGER pause_valid_route BEFORE UPDATE ON known_routes FOR EACH ROW EXEC
before := time.Now()
batchCtx, cancel := context.WithTimeout(ctx, 100*time.Millisecond)
defer cancel()
if _, err := store.ReconfirmRoutes(batchCtx, 2, before); err == nil {
if _, err := store.ReconfirmRoutes(batchCtx, 2, before, ambiguity(t, ctx, store)); err == nil {
t.Fatal("expected batch cancellation")
}
if _, err := pool.Exec(ctx, `DROP TRIGGER pause_valid_route ON known_routes`); err != nil {
Expand All @@ -230,7 +240,7 @@ CREATE TRIGGER pause_valid_route BEFORE UPDATE ON known_routes FOR EACH ROW EXEC
if err := pool.QueryRow(ctx, `SELECT count(*) FROM known_routes`).Scan(&count); err != nil || count != 2 {
t.Fatalf("cancelled batch partially deleted routes: %d, %v", count, err)
}
if n, err := store.ReconfirmRoutes(ctx, 2, before); err != nil || n != 2 {
if n, err := store.ReconfirmRoutes(ctx, 2, before, ambiguity(t, ctx, store)); err != nil || n != 2 {
t.Fatalf("retry did not process both routes: %d, %v", n, err)
}
if err := pool.QueryRow(ctx, `SELECT count(*) FROM known_routes WHERE path_key=int4send(2)`).Scan(&count); err != nil || count != 1 {
Expand Down Expand Up @@ -299,7 +309,7 @@ func TestReconfirmSkipsBusyRoutePostgres(t *testing.T) {
}
batchCtx, cancel := context.WithTimeout(ctx, 500*time.Millisecond)
defer cancel()
if _, err := store.ReconfirmRoutes(batchCtx, 1, time.Now()); err != nil {
if _, err := store.ReconfirmRoutes(batchCtx, 1, time.Now(), ambiguity(t, ctx, store)); err != nil {
t.Fatalf("maintenance should skip the busy route and validate the next one: %v", err)
}
var checked, preserved int
Expand All @@ -312,7 +322,7 @@ func TestReconfirmSkipsBusyRoutePostgres(t *testing.T) {
if err := busy.Rollback(ctx); err != nil {
t.Fatal(err)
}
if _, err := store.ReconfirmRoutes(ctx, 1, time.Now()); err != nil {
if _, err := store.ReconfirmRoutes(ctx, 1, time.Now(), ambiguity(t, ctx, store)); err != nil {
t.Fatal(err)
}
if err := pool.QueryRow(ctx, `SELECT count(*) FROM known_routes WHERE last_reconfirmed_at > '2026-02-01'`).Scan(&checked); err != nil {
Expand Down
2 changes: 1 addition & 1 deletion db/route_evidence.go
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,7 @@ func (s *Store) GetRouteEvidence(ctx context.Context, iata, key string, query ap
}
route := toKnownRoutes([]knownRouteRow{{ID: row.ID, NodeIds: row.NodeIds, HashPrefix: row.HashPrefix, Iata: row.Iata, HopCount: row.HopCount, FirstSeen: row.FirstSeen, LastSeen: row.LastSeen, ObservationCount: row.ObservationCount}}, nodes)[0]
route.PathKey = hex.EncodeToString(row.PathKey)
out := &api.RouteEvidence{Page: api.Page[api.RouteObservation]{Items: []api.RouteObservation{}}, Route: route, WindowStart: query.Since.UnixMilli(), WindowEnd: query.Until.UnixMilli(), GeneratedAt: time.Now().UnixMilli(), MatchType: "saved_path_prefixes"}
out := &api.RouteEvidence{Items: []api.RouteObservation{}, Route: route, WindowStart: query.Since.UnixMilli(), WindowEnd: query.Until.UnixMilli(), GeneratedAt: time.Now().UnixMilli(), MatchType: "saved_path_prefixes"}
width, path, valid := savedRoutePath(row.HashPrefix, row.HopCount)
if !valid {
return out, nil
Expand Down
26 changes: 25 additions & 1 deletion db/routes.go
Original file line number Diff line number Diff line change
Expand Up @@ -287,12 +287,36 @@ func (s *Store) SearchCrossIATARoutes(ctx context.Context, fromHash, fromIATA, t
return results, nil
}

// AmbiguousPrefixes is the per-IATA set of hop prefixes that resolve to more than one node.
type AmbiguousPrefixes struct {
IATAs []string
Lens []int32
Prefixes [][]byte
}

func (s *Store) AmbiguousPrefixes(ctx context.Context) (AmbiguousPrefixes, error) {
rows, err := s.q.AmbiguousPrefixes(ctx)
if err != nil {
return AmbiguousPrefixes{}, err
}
amb := AmbiguousPrefixes{IATAs: make([]string, 0, len(rows)), Lens: make([]int32, 0, len(rows)), Prefixes: make([][]byte, 0, len(rows))}
for _, r := range rows {
amb.IATAs = append(amb.IATAs, r.Iata)
amb.Lens = append(amb.Lens, r.Len)
amb.Prefixes = append(amb.Prefixes, r.Prefix)
}
return amb, nil
}

// ReconfirmRoutes checks the batchSize least-recently-reconfirmed routes,
// deleting stale or ambiguous ones and stamping the survivors.
func (s *Store) ReconfirmRoutes(ctx context.Context, batchSize int32, before time.Time) (int64, error) {
func (s *Store) ReconfirmRoutes(ctx context.Context, batchSize int32, before time.Time, amb AmbiguousPrefixes) (int64, error) {
return s.q.ReconfirmRoutes(ctx, sqlc.ReconfirmRoutesParams{
BatchSize: batchSize,
Before: pgtype.Timestamptz{Time: before, Valid: true},
AmbIata: amb.IATAs,
AmbLen: amb.Lens,
AmbPrefix: amb.Prefixes,
})
}

Expand Down
15 changes: 15 additions & 0 deletions db/sqlc/mock/querier.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

5 changes: 4 additions & 1 deletion db/sqlc/querier.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading
Loading