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 README.md
Original file line number Diff line number Diff line change
Expand Up @@ -242,7 +242,7 @@ telemetry:

# Packet and observation retention.
packets:
retention: 720h # how long to keep packets and observations (default: 30 days)
retention: 168h # packets, observations, channel messages (default: 7 days)

# Presence write coalescing.
# Observer last_seen and packet last_heard_at bumps are batched in memory and
Expand All @@ -260,7 +260,7 @@ websocket:
nodes:
mark_foreign: false # optional indication for repeaters outside configured IATA borders
stale_threshold: 24h # mark a node "stale" in the API after this long unseen (default: 24h)
delete_after: 720h # delete a node entirely after this long unseen (default: 30 days, same default as packets.retention)
delete_after: 720h # delete a node entirely after this long unseen (default: 30 days)
clock_drift_threshold: 5m # |device clock - server clock| above which clockOutOfSync=true for a repeater/room server (default: 5m)

# Optional observer age-out (disabled by default; set e.g. 720h to opt in).
Expand Down
5 changes: 3 additions & 2 deletions config.yaml.example
Original file line number Diff line number Diff line change
Expand Up @@ -118,8 +118,9 @@ telemetry:
#backup:
# enabled: false

# Also covers observations and channel messages.
packets:
retention: 720h # 30 days
retention: 168h # 7 days

# Known-route retention. Routes are distilled packet history (the path a packet
# took, hop by hop), so they may outlive packets.retention; there is no required
Expand Down Expand Up @@ -155,7 +156,7 @@ websocket:
# # uses the union of configured iatas.*.borderFile polygons;
# # requires at least one valid border; no packet filtering
# stale_threshold: 24h # default: 24h
# delete_after: 720h # default: 30 days, same default as packets.retention
# delete_after: 720h # default: 30 days
# clock_drift_threshold: 5m # default: 5m

# Optional observer age-out. Omitted or nonpositive leaves observers intact.
Expand Down
14 changes: 14 additions & 0 deletions db/migrations/037_autovacuum_tuning.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
-- Copyright 2026 Beacon Contributors
-- SPDX-License-Identifier: AGPL-3.0-or-later

-- Default 20% trigger never fired on known_routes, bloating its indexes.
ALTER TABLE known_routes SET (
autovacuum_vacuum_scale_factor = 0.02,
autovacuum_analyze_scale_factor = 0.02,
autovacuum_vacuum_cost_limit = 1000
);
ALTER TABLE packets SET (autovacuum_vacuum_scale_factor = 0.05, autovacuum_vacuum_cost_limit = 1000);
ALTER TABLE packet_observations SET (autovacuum_vacuum_scale_factor = 0.05, autovacuum_vacuum_cost_limit = 1000);
ALTER TABLE nodes SET (autovacuum_vacuum_scale_factor = 0.05);
ALTER TABLE node_iatas SET (autovacuum_vacuum_scale_factor = 0.05);
ALTER TABLE node_short_ids SET (autovacuum_vacuum_scale_factor = 0.05);
10 changes: 9 additions & 1 deletion db/packets.go
Original file line number Diff line number Diff line change
Expand Up @@ -650,6 +650,14 @@ func (s *Store) GetPacketObservationCount(ctx context.Context, packetHash []byte
return s.q.GetPacketObservationCount(ctx, packetHash)
}

// Small because each packet cascades to its observations.
const packetDeleteBatch = 1000

func (s *Store) DeleteOldPackets(ctx context.Context, cutoff time.Time) error {
return s.q.DeleteOldPackets(ctx, pgtype.Timestamptz{Time: cutoff, Valid: true})
return deleteInBatches(ctx, packetDeleteBatch, func(ctx context.Context, n int32) (int64, error) {
return s.q.DeleteOldPackets(ctx, sqlc.DeleteOldPacketsParams{
Cutoff: pgtype.Timestamptz{Time: cutoff, Valid: true},
BatchSize: n,
})
})
}
20 changes: 20 additions & 0 deletions db/packets_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -642,3 +642,23 @@ func TestListPackets_UnfilteredKeepsGlobalQuery(t *testing.T) {
t.Fatalf("unexpected error: %v", err)
}
}

