From a84112f7fbbe49fe29db01123dfc02700fc1eaeb Mon Sep 17 00:00:00 2001 From: MrAlders0n Date: Thu, 23 Jul 2026 13:10:28 -0400 Subject: [PATCH] fix(traces): review fixes for trace_iatas Refresh last_heard on duplicate observations too, capped at hourly, so steady traffic can't age a trace out of the filter (dedup key has no heard_at). Drop the unused last_heard index column so upserts stay HOT, and share one retention cutoff across the cleanup deletes. --- db/migrations/012_trace_iatas.sql | 2 +- db/queries/queries.sql | 4 +++- db/sqlc/querier.go | 1 + db/sqlc/queries.sql.go | 4 +++- internal/background/tasks.go | 9 ++++++--- internal/ingest/packet.go | 3 ++- 6 files changed, 16 insertions(+), 7 deletions(-) diff --git a/db/migrations/012_trace_iatas.sql b/db/migrations/012_trace_iatas.sql index f26e05d..ffc9840 100644 --- a/db/migrations/012_trace_iatas.sql +++ b/db/migrations/012_trace_iatas.sql @@ -8,7 +8,7 @@ CREATE TABLE trace_iatas ( PRIMARY KEY (trace_tag, iata) ); -CREATE INDEX idx_trace_iatas_iata ON trace_iatas(iata, last_heard DESC); +CREATE INDEX idx_trace_iatas_iata ON trace_iatas(iata); -- Seed from retained packets; parallelism off so the join spills to disk, not /dev/shm. SET max_parallel_workers_per_gather = 0; diff --git a/db/queries/queries.sql b/db/queries/queries.sql index 698ab3a..4f1c265 100644 --- a/db/queries/queries.sql +++ b/db/queries/queries.sql @@ -684,10 +684,12 @@ ON CONFLICT (channel_hash, iata) DO UPDATE SET WHERE EXCLUDED.last_heard > channel_iatas.last_heard + INTERVAL '1 hour'; -- name: UpsertTraceIATA :exec +-- Refreshes at most hourly so repeat hears don't churn the row. INSERT INTO trace_iatas (trace_tag, iata, last_heard) VALUES ($1, $2, $3) ON CONFLICT (trace_tag, iata) DO UPDATE SET - last_heard = GREATEST(trace_iatas.last_heard, EXCLUDED.last_heard); + last_heard = EXCLUDED.last_heard +WHERE EXCLUDED.last_heard > trace_iatas.last_heard + INTERVAL '1 hour'; -- name: ListChannels :many -- Channels ordered by last seen, optionally filtered by hash and/or IATAs diff --git a/db/sqlc/querier.go b/db/sqlc/querier.go index 3b158bb..fa93a88 100644 --- a/db/sqlc/querier.go +++ b/db/sqlc/querier.go @@ -236,6 +236,7 @@ type Querier interface { UpsertPacket(ctx context.Context, arg UpsertPacketParams) (UpsertPacketRow, error) UpsertRegion(ctx context.Context, arg UpsertRegionParams) (int32, error) UpsertRegionIATA(ctx context.Context, arg UpsertRegionIATAParams) error + // Refreshes at most hourly so repeat hears don't churn the row. UpsertTraceIATA(ctx context.Context, arg UpsertTraceIATAParams) error // ============================================================ // TRANSPORT CODES diff --git a/db/sqlc/queries.sql.go b/db/sqlc/queries.sql.go index 748207a..1698cbe 100644 --- a/db/sqlc/queries.sql.go +++ b/db/sqlc/queries.sql.go @@ -4072,7 +4072,8 @@ const upsertTraceIATA = `-- name: UpsertTraceIATA :exec INSERT INTO trace_iatas (trace_tag, iata, last_heard) VALUES ($1, $2, $3) ON CONFLICT (trace_tag, iata) DO UPDATE SET - last_heard = GREATEST(trace_iatas.last_heard, EXCLUDED.last_heard) + last_heard = EXCLUDED.last_heard +WHERE EXCLUDED.last_heard > trace_iatas.last_heard + INTERVAL '1 hour' ` type UpsertTraceIATAParams struct { @@ -4081,6 +4082,7 @@ type UpsertTraceIATAParams struct { LastHeard pgtype.Timestamptz `json:"last_heard"` } +// Refreshes at most hourly so repeat hears don't churn the row. func (q *Queries) UpsertTraceIATA(ctx context.Context, arg UpsertTraceIATAParams) error { _, err := q.db.Exec(ctx, upsertTraceIATA, arg.TraceTag, arg.Iata, arg.LastHeard) return err diff --git a/internal/background/tasks.go b/internal/background/tasks.go index 9855b80..8824304 100644 --- a/internal/background/tasks.go +++ b/internal/background/tasks.go @@ -41,13 +41,16 @@ func CleanupTask(store *db.Store, telemetryRetention, packetRetention, interval if err := store.DeleteOldTelemetry(ctx, time.Now().Add(-telemetryRetention)); err != nil { return err } - if err := store.DeleteOldPackets(ctx, time.Now().Add(-packetRetention)); err != nil { + // One cutoff for all three so the IATA tables stay in step + // with the packets they mirror. + cutoff := time.Now().Add(-packetRetention) + if err := store.DeleteOldPackets(ctx, cutoff); err != nil { return err } - if err := store.DeleteOldChannelIATAs(ctx, time.Now().Add(-packetRetention)); err != nil { + if err := store.DeleteOldChannelIATAs(ctx, cutoff); err != nil { return err } - if err := store.DeleteOldTraceIATAs(ctx, time.Now().Add(-packetRetention)); err != nil { + if err := store.DeleteOldTraceIATAs(ctx, cutoff); err != nil { return err } return nil diff --git a/internal/ingest/packet.go b/internal/ingest/packet.go index 9ac0ed8..b51c99e 100644 --- a/internal/ingest/packet.go +++ b/internal/ingest/packet.go @@ -790,7 +790,8 @@ func (w *Worker) handlePacket(ctx context.Context, iata, pubkeyHex string, raw [ } } - if traceTag != nil && inserted { + // Runs on duplicate observations too; the upsert only writes when the row is >1h stale. + if traceTag != nil { if err := w.db.UpsertTraceIATA(ctx, traceTag, iata, heardAt); err != nil { log.Printf("ingest[%s]: db: upsert trace IATA failed from %s/%s: %v", w.cfg.BrokerName, iata, pubkeyHex, err) }