diff --git a/README.md b/README.md index becb8bc4..5418223c 100644 --- a/README.md +++ b/README.md @@ -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 @@ -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). diff --git a/config.yaml.example b/config.yaml.example index dd7bfac4..de12bef4 100644 --- a/config.yaml.example +++ b/config.yaml.example @@ -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 @@ -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. diff --git a/db/migrations/037_autovacuum_tuning.sql b/db/migrations/037_autovacuum_tuning.sql new file mode 100644 index 00000000..f8db967d --- /dev/null +++ b/db/migrations/037_autovacuum_tuning.sql @@ -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); diff --git a/db/packets.go b/db/packets.go index e6512694..8172c66f 100644 --- a/db/packets.go +++ b/db/packets.go @@ -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, + }) + }) } diff --git a/db/packets_test.go b/db/packets_test.go index 6541ba35..d93aac57 100644 --- a/db/packets_test.go +++ b/db/packets_test.go @@ -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) + } +} diff --git a/db/queries/queries.sql b/db/queries/queries.sql index 614c92c1..1f84bdf3 100644 --- a/db/queries/queries.sql +++ b/db/queries/queries.sql @@ -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- @@ -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. diff --git a/db/retention_integration_test.go b/db/retention_integration_test.go new file mode 100644 index 00000000..5e40dbd2 --- /dev/null +++ b/db/retention_integration_test.go @@ -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) + } +} diff --git a/db/routes.go b/db/routes.go index 047bfec6..a629bcbc 100644 --- a/db/routes.go +++ b/db/routes.go @@ -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, + }) }) } diff --git a/db/routes_test.go b/db/routes_test.go index 977b541a..034d57f4 100644 --- a/db/routes_test.go +++ b/db/routes_test.go @@ -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) diff --git a/db/sqlc/mock/querier.go b/db/sqlc/mock/querier.go index d941dc96..4efbfaa5 100644 --- a/db/sqlc/mock/querier.go +++ b/db/sqlc/mock/querier.go @@ -117,25 +117,27 @@ func (mr *MockQuerierMockRecorder) DeleteOldObservers(ctx, lastSeen any) *gomock } // DeleteOldPackets mocks base method. -func (m *MockQuerier) DeleteOldPackets(ctx context.Context, lastHeardAt pgtype.Timestamptz) error { +func (m *MockQuerier) DeleteOldPackets(ctx context.Context, arg db.DeleteOldPacketsParams) (int64, error) { m.ctrl.T.Helper() - ret := m.ctrl.Call(m, "DeleteOldPackets", ctx, lastHeardAt) - ret0, _ := ret[0].(error) - return ret0 + ret := m.ctrl.Call(m, "DeleteOldPackets", ctx, arg) + ret0, _ := ret[0].(int64) + ret1, _ := ret[1].(error) + return ret0, ret1 } // DeleteOldPackets indicates an expected call of DeleteOldPackets. -func (mr *MockQuerierMockRecorder) DeleteOldPackets(ctx, lastHeardAt any) *gomock.Call { +func (mr *MockQuerierMockRecorder) DeleteOldPackets(ctx, arg any) *gomock.Call { mr.mock.ctrl.T.Helper() - return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "DeleteOldPackets", reflect.TypeOf((*MockQuerier)(nil).DeleteOldPackets), ctx, lastHeardAt) + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "DeleteOldPackets", reflect.TypeOf((*MockQuerier)(nil).DeleteOldPackets), ctx, arg) } // DeleteOldRoutes mocks base method. -func (m *MockQuerier) DeleteOldRoutes(ctx context.Context, arg db.DeleteOldRoutesParams) error { +func (m *MockQuerier) DeleteOldRoutes(ctx context.Context, arg db.DeleteOldRoutesParams) (int64, error) { m.ctrl.T.Helper() ret := m.ctrl.Call(m, "DeleteOldRoutes", ctx, arg) - ret0, _ := ret[0].(error) - return ret0 + ret0, _ := ret[0].(int64) + ret1, _ := ret[1].(error) + return ret0, ret1 } // DeleteOldRoutes indicates an expected call of DeleteOldRoutes. diff --git a/db/sqlc/querier.go b/db/sqlc/querier.go index aa4bd874..d3913c42 100644 --- a/db/sqlc/querier.go +++ b/db/sqlc/querier.go @@ -29,12 +29,11 @@ type Querier interface { // Opt-in age-out: preserve retained history and manually recorded ownership. // Bound deletions per cleanup tick and skip observers being updated by ingest. DeleteOldObservers(ctx context.Context, lastSeen pgtype.Timestamptz) ([]uuid.UUID, error) - // Deletes packets and their observations older than the given cutoff. - // packet_observations cascade-delete via FK. - DeleteOldPackets(ctx context.Context, lastHeardAt pgtype.Timestamptz) error - // Deletes routes not observed since the retention cutoff ($1), and rarely-observed - // routes (observation_count < $2) not observed since the grace cutoff ($3). - DeleteOldRoutes(ctx context.Context, arg DeleteOldRoutesParams) error + // One batch of expired packets; observations and channel messages cascade. + DeleteOldPackets(ctx context.Context, arg DeleteOldPacketsParams) (int64, error) + // One batch of routes past retention, or past grace with too few observations. + // GREATEST keeps the scan on idx_known_routes_last_seen. + DeleteOldRoutes(ctx context.Context, arg DeleteOldRoutesParams) (int64, error) // Deletes telemetry rows older than the given cutoff. Called by the cleanup goroutine. DeleteOldTelemetry(ctx context.Context, reportedAt pgtype.Timestamptz) error // Keeps the trace IATA filter in step with packet retention. diff --git a/db/sqlc/queries.sql.go b/db/sqlc/queries.sql.go index 16e76c4f..d454af2d 100644 --- a/db/sqlc/queries.sql.go +++ b/db/sqlc/queries.sql.go @@ -124,34 +124,67 @@ func (q *Queries) DeleteOldObservers(ctx context.Context, lastSeen pgtype.Timest return items, nil } -const deleteOldPackets = `-- name: DeleteOldPackets :exec -DELETE FROM packets WHERE last_heard_at < $1 +const deleteOldPackets = `-- name: DeleteOldPackets :execrows +WITH expired AS ( + SELECT ep.packet_hash + FROM packets ep + WHERE ep.last_heard_at < $1 + ORDER BY ep.last_heard_at + LIMIT $2 + FOR UPDATE OF ep SKIP LOCKED +) +DELETE FROM packets p USING expired e +WHERE p.packet_hash = e.packet_hash ` -// Deletes packets and their observations older than the given cutoff. -// packet_observations cascade-delete via FK. -func (q *Queries) DeleteOldPackets(ctx context.Context, lastHeardAt pgtype.Timestamptz) error { - _, err := q.db.Exec(ctx, deleteOldPackets, lastHeardAt) - return err +type DeleteOldPacketsParams struct { + Cutoff pgtype.Timestamptz `json:"cutoff"` + BatchSize int32 `json:"batch_size"` } -const deleteOldRoutes = `-- name: DeleteOldRoutes :exec -DELETE FROM known_routes -WHERE last_seen < $1 - OR (observation_count < $2 AND last_seen < $3) +// One batch of expired packets; observations and channel messages cascade. +func (q *Queries) DeleteOldPackets(ctx context.Context, arg DeleteOldPacketsParams) (int64, error) { + result, err := q.db.Exec(ctx, deleteOldPackets, arg.Cutoff, arg.BatchSize) + if err != nil { + return 0, err + } + return result.RowsAffected(), nil +} + +const deleteOldRoutes = `-- name: DeleteOldRoutes :execrows +WITH expired AS ( + SELECT r.iata, r.path_key + FROM known_routes r + WHERE r.last_seen < GREATEST($1::timestamptz, $2::timestamptz) + AND (r.last_seen < $1 OR + (r.observation_count < $3 AND r.last_seen < $2)) + LIMIT $4 + 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 ` type DeleteOldRoutesParams struct { - LastSeen pgtype.Timestamptz `json:"last_seen"` - ObservationCount int64 `json:"observation_count"` - LastSeen_2 pgtype.Timestamptz `json:"last_seen_2"` -} - -// Deletes routes not observed since the retention cutoff ($1), and rarely-observed -// routes (observation_count < $2) not observed since the grace cutoff ($3). -func (q *Queries) DeleteOldRoutes(ctx context.Context, arg DeleteOldRoutesParams) error { - _, err := q.db.Exec(ctx, deleteOldRoutes, arg.LastSeen, arg.ObservationCount, arg.LastSeen_2) - return err + RetentionCutoff pgtype.Timestamptz `json:"retention_cutoff"` + GraceCutoff pgtype.Timestamptz `json:"grace_cutoff"` + MinObservations int64 `json:"min_observations"` + BatchSize int32 `json:"batch_size"` +} + +// One batch of routes past retention, or past grace with too few observations. +// GREATEST keeps the scan on idx_known_routes_last_seen. +func (q *Queries) DeleteOldRoutes(ctx context.Context, arg DeleteOldRoutesParams) (int64, error) { + result, err := q.db.Exec(ctx, deleteOldRoutes, + arg.RetentionCutoff, + arg.GraceCutoff, + arg.MinObservations, + arg.BatchSize, + ) + if err != nil { + return 0, err + } + return result.RowsAffected(), nil } const deleteOldTelemetry = `-- name: DeleteOldTelemetry :exec diff --git a/db/store.go b/db/store.go index a22b9600..20866d76 100644 --- a/db/store.go +++ b/db/store.go @@ -149,3 +149,19 @@ func toChannelMessage(id int64, packetHashHex string, channelHash []byte, sender ObservationCount: observationCount, } } + +// deleteInBatches repeats a delete until a batch comes back short. +func deleteInBatches(ctx context.Context, batchSize int32, del func(context.Context, int32) (int64, error)) error { + for { + n, err := del(ctx, batchSize) + if err != nil { + return err + } + if n < int64(batchSize) { + return nil + } + if err := ctx.Err(); err != nil { + return err + } + } +} diff --git a/db/store_test.go b/db/store_test.go index 8bfab391..4d0521be 100644 --- a/db/store_test.go +++ b/db/store_test.go @@ -179,3 +179,31 @@ func TestResolvePathHashes_Mapping(t *testing.T) { t.Errorf("expected Name %s, got %v", name, entries[0].Name) } } + +func TestDeleteInBatches_StopsOnError(t *testing.T) { + boom := errors.New("boom") + calls := 0 + err := deleteInBatches(context.Background(), 10, func(context.Context, int32) (int64, error) { + calls++ + if calls == 2 { + return 0, boom + } + return 10, nil + }) + if !errors.Is(err, boom) || calls != 2 { + t.Fatalf("err=%v calls=%d, want boom after 2 calls", err, calls) + } +} + +func TestDeleteInBatches_StopsWhenCancelled(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + calls := 0 + err := deleteInBatches(ctx, 10, func(context.Context, int32) (int64, error) { + calls++ + cancel() + return 10, nil + }) + if !errors.Is(err, context.Canceled) || calls != 1 { + t.Fatalf("err=%v calls=%d, want context.Canceled after 1 call", err, calls) + } +} diff --git a/internal/config/config.go b/internal/config/config.go index 61786364..5f532726 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -245,8 +245,7 @@ type WebSocketConfig struct { // PacketsConfig controls packet retention behaviour. type PacketsConfig struct { - // Retention is how long packet and observation rows are kept. - // Defaults to 720h (30 days) if not set. + // Retention covers packets, observations and channel messages. Defaults to 168h. Retention duration `yaml:"retention"` } @@ -276,8 +275,7 @@ type NodesConfig struct { // stale=true for it. Defaults to 24h if not set. StaleThreshold duration `yaml:"stale_threshold"` // DeleteAfter is how long since a node's last_seen before the cleanup job deletes the - // node entirely. Defaults to the same 30-day default as packets.retention if not set -- - // independently configurable from it, just the same starting point. + // node entirely. Defaults to 720h (30 days) if not set. DeleteAfter duration `yaml:"delete_after"` } @@ -449,7 +447,7 @@ func Resolve(cfg *Config) ResolvedConfig { r.TelemetryRetention = 28 * 24 * time.Hour } if r.PacketRetention == 0 { - r.PacketRetention = 30 * 24 * time.Hour + r.PacketRetention = 7 * 24 * time.Hour } if r.RouteRetention == 0 { r.RouteRetention = 14 * 24 * time.Hour @@ -488,8 +486,6 @@ func Resolve(cfg *Config) ResolvedConfig { r.NodeStaleThreshold = 24 * time.Hour } if r.NodeDeleteAfter == 0 { - // Same default as packets.retention (30 days) -- independently configurable, just - // the same starting point, not tied to whatever PacketRetention resolves to. r.NodeDeleteAfter = 30 * 24 * time.Hour } return r diff --git a/internal/config/config_test.go b/internal/config/config_test.go index 8ce2828c..4ba8b238 100644 --- a/internal/config/config_test.go +++ b/internal/config/config_test.go @@ -130,8 +130,8 @@ func TestResolve_Defaults(t *testing.T) { if r.TelemetryRetention != 28*24*time.Hour { t.Errorf("expected TelemetryRetention 672h, got %v", r.TelemetryRetention) } - if r.PacketRetention != 30*24*time.Hour { - t.Errorf("expected PacketRetention 720h, got %v", r.PacketRetention) + if r.PacketRetention != 7*24*time.Hour { + t.Errorf("expected PacketRetention 168h, got %v", r.PacketRetention) } if r.MaxConnsPerIP != 5 { t.Errorf("expected MaxConnsPerIP 5, got %d", r.MaxConnsPerIP) @@ -152,7 +152,7 @@ func TestResolve_Defaults(t *testing.T) { t.Errorf("expected NodeStaleThreshold 24h, got %v", r.NodeStaleThreshold) } if r.NodeDeleteAfter != 30*24*time.Hour { - t.Errorf("expected NodeDeleteAfter 720h (same default as PacketRetention), got %v", r.NodeDeleteAfter) + t.Errorf("expected NodeDeleteAfter 720h, got %v", r.NodeDeleteAfter) } if r.ObserverDeleteAfter != 0 { t.Errorf("observer deletion must be disabled by default, got %v", r.ObserverDeleteAfter)