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|PacketSummaries|PacketEndpointsResolveLive|ObserverMetrics|RouteEvidence|RouteEvidenceIndex|AnalyticsRetention|AnalyticsRetentionConcurrent|DeleteOldPacketsBatches|MeshMapperCatalogue|ChannelMessageScopes|NodeLocationReset|Reconfirm[A-Za-z]*)Postgres$' -count=1 -v
run: go test ./db -run '^Test(Signal|Paths|PacketSummaries|PacketEndpointsResolveLive|ObserverMetrics|RouteEvidence|RouteEvidenceIndex|AnalyticsRetention|AnalyticsRetentionConcurrent|DeleteOldPacketsBatches|MeshMapperCatalogue|ChannelMessageScopes|NodeLocationReset|RunMigrations[A-Za-z]*|Reconfirm[A-Za-z]*)Postgres$' -count=1 -v

- name: Verify backup command against PostgreSQL 16
run: |
Expand Down
6 changes: 3 additions & 3 deletions CONTRIBUTING.md
Original file line number Diff line number Diff line change
Expand Up @@ -79,8 +79,8 @@ swag init # if you changed any handler or api type (see below)
All schema changes must include a proper migration path:

- Add a new migration file to `db/migrations/` following the existing naming
convention (e.g. `002_add_observation_count.sql`). Do not modify existing
migration files — append only via new files.
convention (e.g. `002_add_observation_count.sql`). `001_baseline.sql` is the
2.0.0 schema; do not modify existing migration files — append only via new files.
- Update `db/queries/queries.sql` with any new or modified queries
- Re-run `sqlc generate` to regenerate `db/sqlc/`
- Update the store layer in `db/` to expose the new functionality
Expand Down Expand Up @@ -239,7 +239,7 @@ Scopes are optional but helpful for larger codebases. Common scopes: `api`,
```
cmd/beacon/ — main entry point, wiring, startup
db/ — store layer: sqlc-generated code + thin mapping layer
migrations/ — SQL schema (single file, append only)
migrations/ — SQL schema (2.0.0 baseline + numbered files, append only)
queries/ — SQL queries (input to sqlc)
sqlc/ — generated Go code (do not edit by hand)
internal/
Expand Down
8 changes: 1 addition & 7 deletions db/accounts_integration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -50,15 +50,9 @@ func TestAccountsPostgres(t *testing.T) {
t.Fatal(err)
}
defer pool.Close()
sql, err := migrationFiles.ReadFile("migrations/034_accounts.sql")
if err != nil {
if err := RunMigrations(ctx, pool); err != nil {
t.Fatal(err)
}
for i := 0; i < 2; i++ {
if err := applyMigration(ctx, pool, string(sql)); err != nil {
t.Fatal(err)
}
}
store := New(pool, 0, 0)
items, err := store.ListAccounts(ctx)
if err != nil || items == nil || len(items) != 0 {
Expand Down
68 changes: 23 additions & 45 deletions db/analytics_retention_integration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,13 +16,20 @@ import (

var retainedViews = []string{"mv_hourly_iata_stats", "mv_payload_breakdown_by_iata", "mv_top_observers_by_iata", "mv_top_talkers_by_iata", "mv_top_advertisers_by_iata", "mv_observer_activity_hourly", "mv_signal_stats_hourly", "mv_path_stats_hourly"}

func analyticsSnapshot(t *testing.T, ctx context.Context, tx pgx.Tx) map[string]string {
func refreshRetainedViews(t *testing.T, ctx context.Context, tx pgx.Tx) {
t.Helper()
result := map[string]string{}
for _, view := range retainedViews {
if _, err := tx.Exec(ctx, "REFRESH MATERIALIZED VIEW "+view); err != nil {
t.Fatal(err)
}
}
}

func analyticsSnapshot(t *testing.T, ctx context.Context, tx pgx.Tx) map[string]string {
t.Helper()
refreshRetainedViews(t, ctx, tx)
result := map[string]string{}
for _, view := range retainedViews {
var rows string
if err := tx.QueryRow(ctx, "SELECT COALESCE(jsonb_agg(r ORDER BY r::text),'[]')::text FROM (SELECT to_jsonb(v) r FROM "+view+" v) s").Scan(&rows); err != nil {
t.Fatal(err)
Expand All @@ -34,9 +41,8 @@ func analyticsSnapshot(t *testing.T, ctx context.Context, tx pgx.Tx) map[string]

func TestAnalyticsRetentionConcurrentPostgres(t *testing.T) {
ctx, tx := retentionTx(t)
analyticsTables(t, ctx, tx)
applyStatsMigration(t, ctx, tx, "039_analytics_retention.sql")
if _, err := tx.Exec(ctx, `INSERT INTO packets (packet_hash,last_heard_at) VALUES ('\x01',NOW()-interval '5 days'); INSERT INTO observers(id) VALUES ('00000000-0000-0000-0000-000000000001')`); err != nil {
applyBaseline(t, ctx, tx)
if _, err := tx.Exec(ctx, `INSERT INTO packets (packet_hash,payload_type,payload_version,route_type,raw_payload,raw_header,first_heard_at,last_heard_at) VALUES ('\x01',4,0,1,'\x00','\x00',NOW()-interval '5 days',NOW()-interval '5 days'); INSERT INTO observers(id,public_key) VALUES ('00000000-0000-0000-0000-000000000001','\x01')`); err != nil {
t.Fatal(err)
}
var schema string
Expand Down Expand Up @@ -70,7 +76,7 @@ func TestAnalyticsRetentionConcurrentPostgres(t *testing.T) {
}
defer ingest.Rollback(context.Background())
// An in-flight FK insert holds KEY SHARE on the packet until it commits.
if _, err = ingest.Exec(ctx, `INSERT INTO packet_observations(packet_hash,observer_id,iata,heard_at) VALUES ('\x01','00000000-0000-0000-0000-000000000001','YVR',NOW()-interval '5 days')`); err != nil {
if _, err = ingest.Exec(ctx, `INSERT INTO packet_observations(packet_hash,observer_id,iata,heard_at,path_length_byte,hash_size,hop_count) VALUES ('\x01','00000000-0000-0000-0000-000000000001','YVR',NOW()-interval '5 days',0,1,0)`); err != nil {
t.Fatal(err)
}
q := sqlc.New(conn)
Expand Down Expand Up @@ -114,33 +120,11 @@ func TestAnalyticsRetentionConcurrentPostgres(t *testing.T) {
}
}

// Minimal real tables keep this regression runnable against an empty CI database.
func analyticsTables(t *testing.T, ctx context.Context, tx pgx.Tx) {
t.Helper()
isolateStatsSchema(t, ctx, tx)
_, err := tx.Exec(ctx, `
CREATE TABLE packets (packet_hash bytea PRIMARY KEY, payload_type smallint, payload_version smallint,
route_type smallint, raw_payload bytea, raw_header bytea, origin_pubkey bytea,
first_heard_at timestamptz, last_heard_at timestamptz);
CREATE INDEX ON packets(last_heard_at);
CREATE TABLE packet_observations (id bigserial PRIMARY KEY, packet_hash bytea REFERENCES packets ON DELETE CASCADE,
observer_id uuid NOT NULL, iata char(3) NOT NULL, heard_at timestamptz NOT NULL,
path_length_byte smallint, hash_size smallint, hop_count smallint, path_bytes bytea,
snr real, rssi smallint, airtime_ms real, payload_type smallint, UNIQUE(packet_hash,observer_id));
CREATE TABLE observers (id uuid PRIMARY KEY, public_key bytea, display_name text, observer_type text);
CREATE TABLE nodes (id uuid PRIMARY KEY, public_key bytea UNIQUE, node_type smallint, name text);
CREATE TABLE channel_messages (id bigserial PRIMARY KEY, channel_id integer, packet_hash bytea UNIQUE REFERENCES packets ON DELETE CASCADE,
sender_name text, content text, sent_at timestamptz);
`)
if err != nil {
t.Fatal(err)
}
}

func TestAnalyticsRetentionPostgres(t *testing.T) {
ctx, tx := retentionTx(t)
analyticsTables(t, ctx, tx)
applyBaseline(t, ctx, tx)
_, err := tx.Exec(ctx, `
INSERT INTO channels (id,channel_hash) VALUES (1,'\x01');
INSERT INTO observers (id,public_key,display_name,observer_type) SELECT ('00000000-0000-0000-0000-00000000000'||i)::uuid,int4send(i),CASE WHEN i=1 THEN NULL ELSE 'Observer '||i END,CASE WHEN i=1 THEN NULL ELSE 'test' END FROM generate_series(1,3) i;
INSERT INTO nodes (id,public_key,node_type,name) VALUES ('00000000-0000-0000-0000-000000000002','\x02',2,'Repeater');
INSERT INTO packets (packet_hash,payload_type,payload_version,route_type,raw_payload,raw_header,origin_pubkey,first_heard_at,last_heard_at)
Expand All @@ -152,22 +136,17 @@ func TestAnalyticsRetentionPostgres(t *testing.T) {
FROM packets CROSS JOIN (VALUES ('YVR',1),('YVR',2),('YYZ',3)) v(iata,n);
INSERT INTO channel_messages (channel_id,packet_hash,sender_name,content,sent_at)
SELECT 1,packet_hash,'Sender','body that must expire',first_heard_at FROM packets;
CREATE MATERIALIZED VIEW mv_hourly_iata_stats AS SELECT iata,date_trunc('hour',heard_at) AS hour,count(*) AS observation_count,count(DISTINCT packet_hash) AS unique_packets,count(DISTINCT observer_id) AS active_observers FROM packet_observations GROUP BY 1,2;
`)
if err != nil {
t.Fatal(err)
}
for _, migration := range []string{"016_mv_payload_breakdown.sql", "017_mv_top_observers.sql", "018_mv_top_talkers.sql", "019_mv_top_advertisers.sql", "021_mv_top_advertisers_route_type.sql", "032_mv_observer_activity.sql", "035_mv_signal_stats.sql", "036_mv_path_stats.sql"} {
applyStatsMigration(t, ctx, tx, migration)
}
applyStatsMigration(t, ctx, tx, "039_analytics_retention.sql")
for _, name := range []string{"idx_mv_signal_stats_hourly_hour", "idx_mv_path_stats_hourly_hour"} {
var definition string
if err := tx.QueryRow(ctx, "SELECT pg_get_indexdef(to_regclass($1))", name).Scan(&definition); err != nil || !strings.Contains(definition, "(hour)") {
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.
refreshRetainedViews(t, ctx, tx)
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)
Expand All @@ -176,14 +155,14 @@ func TestAnalyticsRetentionPostgres(t *testing.T) {
}
// 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)
t.Fatalf("signal stats: %+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)
t.Fatalf("path stats: %+v", paths)
}
before := analyticsSnapshot(t, ctx, tx)
q := sqlc.New(tx)
Expand Down Expand Up @@ -216,18 +195,17 @@ func TestAnalyticsRetentionPostgres(t *testing.T) {
t.Errorf("%s lost or double-counted history after raw expiry: before=%s after=%s", view, before[view], after[view])
}
}
// Retry and migration journal retry preserve the exact summaries.
// A retry preserves the exact summaries.
if n, err := q.DeleteOldPackets(ctx, sqlc.DeleteOldPacketsParams{Cutoff: pgtype.Timestamptz{Time: cutoff, Valid: true}, BatchSize: 1000}); err != nil || n != 0 {
t.Fatalf("retry: %d %v", n, err)
}
applyStatsMigration(t, ctx, tx, "039_analytics_retention.sql")
for view, rows := range analyticsSnapshot(t, ctx, tx) {
if rows != before[view] {
t.Errorf("retry changed %s", view)
}
}
// Later receptions on a still-live packet must join its old bucket before expiry.
if _, err := tx.Exec(ctx, `INSERT INTO observers (id) VALUES ('00000000-0000-0000-0000-000000000004'); INSERT INTO packet_observations (packet_hash,observer_id,iata,heard_at,path_length_byte,hash_size,hop_count,snr,rssi,payload_type)
if _, err := tx.Exec(ctx, `INSERT INTO observers (id,public_key) VALUES ('00000000-0000-0000-0000-000000000004','\x04'); INSERT INTO packet_observations (packet_hash,observer_id,iata,heard_at,path_length_byte,hash_size,hop_count,snr,rssi,payload_type)
SELECT packet_hash,'00000000-0000-0000-0000-000000000004','YVR',first_heard_at,0,1,0,'NaN',-80,NULL FROM packets;
UPDATE packets SET last_heard_at=first_heard_at`); err != nil {
t.Fatal(err)
Expand Down Expand Up @@ -282,13 +260,13 @@ func TestAnalyticsRetentionPostgres(t *testing.T) {
}
// More than one populated cohort, with shared hourly keys across the boundary.
if _, err := tx.Exec(ctx, `
INSERT INTO observers (id) VALUES ('00000000-0000-0000-0000-000000000001');
INSERT INTO observers (id,public_key) VALUES ('00000000-0000-0000-0000-000000000001','\x01');
INSERT INTO nodes (id,public_key,node_type) VALUES ('00000000-0000-0000-0000-000000000002','\x02',2);
INSERT INTO packets (packet_hash,payload_type,route_type,origin_pubkey,first_heard_at,last_heard_at,raw_payload)
SELECT int4send(i),4,1,'\x02',date_trunc('hour',NOW())-interval '5 days',date_trunc('hour',NOW())-interval '5 days',decode(repeat('ab',1024),'hex') FROM generate_series(10000,11000) i;
INSERT INTO packets (packet_hash,payload_type,payload_version,route_type,origin_pubkey,first_heard_at,last_heard_at,raw_payload,raw_header)
SELECT int4send(i),4,0,1,'\x02',date_trunc('hour',NOW())-interval '5 days',date_trunc('hour',NOW())-interval '5 days',decode(repeat('ab',1024),'hex'),'\x00' FROM generate_series(10000,11000) i;
INSERT INTO packet_observations (packet_hash,observer_id,iata,heard_at,path_length_byte,hash_size,hop_count,payload_type,snr,rssi)
SELECT packet_hash,'00000000-0000-0000-0000-000000000001','YVR',first_heard_at,0,1,0,4,0,-100 FROM packets;
INSERT INTO channel_messages (packet_hash,sender_name,sent_at) SELECT packet_hash,'Sender',first_heard_at FROM packets;
INSERT INTO channel_messages (channel_id,packet_hash,sender_name,sent_at) SELECT 1,packet_hash,'Sender',first_heard_at FROM packets;
`); err != nil {
t.Fatal(err)
}
Expand Down
119 changes: 119 additions & 0 deletions db/baseline_integration_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,119 @@
// Copyright 2026 Beacon Contributors
// SPDX-License-Identifier: AGPL-3.0-or-later

package db

import (
"context"
"errors"
"os"
"strings"
"testing"
"time"

"github.com/google/uuid"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
)

// Views cannot reference temporary tables. Keep fixtures and real migration DDL
// in a unique schema inside the caller's rolled-back transaction instead.
func isolateStatsSchema(t *testing.T, ctx context.Context, tx pgx.Tx) {
t.Helper()
schema := pgx.Identifier{"stats_test_" + uuid.NewString()}.Sanitize()
if _, err := tx.Exec(ctx, "CREATE SCHEMA "+schema+"; SET LOCAL search_path TO "+schema); err != nil {
t.Fatal(err)
}
}

// applyBaseline builds the real schema in an isolated schema, plus the IATAs fixtures use.
func applyBaseline(t *testing.T, ctx context.Context, tx pgx.Tx) {
t.Helper()
isolateStatsSchema(t, ctx, tx)
ddl, err := migrationFiles.ReadFile("migrations/" + baselineMigration)
if err != nil {
t.Fatal(err)
}
if _, err := tx.Exec(ctx, string(ddl)); err != nil {
t.Fatal(err)
}
if _, err := tx.Exec(ctx, "INSERT INTO iata_codes (iata) VALUES ('YVR'),('YYJ'),('YYZ'),('YOW')"); err != nil {
t.Fatal(err)
}
}

// schemaPool returns a pool whose connections all use a fresh, dropped-on-cleanup schema.
func schemaPool(t *testing.T) (context.Context, *pgxpool.Pool) {
t.Helper()
dsn := os.Getenv("BEACON_TEST_POSTGRES_DSN")
if dsn == "" {
t.Skip("set BEACON_TEST_POSTGRES_DSN for PostgreSQL regression")
}
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute)
t.Cleanup(cancel)
setup, err := pgx.Connect(ctx, dsn)
if err != nil {
t.Fatal(err)
}
schema := "migrate_test_" + strings.ReplaceAll(uuid.NewString(), "-", "")
quoted := pgx.Identifier{schema}.Sanitize()
if _, err := setup.Exec(ctx, "CREATE SCHEMA "+quoted); err != nil {
t.Fatal(err)
}
t.Cleanup(func() {
if _, err := setup.Exec(context.Background(), "DROP SCHEMA "+quoted+" CASCADE"); err != nil {
t.Error(err)
}
setup.Close(context.Background())
})
cfg, err := pgxpool.ParseConfig(dsn)
if err != nil {
t.Fatal(err)
}
cfg.ConnConfig.RuntimeParams["search_path"] = schema
cfg.MaxConns = 4
pool, err := pgxpool.NewWithConfig(ctx, cfg)
if err != nil {
t.Fatal(err)
}
t.Cleanup(pool.Close)
return ctx, pool
}

func TestRunMigrationsBaselinePostgres(t *testing.T) {
ctx, pool := schemaPool(t)
for range 2 {
if err := RunMigrations(ctx, pool); err != nil {
t.Fatal(err)
}
}
rows, err := pool.Query(ctx, "SELECT filename FROM schema_migrations ORDER BY filename")
if err != nil {
t.Fatal(err)
}
ledger, err := pgx.CollectRows(rows, pgx.RowTo[string])
if err != nil || len(ledger) != 1 || ledger[0] != baselineMigration {
t.Fatalf("ledger %v, %v", ledger, err)
}
// Stats refresh uses CONCURRENTLY, which needs populated views.
if _, err := pool.Exec(ctx, "REFRESH MATERIALIZED VIEW CONCURRENTLY mv_hourly_iata_stats"); err != nil {
t.Fatal(err)
}
}

func TestRunMigrationsRefusesPreBaselinePostgres(t *testing.T) {
for name, setup := range map[string]string{
"1.x ledger": "CREATE TABLE schema_migrations (filename text PRIMARY KEY); INSERT INTO schema_migrations VALUES ('001_initial_schema.sql')",
"pre-ledger 1.x": "CREATE TABLE packets (packet_hash bytea PRIMARY KEY)",
} {
t.Run(name, func(t *testing.T) {
ctx, pool := schemaPool(t)
if _, err := pool.Exec(ctx, setup); err != nil {
t.Fatal(err)
}
if err := RunMigrations(ctx, pool); !errors.Is(err, errPreBaseline) {
t.Fatalf("got %v, want errPreBaseline", err)
}
})
}
}
6 changes: 2 additions & 4 deletions db/channel_cursor_integration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -56,10 +56,8 @@ func TestChannelCursorPostgres(t *testing.T) {
t.Fatal(err)
}
}
exec(`CREATE TEMP TABLE channels (LIKE public.channels INCLUDING ALL) ON COMMIT DROP;
CREATE TEMP TABLE channel_iatas (LIKE public.channel_iatas INCLUDING ALL) ON COMMIT DROP;
CREATE INDEX ON channels(last_seen DESC,id DESC);
INSERT INTO channels(id,channel_hash,key_fingerprint,last_seen) VALUES
applyBaseline(t, ctx, tx)
exec(`INSERT INTO channels(id,channel_hash,key_fingerprint,last_seen) VALUES
(1,'\xaa','\x01','2026-09-08 12:00:00+00'),(2,'\xaa','\x02','2026-09-08 12:00:00+00'),
(3,'\xaa','\x03','2026-09-08 12:00:00+00'),(4,'\xbb','\x04','2026-09-08 11:59:59+00'),
(5,'\xcc','\x05','2026-09-08 12:00:00.000321+00'),
Expand Down
Loading
Loading