func TestDeleteOldPackets_LoopsUntilShortBatch(t *testing.T) {
ctrl := gomock.NewController(t)
mock := mockdb.NewMockQuerier(ctrl)
store := &Store{q: mock}
cutoff := time.Date(2026, 9, 18, 0, 0, 0, 0, time.UTC)

byCutoff := gomock.Cond(func(p sqlc.DeleteOldPacketsParams) bool {
return p.Cutoff.Time.Equal(cutoff) && p.BatchSize == packetDeleteBatch
})
gomock.InOrder(
mock.EXPECT().DeleteOldPackets(gomock.Any(), byCutoff).Return(int64(packetDeleteBatch), nil),
mock.EXPECT().DeleteOldPackets(gomock.Any(), byCutoff).Return(int64(packetDeleteBatch), nil),
mock.EXPECT().DeleteOldPackets(gomock.Any(), byCutoff).Return(int64(7), nil),
)

if err := store.DeleteOldPackets(context.Background(), cutoff); err != nil {
t.Fatal(err)
}
}
36 changes: 26 additions & 10 deletions db/queries/queries.sql
Original file line number Diff line number Diff line change
Expand Up @@ -658,10 +658,18 @@ ORDER BY po.id ASC
LIMIT $6;


-- name: DeleteOldPackets :exec
-- Deletes packets and their observations older than the given cutoff.
-- packet_observations cascade-delete via FK.
DELETE FROM packets WHERE last_heard_at < $1;
-- name: DeleteOldPackets :execrows
-- One batch of expired packets; observations and channel messages cascade.
WITH expired AS (
SELECT ep.packet_hash
FROM packets ep
WHERE ep.last_heard_at < @cutoff
ORDER BY ep.last_heard_at
LIMIT @batch_size
FOR UPDATE OF ep SKIP LOCKED
)
DELETE FROM packets p USING expired e
WHERE p.packet_hash = e.packet_hash;

-- name: DeleteOldNodes :exec
-- Deletes nodes not seen since the given cutoff. node_iatas and node_neighbors cascade-
Expand All @@ -675,12 +683,20 @@ DELETE FROM nodes
WHERE last_seen < $1
AND id NOT IN (SELECT owner_node_id FROM observer_owners WHERE owner_node_id IS NOT NULL);

-- name: DeleteOldRoutes :exec
-- Deletes routes not observed since the retention cutoff ($1), and rarely-observed
-- routes (observation_count < $2) not observed since the grace cutoff ($3).
DELETE FROM known_routes
WHERE last_seen < $1
OR (observation_count < $2 AND last_seen < $3);
-- name: DeleteOldRoutes :execrows
-- One batch of routes past retention, or past grace with too few observations.
-- GREATEST keeps the scan on idx_known_routes_last_seen.
WITH expired AS (
SELECT r.iata, r.path_key
FROM known_routes r
WHERE r.last_seen < GREATEST(@retention_cutoff::timestamptz, @grace_cutoff::timestamptz)
AND (r.last_seen < @retention_cutoff OR
(r.observation_count < @min_observations AND r.last_seen < @grace_cutoff))
LIMIT @batch_size
FOR UPDATE OF r SKIP LOCKED
)
DELETE FROM known_routes kr USING expired e
WHERE kr.iata = e.iata AND kr.path_key = e.path_key;

-- name: DeleteOldChannelIATAs :exec
-- Keeps the channel IATA filter in step with packet retention.
Expand Down
111 changes: 111 additions & 0 deletions db/retention_integration_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,111 @@
// Copyright 2026 Beacon Contributors
// SPDX-License-Identifier: AGPL-3.0-or-later

package db

