From 70c5c389f9886cb8d8c9d9d34462c805cad6fc6b Mon Sep 17 00:00:00 2001 From: n30nex Date: Sat, 26 Sep 2026 03:46:29 -0400 Subject: [PATCH 1/3] feat(db): retain hourly analytics after raw packet expiry --- .github/workflows/ci.yml | 2 +- README.md | 15 + db/analytics_retention_integration_test.go | 261 +++++++++ db/migrations/039_analytics_retention.sql | 611 +++++++++++++++++++++ db/queries/queries.sql | 15 +- db/retention_integration_test.go | 5 +- db/sqlc/models.go | 158 ++++++ db/sqlc/querier.go | 2 +- db/sqlc/queries.sql.go | 24 +- sqlc.yaml | 8 + 10 files changed, 1067 insertions(+), 34 deletions(-) create mode 100644 db/analytics_retention_integration_test.go create mode 100644 db/migrations/039_analytics_retention.sql diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index beb7284f..7d374920 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -55,7 +55,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)Postgres$' -count=1 -v + run: go test ./db -run '^Test(Signal|Paths|PacketEndpointsResolveLive|AnalyticsRetention|AnalyticsRetentionConcurrent|DeleteOldPacketsBatches)Postgres$' -count=1 -v - name: Verify backup command against PostgreSQL 16 run: | diff --git a/README.md b/README.md index 5418223c..5243068a 100644 --- a/README.md +++ b/README.md @@ -533,3 +533,18 @@ log: `LOG_LEVEL` and `LOG_FORMAT` override file settings; empty settings use `info` and `text`. Invalid values prevent startup. Configuration-loading failures can use the bootstrap text logger before file settings are available; failures after initialization retain error severity at every supported level. Records include a component field. Ingest workers also include their broker name, and HTTP completion records include the validated client address, route, status and duration. Query strings and protocol hello payloads are excluded. Expected ingest skips and routine WebSocket lifecycle details are debug-level. Changing the application's format does not change Caddy/Apache access logs or their fail2ban configuration. Collect/rotate stderr through Docker or systemd. + +### Analytics retention + +Hourly traffic, payload, observer activity, talker, advertiser, Signal and Paths +summaries retain 30 days independently of `packets.retention`. Cleanup saves only +aggregates before deleting each packet batch, in the same transaction. Raw +packets, observations and message bodies still expire under packet retention. +The summaries use UTC hourly buckets and appear on the normal view-refresh cycle. + +Migration 039 starts from data still present; previously deleted history cannot +be reconstructed. Telemetry has its own retention setting. Packet drill-down, +sub-hour observer activity, exact observer comparison and current entity/scope +counts continue to describe retained raw data or current entities, rather than +claiming archived packet detail. Archived summaries expire without waiting for +new packet deletions. A failed archive leaves its entire raw batch intact. diff --git a/db/analytics_retention_integration_test.go b/db/analytics_retention_integration_test.go new file mode 100644 index 00000000..2e0466c4 --- /dev/null +++ b/db/analytics_retention_integration_test.go @@ -0,0 +1,261 @@ +// Copyright 2026 Beacon Contributors +// SPDX-License-Identifier: AGPL-3.0-or-later + +package db + +import ( + "context" + sqlc "github.com/MeshCore-Beacon/beacon-server/db/sqlc" + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgtype" + "os" + "testing" + "time" +) + +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 { + t.Helper() + result := map[string]string{} + for _, view := range retainedViews { + if _, err := tx.Exec(ctx, "REFRESH MATERIALIZED VIEW "+view); err != nil { + t.Fatal(err) + } + 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) + } + result[view] = rows + } + return result +} + +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 { + t.Fatal(err) + } + var schema string + if err := tx.QueryRow(ctx, "SELECT current_schema()").Scan(&schema); err != nil { + t.Fatal(err) + } + if err := tx.Commit(ctx); err != nil { + t.Fatal(err) + } + conn := tx.Conn() + quoted := pgx.Identifier{schema}.Sanitize() + t.Cleanup(func() { + if _, err := conn.Exec(context.Background(), "DROP SCHEMA "+quoted+" CASCADE"); err != nil { + t.Error(err) + } + }) + if _, err := conn.Exec(ctx, "SET search_path TO "+quoted); err != nil { + t.Fatal(err) + } + writer, err := pgx.Connect(ctx, os.Getenv("BEACON_TEST_POSTGRES_DSN")) + if err != nil { + t.Fatal(err) + } + defer writer.Close(context.Background()) + if _, err = writer.Exec(ctx, "SET search_path TO "+quoted); err != nil { + t.Fatal(err) + } + ingest, err := writer.Begin(ctx) + if err != nil { + t.Fatal(err) + } + 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 { + t.Fatal(err) + } + q := sqlc.New(conn) + params := sqlc.DeleteOldPacketsParams{Cutoff: pgtype.Timestamptz{Time: time.Now().Add(-72 * time.Hour), Valid: true}, BatchSize: 1000} + if n, err := q.DeleteOldPackets(ctx, params); err != nil || n != 0 { + t.Fatalf("must skip in-flight reception: %d %v", n, err) + } + if err := ingest.Commit(ctx); err != nil { + t.Fatal(err) + } + if n, err := q.DeleteOldPackets(ctx, params); err != nil || n != 1 { + t.Fatalf("retry after ingestion: %d %v", n, err) + } + var count int + if err := conn.QueryRow(ctx, "SELECT observation_count FROM analytics_hourly_iata_stats").Scan(&count); err != nil || count != 1 { + t.Fatalf("committed reception not archived: %d %v", count, err) + } +} + +// 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) + _, err := tx.Exec(ctx, ` + 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) + SELECT int4send(i),4,0,CASE WHEN i%2=0 THEN 1 ELSE 2 END,'\x00','\x00','\x02',date_trunc('hour',NOW())-interval '5 days', + CASE WHEN i=4 THEN NOW() ELSE date_trunc('hour',NOW())-interval '5 days' END FROM generate_series(1,4) i; + INSERT INTO packet_observations (packet_hash,observer_id,iata,heard_at,path_length_byte,hash_size,hop_count,snr,rssi,airtime_ms,payload_type) + SELECT packet_hash,('00000000-0000-0000-0000-00000000000'||n)::uuid,iata,first_heard_at+make_interval(secs=>n),0,1,0, + CASE WHEN n=1 THEN 0 ELSE NULL END,CASE WHEN n=1 THEN -100 ELSE NULL END,1.5,4 + 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") + before := analyticsSnapshot(t, ctx, tx) + q := sqlc.New(tx) + cutoff := time.Now().Add(-72 * time.Hour) + // Failure in one archive must roll back every archive write and the raw deletion. + if _, err := tx.Exec(ctx, `SAVEPOINT archive_failure; ALTER TABLE analytics_signal_stats_hourly ADD CONSTRAINT fail_archive CHECK(receptions < 0) NOT VALID`); err != nil { + t.Fatal(err) + } + if _, err := q.DeleteOldPackets(ctx, sqlc.DeleteOldPacketsParams{Cutoff: pgtype.Timestamptz{Time: cutoff, Valid: true}, BatchSize: 1}); err == nil { + t.Fatal("expected archive failure") + } + if _, err := tx.Exec(ctx, "ROLLBACK TO archive_failure"); err != nil { + t.Fatal(err) + } + if countRows(t, ctx, tx, "SELECT count(*) FROM packets") != 4 { + t.Fatal("archive failure deleted raw packets") + } + for i := 0; i < 3; i++ { + n, err := q.DeleteOldPackets(ctx, sqlc.DeleteOldPacketsParams{Cutoff: pgtype.Timestamptz{Time: cutoff, Valid: true}, BatchSize: 1}) + if err != nil || n != 1 { + t.Fatalf("delete batch: %d %v", n, err) + } + } + if countRows(t, ctx, tx, "SELECT count(*) FROM packets") != 1 || countRows(t, ctx, tx, "SELECT count(*) FROM channel_messages") != 1 { + t.Fatal("raw packets/messages did not expire") + } + after := analyticsSnapshot(t, ctx, tx) + for _, view := range retainedViews { + if after[view] != before[view] { + 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. + 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) + 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) + } + late := analyticsSnapshot(t, ctx, tx) + if err := (&Store{q: q}).DeleteOldPackets(ctx, cutoff); err != nil { + t.Fatal(err) + } + for view, rows := range analyticsSnapshot(t, ctx, tx) { + if rows != late[view] { + t.Errorf("late reception changed %s: %s != %s", view, rows, late[view]) + } + } + // Archived labels follow a current rename; no entity FK may erase retained history. + if _, err := tx.Exec(ctx, `UPDATE nodes SET name='Renamed'; UPDATE observers SET display_name=NULL,observer_type=NULL;`); err != nil { + t.Fatal(err) + } + analyticsSnapshot(t, ctx, tx) + if countRows(t, ctx, tx, "SELECT count(*) FROM mv_top_advertisers_by_iata WHERE name='Renamed'") != 2 { + t.Fatal("archive name did not follow rename") + } + if _, err := q.GetStatsTopObservers(ctx, sqlc.GetStatsTopObserversParams{Column1: pgtype.Interval{Microseconds: int64(30 * 24 * time.Hour / time.Microsecond), Valid: true}, Limit: 10}); err != nil { + t.Fatal(err) + } + if _, err := tx.Exec(ctx, "DELETE FROM nodes; DELETE FROM observers"); err != nil { + t.Fatal(err) + } + if countRows(t, ctx, tx, "SELECT count(*) FROM analytics_top_advertisers_by_iata") != 2 { + t.Fatal("entity deletion erased archive") + } + // Independent expiry runs even when there are no packets left to delete. + for _, view := range retainedViews { + table := "analytics_" + view[3:] + at := "bucket" + if view == "mv_hourly_iata_stats" || view == "mv_signal_stats_hourly" || view == "mv_path_stats_hourly" { + at = "hour" + } + if _, err := tx.Exec(ctx, "UPDATE "+table+" SET "+at+"="+at+"-interval '31 days'"); err != nil { + t.Fatal(err) + } + } + if err := (&Store{q: q}).DeleteOldPackets(ctx, cutoff); err != nil { + t.Fatal(err) + } + for _, view := range retainedViews { + if countRows(t, ctx, tx, "SELECT count(*) FROM analytics_"+view[3:]) != 0 { + t.Errorf("%s archive did not expire", view) + } + } + // 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 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 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; +`); err != nil { + t.Fatal(err) + } + populated := analyticsSnapshot(t, ctx, tx) + started := time.Now() + if err := (&Store{q: q}).DeleteOldPackets(ctx, cutoff); err != nil { + t.Fatal(err) + } + t.Logf("1001 packets with bodies and observations archived/deleted in %s", time.Since(started)) + for view, rows := range analyticsSnapshot(t, ctx, tx) { + if rows != populated[view] { + t.Errorf("batch boundary changed %s", view) + } + } + var archiveRows int + for _, view := range retainedViews { + archiveRows += countRows(t, ctx, tx, "SELECT count(*) FROM analytics_"+view[3:]) + } + if archiveRows > 12 { + t.Fatalf("expected compact rollups, got %d archive rows for 1001 packets", archiveRows) + } +} diff --git a/db/migrations/039_analytics_retention.sql b/db/migrations/039_analytics_retention.sql new file mode 100644 index 00000000..75d50078 --- /dev/null +++ b/db/migrations/039_analytics_retention.sql @@ -0,0 +1,611 @@ +-- Copyright 2026 Beacon Contributors +-- SPDX-License-Identifier: AGPL-3.0-or-later + +-- Archive only aggregates of deleted packet cohorts, never packet/message bodies +-- or raw paths. Live and archived cohorts are disjoint in every refresh snapshot. +-- Thirty days matches the longest analytics window; raw retention is independent. + + +CREATE TABLE IF NOT EXISTS analytics_hourly_iata_stats ( + iata char(3), + hour timestamptz, + observation_count bigint, + unique_packets bigint, + PRIMARY KEY (iata,hour) +); + +CREATE INDEX IF NOT EXISTS idx_analytics_hourly_iata_stats_time ON analytics_hourly_iata_stats (hour); + +CREATE TABLE IF NOT EXISTS analytics_payload_breakdown_by_iata ( + iata char(3), + payload_type smallint, + bucket timestamptz, + count bigint, + PRIMARY KEY (iata,payload_type,bucket) +); + +CREATE INDEX IF NOT EXISTS idx_analytics_payload_breakdown_by_iata_time ON analytics_payload_breakdown_by_iata (bucket); + +CREATE TABLE IF NOT EXISTS analytics_top_observers_by_iata ( + iata char(3), + observer_id uuid, + bucket timestamptz, + observation_count bigint, + display_name text, + observer_type text, + PRIMARY KEY (iata,observer_id,bucket) +); + +CREATE INDEX IF NOT EXISTS idx_analytics_top_observers_by_iata_time ON analytics_top_observers_by_iata (bucket); + +CREATE TABLE IF NOT EXISTS analytics_top_talkers_by_iata ( + iata char(3), + sender_name text, + bucket timestamptz, + message_count bigint, + last_sent timestamptz, + PRIMARY KEY (iata,sender_name,bucket) +); + +CREATE INDEX IF NOT EXISTS idx_analytics_top_talkers_by_iata_time ON analytics_top_talkers_by_iata (bucket); + +CREATE TABLE IF NOT EXISTS analytics_top_advertisers_by_iata ( + iata char(3), + node_id uuid, + bucket timestamptz, + advert_count bigint, + flood_advert_count bigint, + direct_advert_count bigint, + last_heard timestamptz, + name text, + node_type smallint, + PRIMARY KEY (iata,node_id,bucket) +); + +CREATE INDEX IF NOT EXISTS idx_analytics_top_advertisers_by_iata_time ON analytics_top_advertisers_by_iata (bucket); + +CREATE TABLE IF NOT EXISTS analytics_observer_activity_hourly ( + observer_id uuid, + payload_type smallint, + bucket timestamptz, + observations bigint, + airtime_ms real, + airtime_n bigint, + snr_sum real, + snr_n bigint, + snr_min real, + rssi_sum bigint, + rssi_n bigint, + PRIMARY KEY (observer_id,payload_type,bucket) +); + +CREATE INDEX IF NOT EXISTS idx_analytics_observer_activity_hourly_time ON analytics_observer_activity_hourly (bucket); + +CREATE TABLE IF NOT EXISTS analytics_signal_stats_hourly ( + iata char(3), + hour timestamptz, + kind integer, + snr_bin integer, + rssi_bin integer, + receptions bigint, + snr_samples bigint, + snr_sum double precision, + rssi_samples bigint, + rssi_sum double precision, + PRIMARY KEY (iata,hour,kind,snr_bin,rssi_bin) +); + +CREATE INDEX IF NOT EXISTS idx_analytics_signal_stats_hourly_time ON analytics_signal_stats_hourly (hour); + +CREATE TABLE IF NOT EXISTS analytics_path_stats_hourly ( + iata char(3), + hour timestamptz, + category integer, + hash_bytes integer, + entries integer, + receptions bigint, + PRIMARY KEY (iata,hour,category,hash_bytes,entries) +); + +CREATE INDEX IF NOT EXISTS idx_analytics_path_stats_hourly_time ON analytics_path_stats_hourly (hour); + +DROP MATERIALIZED VIEW IF EXISTS mv_hourly_iata_stats; + +CREATE OR REPLACE VIEW analytics_live_hourly_iata_stats AS +SELECT + iata, + date_trunc('hour', heard_at, 'UTC')::timestamptz AS hour, + COUNT(*) AS observation_count, + COUNT(DISTINCT packet_hash) AS unique_packets +FROM packet_observations +WHERE heard_at > NOW() - INTERVAL '30 days' +GROUP BY iata, date_trunc('hour', heard_at, 'UTC'); + +CREATE MATERIALIZED VIEW mv_hourly_iata_stats AS +WITH combined AS ( + SELECT iata, hour, observation_count, unique_packets FROM analytics_live_hourly_iata_stats + UNION ALL + SELECT iata, hour, observation_count, unique_packets FROM analytics_hourly_iata_stats + WHERE hour >= date_trunc('hour', NOW(), 'UTC') - INTERVAL '30 days' +), observers_by_hour AS ( + SELECT iata, date_trunc('hour', heard_at, 'UTC') AS hour, observer_id + FROM packet_observations WHERE heard_at > NOW() - INTERVAL '30 days' + UNION + SELECT iata, bucket AS hour, observer_id FROM analytics_top_observers_by_iata + WHERE bucket >= date_trunc('hour', NOW(), 'UTC') - INTERVAL '30 days' +), observer_counts AS (SELECT iata, hour, COUNT(*) AS n FROM observers_by_hour GROUP BY iata,hour) +SELECT combined.iata, combined.hour, SUM(observation_count)::bigint AS observation_count, SUM(unique_packets)::bigint AS unique_packets, MAX(o.n)::bigint AS active_observers +FROM combined JOIN observer_counts o USING (iata,hour) GROUP BY combined.iata,combined.hour; + +CREATE UNIQUE INDEX idx_analytics_hourly_iata_stats_view ON mv_hourly_iata_stats (iata,hour); + + +DROP MATERIALIZED VIEW IF EXISTS mv_payload_breakdown_by_iata; + +CREATE OR REPLACE VIEW analytics_live_payload_breakdown_by_iata AS +SELECT + iata, + payload_type, + date_trunc('hour', heard_at, 'UTC')::timestamptz AS bucket, + COUNT(*) AS count +FROM packet_observations +WHERE heard_at > NOW() - INTERVAL '30 days' + AND payload_type IS NOT NULL +GROUP BY iata, payload_type, date_trunc('hour', heard_at, 'UTC'); + +CREATE MATERIALIZED VIEW mv_payload_breakdown_by_iata AS +WITH combined AS ( + SELECT iata, payload_type, bucket, count FROM analytics_live_payload_breakdown_by_iata + UNION ALL + SELECT iata, payload_type, bucket, count FROM analytics_payload_breakdown_by_iata + WHERE bucket >= date_trunc('hour', NOW(), 'UTC') - INTERVAL '30 days' +) +SELECT iata, payload_type, bucket, SUM(count)::bigint AS count +FROM combined GROUP BY iata,payload_type,bucket; + +CREATE UNIQUE INDEX idx_analytics_payload_breakdown_by_iata_view ON mv_payload_breakdown_by_iata (iata,payload_type,bucket); + + +DROP MATERIALIZED VIEW IF EXISTS mv_top_observers_by_iata; + +CREATE OR REPLACE VIEW analytics_live_top_observers_by_iata AS +SELECT + po.iata, + po.observer_id, + o.display_name, + o.observer_type, + date_trunc('hour', po.heard_at, 'UTC')::timestamptz AS bucket, + COUNT(*) AS observation_count +FROM packet_observations po +JOIN observers o ON o.id = po.observer_id +WHERE po.heard_at > NOW() - INTERVAL '30 days' +GROUP BY po.iata, po.observer_id, o.display_name, o.observer_type, date_trunc('hour', po.heard_at, 'UTC'); + +CREATE MATERIALIZED VIEW mv_top_observers_by_iata AS +WITH combined AS ( + SELECT iata, observer_id, bucket, observation_count, display_name, observer_type FROM analytics_live_top_observers_by_iata + UNION ALL + SELECT iata, observer_id, bucket, observation_count, display_name, observer_type FROM analytics_top_observers_by_iata + WHERE bucket >= date_trunc('hour', NOW(), 'UTC') - INTERVAL '30 days' +), totals AS (SELECT iata, observer_id, bucket, SUM(observation_count)::bigint AS observation_count, MAX(display_name) AS display_name, MAX(observer_type) AS observer_type +FROM combined GROUP BY iata,observer_id,bucket) +SELECT t.iata, t.observer_id, COALESCE(o.display_name,t.display_name) AS display_name, + COALESCE(o.observer_type,t.observer_type) AS observer_type, t.bucket, t.observation_count +FROM totals t LEFT JOIN observers o ON o.id=t.observer_id; + +CREATE UNIQUE INDEX idx_analytics_top_observers_by_iata_view ON mv_top_observers_by_iata (iata,observer_id,bucket); + + +DROP MATERIALIZED VIEW IF EXISTS mv_top_talkers_by_iata; + +CREATE OR REPLACE VIEW analytics_live_top_talkers_by_iata AS +SELECT + po.iata, + cm.sender_name, + date_trunc('hour', cm.sent_at, 'UTC')::timestamptz AS bucket, + COUNT(DISTINCT cm.id) AS message_count, + MAX(cm.sent_at) AS last_sent +FROM channel_messages cm +JOIN packet_observations po ON po.packet_hash = cm.packet_hash +WHERE cm.sender_name IS NOT NULL + AND cm.sent_at > NOW() - INTERVAL '30 days' +GROUP BY po.iata, cm.sender_name, date_trunc('hour', cm.sent_at, 'UTC'); + +CREATE MATERIALIZED VIEW mv_top_talkers_by_iata AS +WITH combined AS ( + SELECT iata, sender_name, bucket, message_count, last_sent FROM analytics_live_top_talkers_by_iata + UNION ALL + SELECT iata, sender_name, bucket, message_count, last_sent FROM analytics_top_talkers_by_iata + WHERE bucket >= date_trunc('hour', NOW(), 'UTC') - INTERVAL '30 days' +) +SELECT iata, sender_name, bucket, SUM(message_count)::bigint AS message_count, MAX(last_sent) AS last_sent +FROM combined GROUP BY iata,sender_name,bucket; + +CREATE UNIQUE INDEX idx_analytics_top_talkers_by_iata_view ON mv_top_talkers_by_iata (iata,sender_name,bucket); + + +DROP MATERIALIZED VIEW IF EXISTS mv_top_advertisers_by_iata; + +CREATE OR REPLACE VIEW analytics_live_top_advertisers_by_iata AS +SELECT + po.iata, + n.id AS node_id, + n.name, + n.node_type, + date_trunc('hour', po.heard_at, 'UTC')::timestamptz AS bucket, + COUNT(DISTINCT p.packet_hash) AS advert_count, + COUNT(DISTINCT p.packet_hash) FILTER (WHERE p.route_type IN (0, 1)) AS flood_advert_count, + COUNT(DISTINCT p.packet_hash) FILTER (WHERE p.route_type IN (2, 3)) AS direct_advert_count, + MAX(po.heard_at) AS last_heard +FROM packets p +JOIN packet_observations po ON po.packet_hash = p.packet_hash +JOIN nodes n ON n.public_key = p.origin_pubkey +WHERE p.payload_type = 4 -- ADVERT + AND po.heard_at > NOW() - INTERVAL '30 days' +GROUP BY po.iata, n.id, n.name, n.node_type, date_trunc('hour', po.heard_at, 'UTC'); + +CREATE MATERIALIZED VIEW mv_top_advertisers_by_iata AS +WITH combined AS ( + SELECT iata, node_id, bucket, advert_count, flood_advert_count, direct_advert_count, last_heard, name, node_type FROM analytics_live_top_advertisers_by_iata + UNION ALL + SELECT iata, node_id, bucket, advert_count, flood_advert_count, direct_advert_count, last_heard, name, node_type FROM analytics_top_advertisers_by_iata + WHERE bucket >= date_trunc('hour', NOW(), 'UTC') - INTERVAL '30 days' +), totals AS (SELECT iata, node_id, bucket, SUM(advert_count)::bigint AS advert_count, SUM(flood_advert_count)::bigint AS flood_advert_count, SUM(direct_advert_count)::bigint AS direct_advert_count, MAX(last_heard) AS last_heard, MAX(name) AS name, MAX(node_type) AS node_type +FROM combined GROUP BY iata,node_id,bucket) +SELECT t.iata, t.node_id, COALESCE(n.name,t.name) AS name, COALESCE(n.node_type,t.node_type) AS node_type, + t.bucket, t.advert_count, t.flood_advert_count, t.direct_advert_count, t.last_heard +FROM totals t LEFT JOIN nodes n ON n.id=t.node_id; + +CREATE UNIQUE INDEX idx_analytics_top_advertisers_by_iata_view ON mv_top_advertisers_by_iata (iata,node_id,bucket); + + +DROP MATERIALIZED VIEW IF EXISTS mv_observer_activity_hourly; + +CREATE OR REPLACE VIEW analytics_live_observer_activity_hourly AS +SELECT + observer_id, + payload_type, + date_trunc('hour', heard_at, 'UTC')::timestamptz AS bucket, + COUNT(*)::bigint AS observations, + SUM(airtime_ms)::real AS airtime_ms, + COUNT(airtime_ms)::bigint AS airtime_n, + SUM(snr) FILTER (WHERE NOT (COALESCE(rssi, 0) = 0 AND COALESCE(snr, 0) = 0))::real AS snr_sum, + COUNT(snr) FILTER (WHERE NOT (COALESCE(rssi, 0) = 0 AND COALESCE(snr, 0) = 0))::bigint AS snr_n, + MIN(snr) FILTER (WHERE NOT (COALESCE(rssi, 0) = 0 AND COALESCE(snr, 0) = 0))::real AS snr_min, + SUM(rssi) FILTER (WHERE NOT (COALESCE(rssi, 0) = 0 AND COALESCE(snr, 0) = 0))::bigint AS rssi_sum, + COUNT(rssi) FILTER (WHERE NOT (COALESCE(rssi, 0) = 0 AND COALESCE(snr, 0) = 0))::bigint AS rssi_n +FROM packet_observations +WHERE heard_at > NOW() - INTERVAL '30 days' + AND payload_type IS NOT NULL +GROUP BY observer_id, payload_type, date_trunc('hour', heard_at, 'UTC'); + +CREATE MATERIALIZED VIEW mv_observer_activity_hourly AS +WITH combined AS ( + SELECT observer_id, payload_type, bucket, observations, airtime_ms, airtime_n, snr_sum, snr_n, snr_min, rssi_sum, rssi_n FROM analytics_live_observer_activity_hourly + UNION ALL + SELECT observer_id, payload_type, bucket, observations, airtime_ms, airtime_n, snr_sum, snr_n, snr_min, rssi_sum, rssi_n FROM analytics_observer_activity_hourly + WHERE bucket >= date_trunc('hour', NOW(), 'UTC') - INTERVAL '30 days' +) +SELECT observer_id, payload_type, bucket, SUM(observations)::bigint AS observations, SUM(airtime_ms)::real AS airtime_ms, SUM(airtime_n)::bigint AS airtime_n, SUM(snr_sum)::real AS snr_sum, SUM(snr_n)::bigint AS snr_n, MIN(snr_min)::real AS snr_min, SUM(rssi_sum)::bigint AS rssi_sum, SUM(rssi_n)::bigint AS rssi_n +FROM combined GROUP BY observer_id,payload_type,bucket; + +CREATE UNIQUE INDEX idx_analytics_observer_activity_hourly_view ON mv_observer_activity_hourly (observer_id,payload_type,bucket); + + +DROP MATERIALIZED VIEW IF EXISTS mv_signal_stats_hourly; + +CREATE OR REPLACE VIEW analytics_live_signal_stats_hourly AS +WITH samples AS ( + SELECT iata, date_trunc('hour', heard_at, 'UTC') AS hour, + CASE WHEN NOT (COALESCE(rssi, 0) = 0 AND COALESCE(snr, 0) = 0) + AND snr > '-Infinity'::real AND snr < 'Infinity'::real + THEN snr::double precision END AS snr, + CASE WHEN NOT (COALESCE(rssi, 0) = 0 AND COALESCE(snr, 0) = 0) + THEN rssi::double precision END AS rssi + FROM packet_observations + WHERE heard_at >= date_trunc('hour', NOW(), 'UTC') - INTERVAL '720 hours' + AND heard_at < date_trunc('hour', NOW(), 'UTC') +), binned AS ( + SELECT *, width_bucket(snr, -30, 30, 12) AS snr_bin, + width_bucket(rssi, -140, 0, 14) AS rssi_bin + FROM samples +) +SELECT iata, hour, grouping(snr_bin, rssi_bin)::integer AS kind, + COALESCE(snr_bin, -1)::integer AS snr_bin, + COALESCE(rssi_bin, -1)::integer AS rssi_bin, + count(*)::bigint AS receptions, + count(snr)::bigint AS snr_samples, + COALESCE(sum(snr), 0)::double precision AS snr_sum, + count(rssi)::bigint AS rssi_samples, + COALESCE(sum(rssi), 0)::double precision AS rssi_sum +FROM binned +GROUP BY GROUPING SETS ((iata, hour), (iata, hour, snr_bin), (iata, hour, rssi_bin)); + +CREATE MATERIALIZED VIEW mv_signal_stats_hourly AS +WITH combined AS ( + SELECT iata, hour, kind, snr_bin, rssi_bin, receptions, snr_samples, snr_sum, rssi_samples, rssi_sum FROM analytics_live_signal_stats_hourly + UNION ALL + SELECT iata, hour, kind, snr_bin, rssi_bin, receptions, snr_samples, snr_sum, rssi_samples, rssi_sum FROM analytics_signal_stats_hourly + WHERE hour >= date_trunc('hour', NOW(), 'UTC') - INTERVAL '30 days' + AND hour < date_trunc('hour', NOW(), 'UTC') +) +SELECT iata, hour, kind, snr_bin, rssi_bin, SUM(receptions)::bigint AS receptions, SUM(snr_samples)::bigint AS snr_samples, SUM(snr_sum)::double precision AS snr_sum, SUM(rssi_samples)::bigint AS rssi_samples, SUM(rssi_sum)::double precision AS rssi_sum +FROM combined GROUP BY iata,hour,kind,snr_bin,rssi_bin; + +CREATE UNIQUE INDEX idx_analytics_signal_stats_hourly_view ON mv_signal_stats_hourly (iata,hour,kind,snr_bin,rssi_bin); + + +DROP MATERIALIZED VIEW IF EXISTS mv_path_stats_hourly; + +CREATE OR REPLACE VIEW analytics_live_path_stats_hourly AS +WITH classified AS ( + SELECT iata, heard_at, hash_size, hop_count, + CASE WHEN payload_type = 9 THEN 2 + -- Match meshcore-go IsValidPathLen (1/2/3-byte hashes, max 64 path bytes). + WHEN payload_type IS NULL OR payload_type NOT BETWEEN 0 AND 15 + OR NOT (path_length_byte BETWEEN 0 AND 191 + AND hash_size BETWEEN 1 AND 3 AND hop_count BETWEEN 0 AND 63 + AND hash_size = (path_length_byte >> 6) + 1 + AND hop_count = (path_length_byte & 63) + AND hash_size::integer * hop_count::integer <= 64 + AND COALESCE(octet_length(path_bytes), 0) = hash_size::integer * hop_count::integer) + THEN 3 + WHEN hop_count = 0 THEN 1 + ELSE 0 END::integer AS category + FROM packet_observations + WHERE heard_at >= date_trunc('hour', NOW(), 'UTC') - INTERVAL '720 hours' + AND heard_at < date_trunc('hour', NOW(), 'UTC') +), buckets AS ( + SELECT iata, date_trunc('hour', heard_at, 'UTC') AS hour, category, + CASE WHEN category = 0 THEN hash_size ELSE 0 END::integer AS hash_bytes, + CASE WHEN category = 0 THEN hop_count ELSE 0 END::integer AS entries + FROM classified +) +SELECT iata, hour, category, hash_bytes, entries, count(*)::bigint AS receptions +FROM buckets GROUP BY iata, hour, category, hash_bytes, entries; + +CREATE MATERIALIZED VIEW mv_path_stats_hourly AS +WITH combined AS ( + SELECT iata, hour, category, hash_bytes, entries, receptions FROM analytics_live_path_stats_hourly + UNION ALL + SELECT iata, hour, category, hash_bytes, entries, receptions FROM analytics_path_stats_hourly + WHERE hour >= date_trunc('hour', NOW(), 'UTC') - INTERVAL '30 days' + AND hour < date_trunc('hour', NOW(), 'UTC') +) +SELECT iata, hour, category, hash_bytes, entries, SUM(receptions)::bigint AS receptions +FROM combined GROUP BY iata,hour,category,hash_bytes,entries; + +CREATE UNIQUE INDEX idx_analytics_path_stats_hourly_view ON mv_path_stats_hourly (iata,hour,category,hash_bytes,entries); + + +-- VOLATILE gives the archive statement a new READ COMMITTED snapshot after +-- acquiring parent locks. FK inserts cannot race between that snapshot and the +-- cascade. SKIP LOCKED leaves packets being ingested for a later cleanup tick. +CREATE OR REPLACE FUNCTION archive_delete_packets(cutoff timestamptz, batch_size integer) +RETURNS bigint LANGUAGE plpgsql VOLATILE AS $$ +DECLARE hashes bytea[]; deleted bigint; +BEGIN + IF batch_size < 1 THEN RAISE EXCEPTION 'batch_size must be positive'; END IF; + SELECT array_agg(p.packet_hash) INTO hashes FROM ( + SELECT ep.packet_hash FROM packets ep + WHERE ep.last_heard_at < cutoff + ORDER BY ep.last_heard_at, ep.packet_hash + LIMIT batch_size FOR UPDATE OF ep SKIP LOCKED + ) p; + + WITH expired_observations AS MATERIALIZED ( + SELECT po.* FROM packet_observations po WHERE po.packet_hash = ANY(hashes) + ), +archived_hourly_iata_stats AS ( + INSERT INTO analytics_hourly_iata_stats (iata, hour, observation_count, unique_packets) + SELECT iata, hour, observation_count, unique_packets FROM ( +SELECT + iata, + date_trunc('hour', heard_at, 'UTC')::timestamptz AS hour, + COUNT(*) AS observation_count, + COUNT(DISTINCT packet_hash) AS unique_packets +FROM expired_observations +WHERE heard_at >= date_trunc('hour', NOW(), 'UTC') - INTERVAL '30 days' +GROUP BY iata, date_trunc('hour', heard_at, 'UTC') + ) batch + ON CONFLICT (iata,hour) DO UPDATE SET + observation_count = analytics_hourly_iata_stats.observation_count + EXCLUDED.observation_count, + unique_packets = analytics_hourly_iata_stats.unique_packets + EXCLUDED.unique_packets +), +archived_payload_breakdown_by_iata AS ( + INSERT INTO analytics_payload_breakdown_by_iata (iata, payload_type, bucket, count) + SELECT iata, payload_type, bucket, count FROM ( +SELECT + iata, + payload_type, + date_trunc('hour', heard_at, 'UTC')::timestamptz AS bucket, + COUNT(*) AS count +FROM expired_observations +WHERE heard_at >= date_trunc('hour', NOW(), 'UTC') - INTERVAL '30 days' + AND payload_type IS NOT NULL +GROUP BY iata, payload_type, date_trunc('hour', heard_at, 'UTC') + ) batch + ON CONFLICT (iata,payload_type,bucket) DO UPDATE SET + count = analytics_payload_breakdown_by_iata.count + EXCLUDED.count +), +archived_top_observers_by_iata AS ( + INSERT INTO analytics_top_observers_by_iata (iata, observer_id, bucket, observation_count, display_name, observer_type) + SELECT iata, observer_id, bucket, observation_count, display_name, observer_type FROM ( +SELECT + po.iata, + po.observer_id, + o.display_name, + o.observer_type, + date_trunc('hour', po.heard_at, 'UTC')::timestamptz AS bucket, + COUNT(*) AS observation_count +FROM expired_observations po +JOIN observers o ON o.id = po.observer_id +WHERE po.heard_at >= date_trunc('hour', NOW(), 'UTC') - INTERVAL '30 days' +GROUP BY po.iata, po.observer_id, o.display_name, o.observer_type, date_trunc('hour', po.heard_at, 'UTC') + ) batch + ON CONFLICT (iata,observer_id,bucket) DO UPDATE SET + observation_count = analytics_top_observers_by_iata.observation_count + EXCLUDED.observation_count, + display_name = EXCLUDED.display_name, + observer_type = EXCLUDED.observer_type +), +archived_top_talkers_by_iata AS ( + INSERT INTO analytics_top_talkers_by_iata (iata, sender_name, bucket, message_count, last_sent) + SELECT iata, sender_name, bucket, message_count, last_sent FROM ( +SELECT + po.iata, + cm.sender_name, + date_trunc('hour', cm.sent_at, 'UTC')::timestamptz AS bucket, + COUNT(DISTINCT cm.id) AS message_count, + MAX(cm.sent_at) AS last_sent +FROM channel_messages cm +JOIN expired_observations po ON po.packet_hash = cm.packet_hash +WHERE cm.sender_name IS NOT NULL + AND cm.sent_at >= date_trunc('hour', NOW(), 'UTC') - INTERVAL '30 days' +GROUP BY po.iata, cm.sender_name, date_trunc('hour', cm.sent_at, 'UTC') + ) batch + ON CONFLICT (iata,sender_name,bucket) DO UPDATE SET + message_count = analytics_top_talkers_by_iata.message_count + EXCLUDED.message_count, + last_sent = GREATEST(analytics_top_talkers_by_iata.last_sent, EXCLUDED.last_sent) +), +archived_top_advertisers_by_iata AS ( + INSERT INTO analytics_top_advertisers_by_iata (iata, node_id, bucket, advert_count, flood_advert_count, direct_advert_count, last_heard, name, node_type) + SELECT iata, node_id, bucket, advert_count, flood_advert_count, direct_advert_count, last_heard, name, node_type FROM ( +SELECT + po.iata, + n.id AS node_id, + n.name, + n.node_type, + date_trunc('hour', po.heard_at, 'UTC')::timestamptz AS bucket, + COUNT(DISTINCT p.packet_hash) AS advert_count, + COUNT(DISTINCT p.packet_hash) FILTER (WHERE p.route_type IN (0, 1)) AS flood_advert_count, + COUNT(DISTINCT p.packet_hash) FILTER (WHERE p.route_type IN (2, 3)) AS direct_advert_count, + MAX(po.heard_at) AS last_heard +FROM packets p +JOIN expired_observations po ON po.packet_hash = p.packet_hash +JOIN nodes n ON n.public_key = p.origin_pubkey +WHERE p.payload_type = 4 -- ADVERT + AND po.heard_at >= date_trunc('hour', NOW(), 'UTC') - INTERVAL '30 days' +GROUP BY po.iata, n.id, n.name, n.node_type, date_trunc('hour', po.heard_at, 'UTC') + ) batch + ON CONFLICT (iata,node_id,bucket) DO UPDATE SET + advert_count = analytics_top_advertisers_by_iata.advert_count + EXCLUDED.advert_count, + flood_advert_count = analytics_top_advertisers_by_iata.flood_advert_count + EXCLUDED.flood_advert_count, + direct_advert_count = analytics_top_advertisers_by_iata.direct_advert_count + EXCLUDED.direct_advert_count, + last_heard = GREATEST(analytics_top_advertisers_by_iata.last_heard, EXCLUDED.last_heard), + name = EXCLUDED.name, + node_type = EXCLUDED.node_type +), +archived_observer_activity_hourly AS ( + INSERT INTO analytics_observer_activity_hourly (observer_id, payload_type, bucket, observations, airtime_ms, airtime_n, snr_sum, snr_n, snr_min, rssi_sum, rssi_n) + SELECT observer_id, payload_type, bucket, observations, airtime_ms, airtime_n, snr_sum, snr_n, snr_min, rssi_sum, rssi_n FROM ( +SELECT + observer_id, + payload_type, + date_trunc('hour', heard_at, 'UTC')::timestamptz AS bucket, + COUNT(*)::bigint AS observations, + SUM(airtime_ms)::real AS airtime_ms, + COUNT(airtime_ms)::bigint AS airtime_n, + SUM(snr) FILTER (WHERE NOT (COALESCE(rssi, 0) = 0 AND COALESCE(snr, 0) = 0))::real AS snr_sum, + COUNT(snr) FILTER (WHERE NOT (COALESCE(rssi, 0) = 0 AND COALESCE(snr, 0) = 0))::bigint AS snr_n, + MIN(snr) FILTER (WHERE NOT (COALESCE(rssi, 0) = 0 AND COALESCE(snr, 0) = 0))::real AS snr_min, + SUM(rssi) FILTER (WHERE NOT (COALESCE(rssi, 0) = 0 AND COALESCE(snr, 0) = 0))::bigint AS rssi_sum, + COUNT(rssi) FILTER (WHERE NOT (COALESCE(rssi, 0) = 0 AND COALESCE(snr, 0) = 0))::bigint AS rssi_n +FROM expired_observations +WHERE heard_at >= date_trunc('hour', NOW(), 'UTC') - INTERVAL '30 days' + AND payload_type IS NOT NULL +GROUP BY observer_id, payload_type, date_trunc('hour', heard_at, 'UTC') + ) batch + ON CONFLICT (observer_id,payload_type,bucket) DO UPDATE SET + observations = analytics_observer_activity_hourly.observations + EXCLUDED.observations, + airtime_ms = CASE WHEN analytics_observer_activity_hourly.airtime_ms IS NULL AND EXCLUDED.airtime_ms IS NULL THEN NULL ELSE COALESCE(analytics_observer_activity_hourly.airtime_ms, 0) + COALESCE(EXCLUDED.airtime_ms, 0) END, + airtime_n = analytics_observer_activity_hourly.airtime_n + EXCLUDED.airtime_n, + snr_sum = CASE WHEN analytics_observer_activity_hourly.snr_sum IS NULL AND EXCLUDED.snr_sum IS NULL THEN NULL ELSE COALESCE(analytics_observer_activity_hourly.snr_sum, 0) + COALESCE(EXCLUDED.snr_sum, 0) END, + snr_n = analytics_observer_activity_hourly.snr_n + EXCLUDED.snr_n, + snr_min = LEAST(analytics_observer_activity_hourly.snr_min, EXCLUDED.snr_min), + rssi_sum = CASE WHEN analytics_observer_activity_hourly.rssi_sum IS NULL AND EXCLUDED.rssi_sum IS NULL THEN NULL ELSE COALESCE(analytics_observer_activity_hourly.rssi_sum, 0) + COALESCE(EXCLUDED.rssi_sum, 0) END, + rssi_n = analytics_observer_activity_hourly.rssi_n + EXCLUDED.rssi_n +), +archived_signal_stats_hourly AS ( + INSERT INTO analytics_signal_stats_hourly (iata, hour, kind, snr_bin, rssi_bin, receptions, snr_samples, snr_sum, rssi_samples, rssi_sum) + SELECT iata, hour, kind, snr_bin, rssi_bin, receptions, snr_samples, snr_sum, rssi_samples, rssi_sum FROM ( +WITH samples AS ( + SELECT iata, date_trunc('hour', heard_at, 'UTC') AS hour, + CASE WHEN NOT (COALESCE(rssi, 0) = 0 AND COALESCE(snr, 0) = 0) + AND snr > '-Infinity'::real AND snr < 'Infinity'::real + THEN snr::double precision END AS snr, + CASE WHEN NOT (COALESCE(rssi, 0) = 0 AND COALESCE(snr, 0) = 0) + THEN rssi::double precision END AS rssi + FROM expired_observations + WHERE heard_at >= date_trunc('hour', NOW(), 'UTC') - INTERVAL '720 hours' + +), binned AS ( + SELECT *, width_bucket(snr, -30, 30, 12) AS snr_bin, + width_bucket(rssi, -140, 0, 14) AS rssi_bin + FROM samples +) +SELECT iata, hour, grouping(snr_bin, rssi_bin)::integer AS kind, + COALESCE(snr_bin, -1)::integer AS snr_bin, + COALESCE(rssi_bin, -1)::integer AS rssi_bin, + count(*)::bigint AS receptions, + count(snr)::bigint AS snr_samples, + COALESCE(sum(snr), 0)::double precision AS snr_sum, + count(rssi)::bigint AS rssi_samples, + COALESCE(sum(rssi), 0)::double precision AS rssi_sum +FROM binned +GROUP BY GROUPING SETS ((iata, hour), (iata, hour, snr_bin), (iata, hour, rssi_bin)) + ) batch + ON CONFLICT (iata,hour,kind,snr_bin,rssi_bin) DO UPDATE SET + receptions = analytics_signal_stats_hourly.receptions + EXCLUDED.receptions, + snr_samples = analytics_signal_stats_hourly.snr_samples + EXCLUDED.snr_samples, + snr_sum = analytics_signal_stats_hourly.snr_sum + EXCLUDED.snr_sum, + rssi_samples = analytics_signal_stats_hourly.rssi_samples + EXCLUDED.rssi_samples, + rssi_sum = analytics_signal_stats_hourly.rssi_sum + EXCLUDED.rssi_sum +), +archived_path_stats_hourly AS ( + INSERT INTO analytics_path_stats_hourly (iata, hour, category, hash_bytes, entries, receptions) + SELECT iata, hour, category, hash_bytes, entries, receptions FROM ( +WITH classified AS ( + SELECT iata, heard_at, hash_size, hop_count, + CASE WHEN payload_type = 9 THEN 2 + -- Match meshcore-go IsValidPathLen (1/2/3-byte hashes, max 64 path bytes). + WHEN payload_type IS NULL OR payload_type NOT BETWEEN 0 AND 15 + OR NOT (path_length_byte BETWEEN 0 AND 191 + AND hash_size BETWEEN 1 AND 3 AND hop_count BETWEEN 0 AND 63 + AND hash_size = (path_length_byte >> 6) + 1 + AND hop_count = (path_length_byte & 63) + AND hash_size::integer * hop_count::integer <= 64 + AND COALESCE(octet_length(path_bytes), 0) = hash_size::integer * hop_count::integer) + THEN 3 + WHEN hop_count = 0 THEN 1 + ELSE 0 END::integer AS category + FROM expired_observations + WHERE heard_at >= date_trunc('hour', NOW(), 'UTC') - INTERVAL '720 hours' + +), buckets AS ( + SELECT iata, date_trunc('hour', heard_at, 'UTC') AS hour, category, + CASE WHEN category = 0 THEN hash_size ELSE 0 END::integer AS hash_bytes, + CASE WHEN category = 0 THEN hop_count ELSE 0 END::integer AS entries + FROM classified +) +SELECT iata, hour, category, hash_bytes, entries, count(*)::bigint AS receptions +FROM buckets GROUP BY iata, hour, category, hash_bytes, entries + ) batch + ON CONFLICT (iata,hour,category,hash_bytes,entries) DO UPDATE SET + receptions = analytics_path_stats_hourly.receptions + EXCLUDED.receptions +) + DELETE FROM packets WHERE packet_hash = ANY(hashes); + GET DIAGNOSTICS deleted = ROW_COUNT; + + -- Once per completed cleanup, including when there were no raw packets to delete. + IF deleted < batch_size THEN + DELETE FROM analytics_hourly_iata_stats WHERE hour < date_trunc('hour', NOW(), 'UTC') - INTERVAL '30 days'; + DELETE FROM analytics_payload_breakdown_by_iata WHERE bucket < date_trunc('hour', NOW(), 'UTC') - INTERVAL '30 days'; + DELETE FROM analytics_top_observers_by_iata WHERE bucket < date_trunc('hour', NOW(), 'UTC') - INTERVAL '30 days'; + DELETE FROM analytics_top_talkers_by_iata WHERE bucket < date_trunc('hour', NOW(), 'UTC') - INTERVAL '30 days'; + DELETE FROM analytics_top_advertisers_by_iata WHERE bucket < date_trunc('hour', NOW(), 'UTC') - INTERVAL '30 days'; + DELETE FROM analytics_observer_activity_hourly WHERE bucket < date_trunc('hour', NOW(), 'UTC') - INTERVAL '30 days'; + DELETE FROM analytics_signal_stats_hourly WHERE hour < date_trunc('hour', NOW(), 'UTC') - INTERVAL '30 days'; + DELETE FROM analytics_path_stats_hourly WHERE hour < date_trunc('hour', NOW(), 'UTC') - INTERVAL '30 days'; + END IF; + RETURN deleted; +END $$; diff --git a/db/queries/queries.sql b/db/queries/queries.sql index ccc8492f..29690a16 100644 --- a/db/queries/queries.sql +++ b/db/queries/queries.sql @@ -661,18 +661,9 @@ ORDER BY po.id ASC LIMIT $6; --- 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: DeleteOldPackets :one +-- Locks, archives and cascades one bounded packet cohort in a single transaction. +SELECT archive_delete_packets(@cutoff::timestamptz, @batch_size::integer)::bigint; -- name: DeleteOldNodes :exec -- Deletes nodes not seen since the given cutoff. node_iatas and node_neighbors cascade- diff --git a/db/retention_integration_test.go b/db/retention_integration_test.go index 5e40dbd2..d6307176 100644 --- a/db/retention_integration_test.go +++ b/db/retention_integration_test.go @@ -46,10 +46,9 @@ func countRows(t *testing.T, ctx context.Context, tx pgx.Tx, sql string) int { // More than one batch of old packets; observations go with them. func TestDeleteOldPacketsBatchesPostgres(t *testing.T) { ctx, tx := retentionTx(t) + analyticsTables(t, ctx, tx) + applyStatsMigration(t, ctx, tx, "039_analytics_retention.sql") _, 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, diff --git a/db/sqlc/models.go b/db/sqlc/models.go index b5ffb054..2c6e20f6 100644 --- a/db/sqlc/models.go +++ b/db/sqlc/models.go @@ -16,6 +16,164 @@ type Account struct { DeactivatedAt pgtype.Timestamptz `json:"deactivated_at"` } +type AnalyticsHourlyIataStat struct { + Iata string `json:"iata"` + Hour pgtype.Timestamptz `json:"hour"` + ObservationCount *int64 `json:"observation_count"` + UniquePackets *int64 `json:"unique_packets"` +} + +type AnalyticsLiveHourlyIataStat struct { + Iata string `json:"iata"` + Hour pgtype.Timestamptz `json:"hour"` + ObservationCount int64 `json:"observation_count"` + UniquePackets int64 `json:"unique_packets"` +} + +type AnalyticsLiveObserverActivityHourly struct { + ObserverID uuid.UUID `json:"observer_id"` + PayloadType *int16 `json:"payload_type"` + Bucket pgtype.Timestamptz `json:"bucket"` + Observations int64 `json:"observations"` + AirtimeMs float32 `json:"airtime_ms"` + AirtimeN int64 `json:"airtime_n"` + SnrSum float32 `json:"snr_sum"` + SnrN int64 `json:"snr_n"` + SnrMin float32 `json:"snr_min"` + RssiSum int64 `json:"rssi_sum"` + RssiN int64 `json:"rssi_n"` +} + +type AnalyticsLivePathStatsHourly struct { + Iata string `json:"iata"` + Hour interface{} `json:"hour"` + Category int32 `json:"category"` + HashBytes int32 `json:"hash_bytes"` + Entries int32 `json:"entries"` + Receptions int64 `json:"receptions"` +} + +type AnalyticsLivePayloadBreakdownByIatum struct { + Iata string `json:"iata"` + PayloadType *int16 `json:"payload_type"` + Bucket pgtype.Timestamptz `json:"bucket"` + Count int64 `json:"count"` +} + +type AnalyticsLiveSignalStatsHourly struct { + Iata string `json:"iata"` + Hour interface{} `json:"hour"` + Kind int32 `json:"kind"` + SnrBin int32 `json:"snr_bin"` + RssiBin int32 `json:"rssi_bin"` + Receptions int64 `json:"receptions"` + SnrSamples int64 `json:"snr_samples"` + SnrSum float64 `json:"snr_sum"` + RssiSamples int64 `json:"rssi_samples"` + RssiSum float64 `json:"rssi_sum"` +} + +type AnalyticsLiveTopAdvertisersByIatum struct { + Iata string `json:"iata"` + NodeID uuid.UUID `json:"node_id"` + Name *string `json:"name"` + NodeType int16 `json:"node_type"` + Bucket pgtype.Timestamptz `json:"bucket"` + AdvertCount int64 `json:"advert_count"` + FloodAdvertCount int64 `json:"flood_advert_count"` + DirectAdvertCount int64 `json:"direct_advert_count"` + LastHeard interface{} `json:"last_heard"` +} + +type AnalyticsLiveTopObserversByIatum struct { + Iata string `json:"iata"` + ObserverID uuid.UUID `json:"observer_id"` + DisplayName *string `json:"display_name"` + ObserverType *string `json:"observer_type"` + Bucket pgtype.Timestamptz `json:"bucket"` + ObservationCount int64 `json:"observation_count"` +} + +type AnalyticsLiveTopTalkersByIatum struct { + Iata string `json:"iata"` + SenderName *string `json:"sender_name"` + Bucket pgtype.Timestamptz `json:"bucket"` + MessageCount int64 `json:"message_count"` + LastSent interface{} `json:"last_sent"` +} + +type AnalyticsObserverActivityHourly struct { + ObserverID uuid.UUID `json:"observer_id"` + PayloadType int16 `json:"payload_type"` + Bucket pgtype.Timestamptz `json:"bucket"` + Observations *int64 `json:"observations"` + AirtimeMs *float32 `json:"airtime_ms"` + AirtimeN *int64 `json:"airtime_n"` + SnrSum *float32 `json:"snr_sum"` + SnrN *int64 `json:"snr_n"` + SnrMin *float32 `json:"snr_min"` + RssiSum *int64 `json:"rssi_sum"` + RssiN *int64 `json:"rssi_n"` +} + +type AnalyticsPathStatsHourly struct { + Iata string `json:"iata"` + Hour pgtype.Timestamptz `json:"hour"` + Category int32 `json:"category"` + HashBytes int32 `json:"hash_bytes"` + Entries int32 `json:"entries"` + Receptions *int64 `json:"receptions"` +} + +type AnalyticsPayloadBreakdownByIatum struct { + Iata string `json:"iata"` + PayloadType int16 `json:"payload_type"` + Bucket pgtype.Timestamptz `json:"bucket"` + Count *int64 `json:"count"` +} + +type AnalyticsSignalStatsHourly struct { + Iata string `json:"iata"` + Hour pgtype.Timestamptz `json:"hour"` + Kind int32 `json:"kind"` + SnrBin int32 `json:"snr_bin"` + RssiBin int32 `json:"rssi_bin"` + Receptions *int64 `json:"receptions"` + SnrSamples *int64 `json:"snr_samples"` + SnrSum *float64 `json:"snr_sum"` + RssiSamples *int64 `json:"rssi_samples"` + RssiSum *float64 `json:"rssi_sum"` +} + +type AnalyticsTopAdvertisersByIatum struct { + Iata string `json:"iata"` + NodeID uuid.UUID `json:"node_id"` + Bucket pgtype.Timestamptz `json:"bucket"` + AdvertCount *int64 `json:"advert_count"` + FloodAdvertCount *int64 `json:"flood_advert_count"` + DirectAdvertCount *int64 `json:"direct_advert_count"` + LastHeard pgtype.Timestamptz `json:"last_heard"` + Name *string `json:"name"` + NodeType *int16 `json:"node_type"` +} + +type AnalyticsTopObserversByIatum struct { + Iata string `json:"iata"` + ObserverID uuid.UUID `json:"observer_id"` + Bucket pgtype.Timestamptz `json:"bucket"` + ObservationCount *int64 `json:"observation_count"` + DisplayName *string `json:"display_name"` + ObserverType *string `json:"observer_type"` +} + +type AnalyticsTopTalkersByIatum struct { + Iata string `json:"iata"` + SenderName string `json:"sender_name"` + Bucket pgtype.Timestamptz `json:"bucket"` + MessageCount *int64 `json:"message_count"` + LastSent pgtype.Timestamptz `json:"last_sent"` +} + type Channel struct { ID int32 `json:"id"` ChannelHash []byte `json:"channel_hash"` diff --git a/db/sqlc/querier.go b/db/sqlc/querier.go index 1547c1c9..38e6c971 100644 --- a/db/sqlc/querier.go +++ b/db/sqlc/querier.go @@ -29,7 +29,7 @@ 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) - // One batch of expired packets; observations and channel messages cascade. + // Locks, archives and cascades one bounded packet cohort in a single transaction. 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. diff --git a/db/sqlc/queries.sql.go b/db/sqlc/queries.sql.go index c047b44f..6a0eee13 100644 --- a/db/sqlc/queries.sql.go +++ b/db/sqlc/queries.sql.go @@ -124,17 +124,8 @@ func (q *Queries) DeleteOldObservers(ctx context.Context, lastSeen pgtype.Timest return items, nil } -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 +const deleteOldPackets = `-- name: DeleteOldPackets :one +SELECT archive_delete_packets($1::timestamptz, $2::integer)::bigint ` type DeleteOldPacketsParams struct { @@ -142,13 +133,12 @@ type DeleteOldPacketsParams struct { BatchSize int32 `json:"batch_size"` } -// One batch of expired packets; observations and channel messages cascade. +// Locks, archives and cascades one bounded packet cohort in a single transaction. 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 + row := q.db.QueryRow(ctx, deleteOldPackets, arg.Cutoff, arg.BatchSize) + var column_1 int64 + err := row.Scan(&column_1) + return column_1, err } const deleteOldRoutes = `-- name: DeleteOldRoutes :execrows diff --git a/sqlc.yaml b/sqlc.yaml index 79431093..032c9911 100644 --- a/sqlc.yaml +++ b/sqlc.yaml @@ -22,3 +22,11 @@ sql: go_type: "encoding/json.RawMessage" - db_type: "bpchar" go_type: "string" + + # Coalesced archive labels can still be NULL; preserve the API pointers. + - column: "mv_top_observers_by_iata.display_name" + go_type: {type: "string", pointer: true} + - column: "mv_top_observers_by_iata.observer_type" + go_type: {type: "string", pointer: true} + - column: "mv_top_advertisers_by_iata.name" + go_type: {type: "string", pointer: true} From 2ab36356fa023abb247db45c8a33de4410507f64 Mon Sep 17 00:00:00 2001 From: n30nex Date: Tue, 29 Sep 2026 18:27:39 -0400 Subject: [PATCH 2/3] fix: address maintainer review for #167 --- db/analytics_retention_integration_test.go | 34 +++++++++++++++ db/migrations/039_analytics_retention.sql | 48 +++++++++++----------- db/sqlc/models.go | 34 +++++++-------- sqlc.yaml | 3 ++ 4 files changed, 79 insertions(+), 40 deletions(-) diff --git a/db/analytics_retention_integration_test.go b/db/analytics_retention_integration_test.go index 2e0466c4..10e33339 100644 --- a/db/analytics_retention_integration_test.go +++ b/db/analytics_retention_integration_test.go @@ -9,6 +9,7 @@ import ( "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgtype" "os" + "strings" "testing" "time" ) @@ -87,6 +88,30 @@ func TestAnalyticsRetentionConcurrentPostgres(t *testing.T) { if err := conn.QueryRow(ctx, "SELECT observation_count FROM analytics_hourly_iata_stats").Scan(&count); err != nil || count != 1 { t.Fatalf("committed reception not archived: %d %v", count, err) } + // A second cleanup runner must wait before touching shared aggregate keys. + first, err := writer.Begin(ctx) + if err != nil { + t.Fatal(err) + } + defer first.Rollback(context.Background()) + if _, err := sqlc.New(first).DeleteOldPackets(ctx, params); err != nil { + t.Fatal(err) + } + if _, err := conn.Exec(ctx, "SET statement_timeout='300ms'"); err != nil { + t.Fatal(err) + } + if _, err := q.DeleteOldPackets(ctx, params); err == nil { + t.Fatalf("overlapping cleanup did not wait for the first transaction: %v", err) + } + if _, err := conn.Exec(ctx, "SET statement_timeout=0"); err != nil { + t.Fatal(err) + } + if err := first.Commit(ctx); err != nil { + t.Fatal(err) + } + if n, err := q.DeleteOldPackets(ctx, params); err != nil || n != 0 { + t.Fatalf("cleanup lock was not released: %d %v", n, err) + } } // Minimal real tables keep this regression runnable against an empty CI database. @@ -136,6 +161,12 @@ func TestAnalyticsRetentionPostgres(t *testing.T) { 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) + } + } before := analyticsSnapshot(t, ctx, tx) q := sqlc.New(tx) cutoff := time.Now().Add(-72 * time.Hour) @@ -184,6 +215,9 @@ func TestAnalyticsRetentionPostgres(t *testing.T) { t.Fatal(err) } late := analyticsSnapshot(t, ctx, tx) + if countRows(t, ctx, tx, "SELECT COALESCE(SUM(observations),0) FROM mv_observer_activity_hourly WHERE payload_type=-1") != 1 { + t.Fatal("unknown-payload observation was omitted") + } if err := (&Store{q: q}).DeleteOldPackets(ctx, cutoff); err != nil { t.Fatal(err) } diff --git a/db/migrations/039_analytics_retention.sql b/db/migrations/039_analytics_retention.sql index 75d50078..069204ad 100644 --- a/db/migrations/039_analytics_retention.sql +++ b/db/migrations/039_analytics_retention.sql @@ -9,8 +9,8 @@ CREATE TABLE IF NOT EXISTS analytics_hourly_iata_stats ( iata char(3), hour timestamptz, - observation_count bigint, - unique_packets bigint, + observation_count bigint NOT NULL DEFAULT 0, + unique_packets bigint NOT NULL DEFAULT 0, PRIMARY KEY (iata,hour) ); @@ -20,7 +20,7 @@ CREATE TABLE IF NOT EXISTS analytics_payload_breakdown_by_iata ( iata char(3), payload_type smallint, bucket timestamptz, - count bigint, + count bigint NOT NULL DEFAULT 0, PRIMARY KEY (iata,payload_type,bucket) ); @@ -30,7 +30,7 @@ CREATE TABLE IF NOT EXISTS analytics_top_observers_by_iata ( iata char(3), observer_id uuid, bucket timestamptz, - observation_count bigint, + observation_count bigint NOT NULL DEFAULT 0, display_name text, observer_type text, PRIMARY KEY (iata,observer_id,bucket) @@ -42,7 +42,7 @@ CREATE TABLE IF NOT EXISTS analytics_top_talkers_by_iata ( iata char(3), sender_name text, bucket timestamptz, - message_count bigint, + message_count bigint NOT NULL DEFAULT 0, last_sent timestamptz, PRIMARY KEY (iata,sender_name,bucket) ); @@ -53,9 +53,9 @@ CREATE TABLE IF NOT EXISTS analytics_top_advertisers_by_iata ( iata char(3), node_id uuid, bucket timestamptz, - advert_count bigint, - flood_advert_count bigint, - direct_advert_count bigint, + advert_count bigint NOT NULL DEFAULT 0, + flood_advert_count bigint NOT NULL DEFAULT 0, + direct_advert_count bigint NOT NULL DEFAULT 0, last_heard timestamptz, name text, node_type smallint, @@ -68,14 +68,14 @@ CREATE TABLE IF NOT EXISTS analytics_observer_activity_hourly ( observer_id uuid, payload_type smallint, bucket timestamptz, - observations bigint, + observations bigint NOT NULL DEFAULT 0, airtime_ms real, - airtime_n bigint, + airtime_n bigint NOT NULL DEFAULT 0, snr_sum real, - snr_n bigint, + snr_n bigint NOT NULL DEFAULT 0, snr_min real, rssi_sum bigint, - rssi_n bigint, + rssi_n bigint NOT NULL DEFAULT 0, PRIMARY KEY (observer_id,payload_type,bucket) ); @@ -87,10 +87,10 @@ CREATE TABLE IF NOT EXISTS analytics_signal_stats_hourly ( kind integer, snr_bin integer, rssi_bin integer, - receptions bigint, - snr_samples bigint, + receptions bigint NOT NULL DEFAULT 0, + snr_samples bigint NOT NULL DEFAULT 0, snr_sum double precision, - rssi_samples bigint, + rssi_samples bigint NOT NULL DEFAULT 0, rssi_sum double precision, PRIMARY KEY (iata,hour,kind,snr_bin,rssi_bin) ); @@ -103,7 +103,7 @@ CREATE TABLE IF NOT EXISTS analytics_path_stats_hourly ( category integer, hash_bytes integer, entries integer, - receptions bigint, + receptions bigint NOT NULL DEFAULT 0, PRIMARY KEY (iata,hour,category,hash_bytes,entries) ); @@ -260,11 +260,10 @@ CREATE UNIQUE INDEX idx_analytics_top_advertisers_by_iata_view ON mv_top_adverti DROP MATERIALIZED VIEW IF EXISTS mv_observer_activity_hourly; - CREATE OR REPLACE VIEW analytics_live_observer_activity_hourly AS SELECT observer_id, - payload_type, + COALESCE(payload_type, -1)::smallint AS payload_type, date_trunc('hour', heard_at, 'UTC')::timestamptz AS bucket, COUNT(*)::bigint AS observations, SUM(airtime_ms)::real AS airtime_ms, @@ -276,8 +275,7 @@ SELECT COUNT(rssi) FILTER (WHERE NOT (COALESCE(rssi, 0) = 0 AND COALESCE(snr, 0) = 0))::bigint AS rssi_n FROM packet_observations WHERE heard_at > NOW() - INTERVAL '30 days' - AND payload_type IS NOT NULL -GROUP BY observer_id, payload_type, date_trunc('hour', heard_at, 'UTC'); +GROUP BY observer_id, COALESCE(payload_type, -1), date_trunc('hour', heard_at, 'UTC'); CREATE MATERIALIZED VIEW mv_observer_activity_hourly AS WITH combined AS ( @@ -292,6 +290,7 @@ FROM combined GROUP BY observer_id,payload_type,bucket; CREATE UNIQUE INDEX idx_analytics_observer_activity_hourly_view ON mv_observer_activity_hourly (observer_id,payload_type,bucket); + DROP MATERIALIZED VIEW IF EXISTS mv_signal_stats_hourly; CREATE OR REPLACE VIEW analytics_live_signal_stats_hourly AS @@ -333,6 +332,7 @@ SELECT iata, hour, kind, snr_bin, rssi_bin, SUM(receptions)::bigint AS reception FROM combined GROUP BY iata,hour,kind,snr_bin,rssi_bin; CREATE UNIQUE INDEX idx_analytics_signal_stats_hourly_view ON mv_signal_stats_hourly (iata,hour,kind,snr_bin,rssi_bin); +CREATE INDEX idx_mv_signal_stats_hourly_hour ON mv_signal_stats_hourly (hour); DROP MATERIALIZED VIEW IF EXISTS mv_path_stats_hourly; @@ -376,6 +376,7 @@ SELECT iata, hour, category, hash_bytes, entries, SUM(receptions)::bigint AS rec FROM combined GROUP BY iata,hour,category,hash_bytes,entries; CREATE UNIQUE INDEX idx_analytics_path_stats_hourly_view ON mv_path_stats_hourly (iata,hour,category,hash_bytes,entries); +CREATE INDEX idx_mv_path_stats_hourly_hour ON mv_path_stats_hourly (hour); -- VOLATILE gives the archive statement a new READ COMMITTED snapshot after @@ -385,6 +386,8 @@ CREATE OR REPLACE FUNCTION archive_delete_packets(cutoff timestamptz, batch_size RETURNS bigint LANGUAGE plpgsql VOLATILE AS $$ DECLARE hashes bytea[]; deleted bigint; BEGIN + -- Serialize cleanup runners before overlapping aggregate upserts. + PERFORM pg_advisory_xact_lock(hashtext('beacon.archive_delete_packets')); IF batch_size < 1 THEN RAISE EXCEPTION 'batch_size must be positive'; END IF; SELECT array_agg(p.packet_hash) INTO hashes FROM ( SELECT ep.packet_hash FROM packets ep @@ -500,7 +503,7 @@ archived_observer_activity_hourly AS ( SELECT observer_id, payload_type, bucket, observations, airtime_ms, airtime_n, snr_sum, snr_n, snr_min, rssi_sum, rssi_n FROM ( SELECT observer_id, - payload_type, + COALESCE(payload_type, -1)::smallint AS payload_type, date_trunc('hour', heard_at, 'UTC')::timestamptz AS bucket, COUNT(*)::bigint AS observations, SUM(airtime_ms)::real AS airtime_ms, @@ -512,8 +515,7 @@ SELECT COUNT(rssi) FILTER (WHERE NOT (COALESCE(rssi, 0) = 0 AND COALESCE(snr, 0) = 0))::bigint AS rssi_n FROM expired_observations WHERE heard_at >= date_trunc('hour', NOW(), 'UTC') - INTERVAL '30 days' - AND payload_type IS NOT NULL -GROUP BY observer_id, payload_type, date_trunc('hour', heard_at, 'UTC') +GROUP BY observer_id, COALESCE(payload_type, -1), date_trunc('hour', heard_at, 'UTC') ) batch ON CONFLICT (observer_id,payload_type,bucket) DO UPDATE SET observations = analytics_observer_activity_hourly.observations + EXCLUDED.observations, diff --git a/db/sqlc/models.go b/db/sqlc/models.go index 2c6e20f6..8a24d538 100644 --- a/db/sqlc/models.go +++ b/db/sqlc/models.go @@ -19,8 +19,8 @@ type Account struct { type AnalyticsHourlyIataStat struct { Iata string `json:"iata"` Hour pgtype.Timestamptz `json:"hour"` - ObservationCount *int64 `json:"observation_count"` - UniquePackets *int64 `json:"unique_packets"` + ObservationCount int64 `json:"observation_count"` + UniquePackets int64 `json:"unique_packets"` } type AnalyticsLiveHourlyIataStat struct { @@ -32,7 +32,7 @@ type AnalyticsLiveHourlyIataStat struct { type AnalyticsLiveObserverActivityHourly struct { ObserverID uuid.UUID `json:"observer_id"` - PayloadType *int16 `json:"payload_type"` + PayloadType int16 `json:"payload_type"` Bucket pgtype.Timestamptz `json:"bucket"` Observations int64 `json:"observations"` AirtimeMs float32 `json:"airtime_ms"` @@ -106,14 +106,14 @@ type AnalyticsObserverActivityHourly struct { ObserverID uuid.UUID `json:"observer_id"` PayloadType int16 `json:"payload_type"` Bucket pgtype.Timestamptz `json:"bucket"` - Observations *int64 `json:"observations"` + Observations int64 `json:"observations"` AirtimeMs *float32 `json:"airtime_ms"` - AirtimeN *int64 `json:"airtime_n"` + AirtimeN int64 `json:"airtime_n"` SnrSum *float32 `json:"snr_sum"` - SnrN *int64 `json:"snr_n"` + SnrN int64 `json:"snr_n"` SnrMin *float32 `json:"snr_min"` RssiSum *int64 `json:"rssi_sum"` - RssiN *int64 `json:"rssi_n"` + RssiN int64 `json:"rssi_n"` } type AnalyticsPathStatsHourly struct { @@ -122,14 +122,14 @@ type AnalyticsPathStatsHourly struct { Category int32 `json:"category"` HashBytes int32 `json:"hash_bytes"` Entries int32 `json:"entries"` - Receptions *int64 `json:"receptions"` + Receptions int64 `json:"receptions"` } type AnalyticsPayloadBreakdownByIatum struct { Iata string `json:"iata"` PayloadType int16 `json:"payload_type"` Bucket pgtype.Timestamptz `json:"bucket"` - Count *int64 `json:"count"` + Count int64 `json:"count"` } type AnalyticsSignalStatsHourly struct { @@ -138,10 +138,10 @@ type AnalyticsSignalStatsHourly struct { Kind int32 `json:"kind"` SnrBin int32 `json:"snr_bin"` RssiBin int32 `json:"rssi_bin"` - Receptions *int64 `json:"receptions"` - SnrSamples *int64 `json:"snr_samples"` + Receptions int64 `json:"receptions"` + SnrSamples int64 `json:"snr_samples"` SnrSum *float64 `json:"snr_sum"` - RssiSamples *int64 `json:"rssi_samples"` + RssiSamples int64 `json:"rssi_samples"` RssiSum *float64 `json:"rssi_sum"` } @@ -149,9 +149,9 @@ type AnalyticsTopAdvertisersByIatum struct { Iata string `json:"iata"` NodeID uuid.UUID `json:"node_id"` Bucket pgtype.Timestamptz `json:"bucket"` - AdvertCount *int64 `json:"advert_count"` - FloodAdvertCount *int64 `json:"flood_advert_count"` - DirectAdvertCount *int64 `json:"direct_advert_count"` + AdvertCount int64 `json:"advert_count"` + FloodAdvertCount int64 `json:"flood_advert_count"` + DirectAdvertCount int64 `json:"direct_advert_count"` LastHeard pgtype.Timestamptz `json:"last_heard"` Name *string `json:"name"` NodeType *int16 `json:"node_type"` @@ -161,7 +161,7 @@ type AnalyticsTopObserversByIatum struct { Iata string `json:"iata"` ObserverID uuid.UUID `json:"observer_id"` Bucket pgtype.Timestamptz `json:"bucket"` - ObservationCount *int64 `json:"observation_count"` + ObservationCount int64 `json:"observation_count"` DisplayName *string `json:"display_name"` ObserverType *string `json:"observer_type"` } @@ -170,7 +170,7 @@ type AnalyticsTopTalkersByIatum struct { Iata string `json:"iata"` SenderName string `json:"sender_name"` Bucket pgtype.Timestamptz `json:"bucket"` - MessageCount *int64 `json:"message_count"` + MessageCount int64 `json:"message_count"` LastSent pgtype.Timestamptz `json:"last_sent"` } diff --git a/sqlc.yaml b/sqlc.yaml index 032c9911..b5e86b89 100644 --- a/sqlc.yaml +++ b/sqlc.yaml @@ -30,3 +30,6 @@ sql: go_type: {type: "string", pointer: true} - column: "mv_top_advertisers_by_iata.name" go_type: {type: "string", pointer: true} + + - column: "mv_observer_activity_hourly.payload_type" + go_type: {type: "int16", pointer: true} From 0c19e95132760eda1bdcf550ae1b01865bbf3c41 Mon Sep 17 00:00:00 2001 From: n30nex Date: Tue, 29 Sep 2026 21:14:51 -0400 Subject: [PATCH 3/3] fix: address release review follow-ups for #167 --- README.md | 2 ++ 1 file changed, 2 insertions(+) diff --git a/README.md b/README.md index 5243068a..823c04be 100644 --- a/README.md +++ b/README.md @@ -541,6 +541,8 @@ summaries retain 30 days independently of `packets.retention`. Cleanup saves onl aggregates before deleting each packet batch, in the same transaction. Raw packets, observations and message bodies still expire under packet retention. The summaries use UTC hourly buckets and appear on the normal view-refresh cycle. +`mv_hourly_iata_stats` now covers 30 days instead of migration 001's seven days; +`/stats/observations?since=` can therefore return retained summaries older than a week. Migration 039 starts from data still present; previously deleted history cannot be reconstructed. Telemetry has its own retention setting. Packet drill-down,