import (
"context"
"os"
"testing"
"time"

sqlc "github.com/MeshCore-Beacon/beacon-server/db/sqlc"
"github.com/jackc/pgx/v5"
)

func retentionTx(t *testing.T) (context.Context, pgx.Tx) {
t.Helper()
dsn := os.Getenv("BEACON_TEST_POSTGRES_DSN")
if dsn == "" {
t.Skip("set BEACON_TEST_POSTGRES_DSN for the PostgreSQL regression test")
}
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute)
t.Cleanup(cancel)
conn, err := pgx.Connect(ctx, dsn)
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { conn.Close(context.Background()) })
tx, err := conn.Begin(ctx)
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { tx.Rollback(context.Background()) })
return ctx, tx
}

func countRows(t *testing.T, ctx context.Context, tx pgx.Tx, sql string) int {
t.Helper()
var n int
if err := tx.QueryRow(ctx, sql).Scan(&n); err != nil {
t.Fatal(err)
}
return n
}

// More than one batch of old packets; observations go with them.
func TestDeleteOldPacketsBatchesPostgres(t *testing.T) {
ctx, tx := retentionTx(t)
_, err := tx.Exec(ctx, `
CREATE TEMP TABLE packets (LIKE public.packets INCLUDING ALL) ON COMMIT DROP;
CREATE TEMP TABLE packet_observations (LIKE public.packet_observations INCLUDING ALL) ON COMMIT DROP;
ALTER TABLE packet_observations ADD FOREIGN KEY (packet_hash) REFERENCES packets(packet_hash) ON DELETE CASCADE;
INSERT INTO packets (packet_hash, payload_type, payload_version, route_type, raw_payload, raw_header, first_heard_at, last_heard_at)
SELECT int4send(i), 4, 0, 1, '\x00', '\x00',
CASE WHEN i <= 2500 THEN '2026-01-01'::timestamptz ELSE '2026-02-01'::timestamptz END,
CASE WHEN i <= 2500 THEN '2026-01-01'::timestamptz ELSE '2026-02-01'::timestamptz END
FROM generate_series(1, 2510) i;
INSERT INTO packet_observations (packet_hash, observer_id, iata, heard_at, path_length_byte, hash_size, hop_count)
SELECT packet_hash, '00000000-0000-0000-0000-000000000001', 'YVR', last_heard_at, 0, 1, 0 FROM packets;`)
if err != nil {
t.Fatal(err)
}
store := &Store{q: sqlc.New(tx)}
if err := store.DeleteOldPackets(ctx, time.Date(2026, 1, 15, 0, 0, 0, 0, time.UTC)); err != nil {
t.Fatal(err)
}
if n := countRows(t, ctx, tx, `SELECT count(*) FROM packets`); n != 10 {
t.Errorf("packets left = %d, want 10", n)
}
if n := countRows(t, ctx, tx, `SELECT count(*) FROM packet_observations`); n != 10 {
t.Errorf("observations left = %d, want 10", n)
}
}

// Covers grace set longer than retention too.
func TestDeleteOldRoutesBatchesPostgres(t *testing.T) {
ctx, tx := retentionTx(t)
// 1-10500 past retention, 10501-10600 rare, 10601-10700 well observed, rest fresh.
_, err := tx.Exec(ctx, `
CREATE TEMP TABLE known_routes (LIKE public.known_routes INCLUDING ALL) ON COMMIT DROP;
INSERT INTO known_routes (id, path_key, node_ids, hash_prefix, iata, hop_count, last_seen, observation_count)
SELECT i, decode(md5(i::text), 'hex'), ARRAY[md5(i::text)::uuid], ARRAY['\x01'::bytea], 'YYZ', 2,
CASE WHEN i <= 10500 THEN '2026-01-01'::timestamptz
WHEN i <= 10700 THEN '2026-01-20'::timestamptz
ELSE '2026-02-01'::timestamptz END,
CASE WHEN i BETWEEN 10501 AND 10600 THEN 1 ELSE 10 END
FROM generate_series(1, 10710) i;`)
if err != nil {
t.Fatal(err)
}
store := &Store{q: sqlc.New(tx)}
retention := time.Date(2026, 1, 10, 0, 0, 0, 0, time.UTC)
grace := time.Date(2026, 1, 25, 0, 0, 0, 0, time.UTC)
if err := store.DeleteOldRoutes(ctx, retention, 3, grace); err != nil {
t.Fatal(err)
}
if n := countRows(t, ctx, tx, `SELECT count(*) FROM known_routes`); n != 110 {
t.Errorf("routes left = %d, want 110 (well-observed + fresh)", n)
}
if n := countRows(t, ctx, tx, `SELECT count(*) FROM known_routes WHERE observation_count = 1`); n != 0 {
t.Errorf("rare routes past grace left = %d, want 0", n)
}

// Grace longer than retention: retention alone decides.
if err := store.DeleteOldRoutes(ctx, time.Date(2026, 1, 25, 0, 0, 0, 0, time.UTC), 3, retention); err != nil {
t.Fatal(err)
}
if n := countRows(t, ctx, tx, `SELECT count(*) FROM known_routes`); n != 10 {
t.Errorf("routes left = %d, want 10 fresh", n)
}
}
13 changes: 9 additions & 4 deletions db/routes.go
Original file line number Diff line number Diff line change
Expand Up @@ -292,13 +292,18 @@ func (s *Store) ReconfirmRoutes(ctx context.Context, batchSize int32) error {
return s.q.ReconfirmRoutes(ctx, batchSize)
}

const routeDeleteBatch = 10000

// DeleteOldRoutes prunes routes per the retention rule: unconditionally past
// retentionCutoff, and past graceCutoff when observed fewer than minObservations times.
func (s *Store) DeleteOldRoutes(ctx context.Context, retentionCutoff time.Time, minObservations int64, graceCutoff time.Time) error {
return s.q.DeleteOldRoutes(ctx, sqlc.DeleteOldRoutesParams{
LastSeen: pgtype.Timestamptz{Time: retentionCutoff, Valid: true},
ObservationCount: minObservations,
LastSeen_2: pgtype.Timestamptz{Time: graceCutoff, Valid: true},
return deleteInBatches(ctx, routeDeleteBatch, func(ctx context.Context, n int32) (int64, error) {
return s.q.DeleteOldRoutes(ctx, sqlc.DeleteOldRoutesParams{
RetentionCutoff: pgtype.Timestamptz{Time: retentionCutoff, Valid: true},
GraceCutoff: pgtype.Timestamptz{Time: graceCutoff, Valid: true},
MinObservations: minObservations,
BatchSize: n,
})
})
}

Expand Down
5 changes: 3 additions & 2 deletions db/routes_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -60,8 +60,9 @@ func TestDeleteOldRoutes_PassesCutoffs(t *testing.T) {
grace := time.Date(2026, 7, 30, 0, 0, 0, 0, time.UTC)

mock.EXPECT().DeleteOldRoutes(gomock.Any(), gomock.Cond(func(p sqlc.DeleteOldRoutesParams) bool {
return p.LastSeen.Time.Equal(retention) && p.ObservationCount == 3 && p.LastSeen_2.Time.Equal(grace)
})).Return(nil)
return p.RetentionCutoff.Time.Equal(retention) && p.MinObservations == 3 && p.GraceCutoff.Time.Equal(grace) &&
p.BatchSize == routeDeleteBatch
})).Return(int64(0), nil)

if err := store.DeleteOldRoutes(context.Background(), retention, 3, grace); err != nil {
t.Fatal(err)
Expand Down
20 changes: 11 additions & 9 deletions db/sqlc/mock/querier.go

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

11 changes: 5 additions & 6 deletions db/sqlc/querier.go

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

Loading
Loading