From d65b0676798103eab3dffa48defb1a27052df9a5 Mon Sep 17 00:00:00 2001 From: MrAlders0n Date: Tue, 21 Jul 2026 10:41:02 -0400 Subject: [PATCH 1/8] perf(channels): track channel activity per IATA in its own table The IATA filter ran a correlated EXISTS over packets/observations with ILIKE, which skipped the iata index and took ~7s live. Keep a small channel_iatas table at ingest (like node_iatas) and filter against it. Also honor the iatas= param so multi-site regions stop getting the global list. Filter now ages out with packet retention rather than matching any retained packet. --- db/channels.go | 22 +++++++-- db/channels_test.go | 42 +++++++++++----- db/migrations/011_channel_iatas.sql | 19 ++++++++ db/queries/queries.sql | 31 +++++++----- db/sqlc/mock/querier.go | 28 +++++++++++ db/sqlc/models.go | 6 +++ db/sqlc/querier.go | 8 +-- db/sqlc/queries.sql.go | 59 ++++++++++++++++------- docs/docs.go | 22 ++++++--- docs/swagger.json | 6 +++ docs/swagger.yaml | 4 ++ internal/api/handlers/channels.go | 5 +- internal/api/handlers/channels_test.go | 36 +++++++++++++- internal/api/handlers/stub_reader_test.go | 6 +-- internal/api/reader.go | 5 +- internal/background/tasks.go | 3 ++ internal/cache/cache_test.go | 2 +- internal/cache/reader.go | 4 +- internal/ingest/ingest.go | 3 ++ internal/ingest/ingest_test.go | 9 +++- internal/ingest/packet.go | 6 +++ internal/ingest/side_effects_test.go | 55 +++++++++++++++++++++ 22 files changed, 313 insertions(+), 68 deletions(-) create mode 100644 db/migrations/011_channel_iatas.sql diff --git a/db/channels.go b/db/channels.go index 094940d..9c652fc 100644 --- a/db/channels.go +++ b/db/channels.go @@ -47,16 +47,28 @@ func (s *Store) UpsertChannelHashOnly(ctx context.Context, channelHash []byte) ( return int(rowID), nil } -func (s *Store) ListChannels(ctx context.Context, limit int32, hash []byte, iata string, cursor int64) (api.Page[api.ChannelSummary], error) { +func (s *Store) UpsertChannelIATA(ctx context.Context, channelHash []byte, iata string, heardAt time.Time) error { + return s.q.UpsertChannelIATA(ctx, sqlc.UpsertChannelIATAParams{ + ChannelHash: channelHash, + Iata: iata, + LastHeard: pgtype.Timestamptz{Time: heardAt, Valid: true}, + }) +} + +func (s *Store) DeleteOldChannelIATAs(ctx context.Context, cutoff time.Time) error { + return s.q.DeleteOldChannelIATAs(ctx, pgtype.Timestamptz{Time: cutoff, Valid: true}) +} + +func (s *Store) ListChannels(ctx context.Context, limit int32, hash []byte, iatas []string, cursor int64) (api.Page[api.ChannelSummary], error) { var cursorTS pgtype.Timestamptz if cursor > 0 { cursorTS = pgtype.Timestamptz{Time: time.UnixMilli(cursor), Valid: true} } rows, err := s.q.ListChannels(ctx, sqlc.ListChannelsParams{ - Column1: hash, - Column2: iata, - Column3: cursorTS, - Limit: limit + 1, + ChannelHash: hash, + Iatas: iatas, + CursorTs: cursorTS, + PageLimit: limit + 1, }) if err != nil { return api.Page[api.ChannelSummary]{}, err diff --git a/db/channels_test.go b/db/channels_test.go index 66ee741..8fdea06 100644 --- a/db/channels_test.go +++ b/db/channels_test.go @@ -21,15 +21,15 @@ func TestListChannels_Empty(t *testing.T) { mock.EXPECT(). ListChannels(gomock.Any(), sqlc.ListChannelsParams{ - Column1: nil, - Column2: "", - Column3: pgtype.Timestamptz{}, - Limit: 11, + ChannelHash: nil, + Iatas: nil, + CursorTs: pgtype.Timestamptz{}, + PageLimit: 11, }). Return([]sqlc.Channel{}, nil) store := &Store{q: mock} - page, err := store.ListChannels(context.Background(), 10, nil, "", 0) + page, err := store.ListChannels(context.Background(), 10, nil, nil, 0) if err != nil { t.Fatalf("unexpected error: %v", err) } @@ -63,15 +63,15 @@ func TestListChannels_Pagination(t *testing.T) { mock.EXPECT(). ListChannels(gomock.Any(), sqlc.ListChannelsParams{ - Column1: nil, - Column2: "", - Column3: pgtype.Timestamptz{}, - Limit: 3, // limit+1 + ChannelHash: nil, + Iatas: nil, + CursorTs: pgtype.Timestamptz{}, + PageLimit: 3, // limit+1 }). Return(rows, nil) store := &Store{q: mock} - page, err := store.ListChannels(context.Background(), 2, nil, "", 0) + page, err := store.ListChannels(context.Background(), 2, nil, nil, 0) if err != nil { t.Fatalf("unexpected error: %v", err) } @@ -95,12 +95,32 @@ func TestListChannels_DBError(t *testing.T) { Return(nil, errors.New("db error")) store := &Store{q: mock} - _, err := store.ListChannels(context.Background(), 10, nil, "", 0) + _, err := store.ListChannels(context.Background(), 10, nil, nil, 0) if err == nil { t.Fatal("expected error, got nil") } } +func TestListChannels_IATAFilter(t *testing.T) { + ctrl := gomock.NewController(t) + mock := mockdb.NewMockQuerier(ctrl) + + mock.EXPECT(). + ListChannels(gomock.Any(), sqlc.ListChannelsParams{ + ChannelHash: nil, + Iatas: []string{"YOW", "YYZ"}, + CursorTs: pgtype.Timestamptz{}, + PageLimit: 11, + }). + Return([]sqlc.Channel{}, nil) + + store := &Store{q: mock} + _, err := store.ListChannels(context.Background(), 10, nil, []string{"YOW", "YYZ"}, 0) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } +} + func TestGetChannel_Basic(t *testing.T) { ctrl := gomock.NewController(t) mock := mockdb.NewMockQuerier(ctrl) diff --git a/db/migrations/011_channel_iatas.sql b/db/migrations/011_channel_iatas.sql new file mode 100644 index 0000000..833a249 --- /dev/null +++ b/db/migrations/011_channel_iatas.sql @@ -0,0 +1,19 @@ +-- Per-IATA channel activity so the IATA filter skips the ~7s EXISTS over packets. +-- Keyed by raw hash (channels can share one), so no FK to channels. + +CREATE TABLE channel_iatas ( + channel_hash BYTEA NOT NULL, + iata CHAR(3) NOT NULL REFERENCES iata_codes(iata) ON DELETE CASCADE, + last_heard TIMESTAMPTZ NOT NULL DEFAULT NOW(), + PRIMARY KEY (channel_hash, iata) +); + +CREATE INDEX idx_channel_iatas_iata ON channel_iatas(iata, last_heard DESC); + +-- Seed from retained packets so the filter works right away. +INSERT INTO channel_iatas (channel_hash, iata, last_heard) +SELECT p.channel_hash, po.iata, MAX(po.heard_at) +FROM packets p +JOIN packet_observations po ON po.packet_hash = p.packet_hash +WHERE p.channel_hash IS NOT NULL +GROUP BY p.channel_hash, po.iata; diff --git a/db/queries/queries.sql b/db/queries/queries.sql index 610525e..57e824f 100644 --- a/db/queries/queries.sql +++ b/db/queries/queries.sql @@ -494,6 +494,10 @@ LIMIT $6; -- packet_observations cascade-delete via FK. DELETE FROM packets WHERE last_heard_at < $1; +-- name: DeleteOldChannelIATAs :exec +-- Keeps the channel IATA filter in step with packet retention. +DELETE FROM channel_iatas WHERE last_heard < $1; + -- ============================================================ -- PACKET OBSERVATIONS -- ============================================================ @@ -667,22 +671,25 @@ ON CONFLICT (channel_hash) WHERE key_fingerprint IS NULL DO UPDATE SET last_seen = NOW() RETURNING id; +-- name: UpsertChannelIATA :exec +INSERT INTO channel_iatas (channel_hash, iata, last_heard) +VALUES ($1, $2, $3) +ON CONFLICT (channel_hash, iata) DO UPDATE SET + last_heard = GREATEST(channel_iatas.last_heard, EXCLUDED.last_heard); + -- name: ListChannels :many --- Returns channels ordered by last seen, optionally filtered by hash and/or IATA. --- Pass NULL for hash to skip hash filtering. Pass empty string for iata to skip IATA filtering. --- IATA filter returns channels that have active packets in that IATA (case-insensitive). +-- Channels ordered by last seen, optionally filtered by hash and/or IATAs +-- (membership via channel_iatas). NULL hash / empty array skip those filters. -- Pass cursor=0 to start from the beginning (cursor is last_seen epoch ms). -SELECT DISTINCT c.* FROM channels c -WHERE ($1::bytea IS NULL OR c.channel_hash = $1) - AND ($2 = '' OR EXISTS ( - SELECT 1 FROM packets p - JOIN packet_observations po ON po.packet_hash = p.packet_hash - WHERE p.channel_hash = c.channel_hash - AND po.iata ILIKE $2 +SELECT c.* FROM channels c +WHERE (@channel_hash::bytea IS NULL OR c.channel_hash = @channel_hash) + AND (COALESCE(cardinality(@iatas::bpchar[]), 0) = 0 OR c.channel_hash IN ( + SELECT ci.channel_hash FROM channel_iatas ci + WHERE ci.iata = ANY(@iatas::bpchar[]) )) - AND ($3::timestamptz IS NULL OR c.last_seen < $3) + AND (@cursor_ts::timestamptz IS NULL OR c.last_seen < @cursor_ts) ORDER BY c.last_seen DESC -LIMIT $4; +LIMIT @page_limit; -- name: GetChannelByID :one SELECT * FROM channels WHERE id = $1; diff --git a/db/sqlc/mock/querier.go b/db/sqlc/mock/querier.go index 3122b6c..7ccdc26 100644 --- a/db/sqlc/mock/querier.go +++ b/db/sqlc/mock/querier.go @@ -43,6 +43,20 @@ func (m *MockQuerier) EXPECT() *MockQuerierMockRecorder { return m.recorder } +// DeleteOldChannelIATAs mocks base method. +func (m *MockQuerier) DeleteOldChannelIATAs(ctx context.Context, lastHeard pgtype.Timestamptz) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "DeleteOldChannelIATAs", ctx, lastHeard) + ret0, _ := ret[0].(error) + return ret0 +} + +// DeleteOldChannelIATAs indicates an expected call of DeleteOldChannelIATAs. +func (mr *MockQuerierMockRecorder) DeleteOldChannelIATAs(ctx, lastHeard any) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "DeleteOldChannelIATAs", reflect.TypeOf((*MockQuerier)(nil).DeleteOldChannelIATAs), ctx, lastHeard) +} + // DeleteOldPackets mocks base method. func (m *MockQuerier) DeleteOldPackets(ctx context.Context, lastHeardAt pgtype.Timestamptz) error { m.ctrl.T.Helper() @@ -1213,6 +1227,20 @@ func (mr *MockQuerierMockRecorder) UpsertChannelHashOnly(ctx, channelHash any) * return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "UpsertChannelHashOnly", reflect.TypeOf((*MockQuerier)(nil).UpsertChannelHashOnly), ctx, channelHash) } +// UpsertChannelIATA mocks base method. +func (m *MockQuerier) UpsertChannelIATA(ctx context.Context, arg db.UpsertChannelIATAParams) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "UpsertChannelIATA", ctx, arg) + ret0, _ := ret[0].(error) + return ret0 +} + +// UpsertChannelIATA indicates an expected call of UpsertChannelIATA. +func (mr *MockQuerierMockRecorder) UpsertChannelIATA(ctx, arg any) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "UpsertChannelIATA", reflect.TypeOf((*MockQuerier)(nil).UpsertChannelIATA), ctx, arg) +} + // UpsertIATA mocks base method. func (m *MockQuerier) UpsertIATA(ctx context.Context, iata string) error { m.ctrl.T.Helper() diff --git a/db/sqlc/models.go b/db/sqlc/models.go index 833107e..e31b389 100644 --- a/db/sqlc/models.go +++ b/db/sqlc/models.go @@ -23,6 +23,12 @@ type Channel struct { MessageCount *int64 `json:"message_count"` } +type ChannelIata struct { + ChannelHash []byte `json:"channel_hash"` + Iata string `json:"iata"` + LastHeard pgtype.Timestamptz `json:"last_heard"` +} + type ChannelKey struct { ChannelID int32 `json:"channel_id"` KeyBytes []byte `json:"key_bytes"` diff --git a/db/sqlc/querier.go b/db/sqlc/querier.go index 9ef83b0..9b05de4 100644 --- a/db/sqlc/querier.go +++ b/db/sqlc/querier.go @@ -12,6 +12,8 @@ import ( ) type Querier interface { + // Keeps the channel IATA filter in step with packet retention. + DeleteOldChannelIATAs(ctx context.Context, lastHeard pgtype.Timestamptz) 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 @@ -97,9 +99,8 @@ type Querier interface { // Pass empty string for iata or scope to skip those filters. // Pass cursor=0 to start from the beginning. ListChannelMessagesByHash(ctx context.Context, arg ListChannelMessagesByHashParams) ([]ListChannelMessagesByHashRow, error) - // Returns channels ordered by last seen, optionally filtered by hash and/or IATA. - // Pass NULL for hash to skip hash filtering. Pass empty string for iata to skip IATA filtering. - // IATA filter returns channels that have active packets in that IATA (case-insensitive). + // Channels ordered by last seen, optionally filtered by hash and/or IATAs + // (membership via channel_iatas). NULL hash / empty array skip those filters. // Pass cursor=0 to start from the beginning (cursor is last_seen epoch ms). ListChannels(ctx context.Context, arg ListChannelsParams) ([]Channel, error) ListIATAs(ctx context.Context) ([]IataCode, error) @@ -181,6 +182,7 @@ type Querier interface { // hash-only records (key unknown). Returns the channel row. UpsertChannel(ctx context.Context, arg UpsertChannelParams) (Channel, error) UpsertChannelHashOnly(ctx context.Context, channelHash []byte) (int32, error) + UpsertChannelIATA(ctx context.Context, arg UpsertChannelIATAParams) error // Copyright 2026 Beacon Contributors // SPDX-License-Identifier: agpl // ============================================================ diff --git a/db/sqlc/queries.sql.go b/db/sqlc/queries.sql.go index d2332a0..e7080c9 100644 --- a/db/sqlc/queries.sql.go +++ b/db/sqlc/queries.sql.go @@ -12,6 +12,16 @@ import ( "github.com/jackc/pgx/v5/pgtype" ) +const deleteOldChannelIATAs = `-- name: DeleteOldChannelIATAs :exec +DELETE FROM channel_iatas WHERE last_heard < $1 +` + +// Keeps the channel IATA filter in step with packet retention. +func (q *Queries) DeleteOldChannelIATAs(ctx context.Context, lastHeard pgtype.Timestamptz) error { + _, err := q.db.Exec(ctx, deleteOldChannelIATAs, lastHeard) + return err +} + const deleteOldPackets = `-- name: DeleteOldPackets :exec DELETE FROM packets WHERE last_heard_at < $1 ` @@ -1883,13 +1893,11 @@ func (q *Queries) ListChannelMessagesByHash(ctx context.Context, arg ListChannel } const listChannels = `-- name: ListChannels :many -SELECT DISTINCT c.id, c.channel_hash, c.key_fingerprint, c.name, c.hashtag, c.is_hashtag, c.is_public, c.key_known, c.first_seen, c.last_seen, c.message_count FROM channels c +SELECT c.id, c.channel_hash, c.key_fingerprint, c.name, c.hashtag, c.is_hashtag, c.is_public, c.key_known, c.first_seen, c.last_seen, c.message_count FROM channels c WHERE ($1::bytea IS NULL OR c.channel_hash = $1) - AND ($2 = '' OR EXISTS ( - SELECT 1 FROM packets p - JOIN packet_observations po ON po.packet_hash = p.packet_hash - WHERE p.channel_hash = c.channel_hash - AND po.iata ILIKE $2 + AND (COALESCE(cardinality($2::bpchar[]), 0) = 0 OR c.channel_hash IN ( + SELECT ci.channel_hash FROM channel_iatas ci + WHERE ci.iata = ANY($2::bpchar[]) )) AND ($3::timestamptz IS NULL OR c.last_seen < $3) ORDER BY c.last_seen DESC @@ -1897,22 +1905,21 @@ LIMIT $4 ` type ListChannelsParams struct { - Column1 []byte `json:"column_1"` - Column2 interface{} `json:"column_2"` - Column3 pgtype.Timestamptz `json:"column_3"` - Limit int32 `json:"limit"` + ChannelHash []byte `json:"channel_hash"` + Iatas []string `json:"iatas"` + CursorTs pgtype.Timestamptz `json:"cursor_ts"` + PageLimit int32 `json:"page_limit"` } -// Returns channels ordered by last seen, optionally filtered by hash and/or IATA. -// Pass NULL for hash to skip hash filtering. Pass empty string for iata to skip IATA filtering. -// IATA filter returns channels that have active packets in that IATA (case-insensitive). +// Channels ordered by last seen, optionally filtered by hash and/or IATAs +// (membership via channel_iatas). NULL hash / empty array skip those filters. // Pass cursor=0 to start from the beginning (cursor is last_seen epoch ms). func (q *Queries) ListChannels(ctx context.Context, arg ListChannelsParams) ([]Channel, error) { rows, err := q.db.Query(ctx, listChannels, - arg.Column1, - arg.Column2, - arg.Column3, - arg.Limit, + arg.ChannelHash, + arg.Iatas, + arg.CursorTs, + arg.PageLimit, ) if err != nil { return nil, err @@ -3556,6 +3563,24 @@ func (q *Queries) UpsertChannelHashOnly(ctx context.Context, channelHash []byte) return id, err } +const upsertChannelIATA = `-- name: UpsertChannelIATA :exec +INSERT INTO channel_iatas (channel_hash, iata, last_heard) +VALUES ($1, $2, $3) +ON CONFLICT (channel_hash, iata) DO UPDATE SET + last_heard = GREATEST(channel_iatas.last_heard, EXCLUDED.last_heard) +` + +type UpsertChannelIATAParams struct { + ChannelHash []byte `json:"channel_hash"` + Iata string `json:"iata"` + LastHeard pgtype.Timestamptz `json:"last_heard"` +} + +func (q *Queries) UpsertChannelIATA(ctx context.Context, arg UpsertChannelIATAParams) error { + _, err := q.db.Exec(ctx, upsertChannelIATA, arg.ChannelHash, arg.Iata, arg.LastHeard) + return err +} + const upsertIATA = `-- name: UpsertIATA :exec diff --git a/docs/docs.go b/docs/docs.go index 288eb58..0f8b74d 100644 --- a/docs/docs.go +++ b/docs/docs.go @@ -4,11 +4,11 @@ package docs import "github.com/swaggo/swag" const docTemplate = `{ - "schemes": [[ marshal .Schemes ]], + "schemes": {{ marshal .Schemes }}, "swagger": "2.0", "info": { - "description": "[[escape .Description]]", - "title": "[[.Title]]", + "description": "{{escape .Description}}", + "title": "{{.Title}}", "termsOfService": "https://github.com/MeshCore-Beacon/beacon-server", "contact": { "name": "MeshCore Beacon", @@ -17,10 +17,10 @@ const docTemplate = `{ "license": { "name": "AGPL-3-or-later" }, - "version": "[[.Version]]" + "version": "{{.Version}}" }, - "host": "[[.Host]]", - "basePath": "[[.BasePath]]", + "host": "{{.Host}}", + "basePath": "{{.BasePath}}", "paths": { "/brokers": { "get": { @@ -66,6 +66,12 @@ const docTemplate = `{ "name": "iata", "in": "query" }, + { + "type": "string", + "description": "Filter by IATA code(s), comma-separated e.g. YOW or YOW,YYZ", + "name": "iatas", + "in": "query" + }, { "type": "integer", "description": "last_seen epoch ms of last item for pagination", @@ -3768,8 +3774,8 @@ var SwaggerInfo = &swag.Spec{ Description: "MeshCore network observation backend. Ingests LoRa packets from MQTT brokers, stores in PostgreSQL, and streams live events via WebSocket.", InfoInstanceName: "swagger", SwaggerTemplate: docTemplate, - LeftDelim: "[[", - RightDelim: "]]", + LeftDelim: "{{", + RightDelim: "}}", } func init() { diff --git a/docs/swagger.json b/docs/swagger.json index 479b58b..984ff47 100644 --- a/docs/swagger.json +++ b/docs/swagger.json @@ -64,6 +64,12 @@ "name": "iata", "in": "query" }, + { + "type": "string", + "description": "Filter by IATA code(s), comma-separated e.g. YOW or YOW,YYZ", + "name": "iatas", + "in": "query" + }, { "type": "integer", "description": "last_seen epoch ms of last item for pagination", diff --git a/docs/swagger.yaml b/docs/swagger.yaml index d63e219..f846b68 100644 --- a/docs/swagger.yaml +++ b/docs/swagger.yaml @@ -1102,6 +1102,10 @@ paths: in: query name: iata type: string + - description: Filter by IATA code(s), comma-separated e.g. YOW or YOW,YYZ + in: query + name: iatas + type: string - description: last_seen epoch ms of last item for pagination in: query name: cursor diff --git a/internal/api/handlers/channels.go b/internal/api/handlers/channels.go index ed0d533..d1652f7 100644 --- a/internal/api/handlers/channels.go +++ b/internal/api/handlers/channels.go @@ -36,6 +36,7 @@ func ChannelsRouter(reader api.Reader) http.Handler { // @Produce json // @Param hash query string false "Single-byte channel hash (hex)" // @Param iata query string false "Filter by IATA code (case-insensitive)" +// @Param iatas query string false "Filter by IATA code(s), comma-separated e.g. YOW or YOW,YYZ" // @Param cursor query int false "last_seen epoch ms of last item for pagination" // @Param limit query int false "Max results (default 50)" // @Success 200 {object} api.Page[api.ChannelSummary] @@ -53,7 +54,7 @@ func listChannels(reader api.Reader) http.HandlerFunc { } limit = l } - iata := r.URL.Query().Get("iata") + iatas := parseIATAs(r) var cursor int64 if cursorParam := r.URL.Query().Get("cursor"); cursorParam != "" { c, err := strconv.ParseInt(cursorParam, 10, 64) @@ -76,7 +77,7 @@ func listChannels(reader api.Reader) http.HandlerFunc { } hashHex = h } - channels, err := reader.ListChannels(r.Context(), int32(limit), hashHex, iata, cursor) + channels, err := reader.ListChannels(r.Context(), int32(limit), hashHex, iatas, cursor) if err != nil { respondError(w, http.StatusInternalServerError, "internal server error") return diff --git a/internal/api/handlers/channels_test.go b/internal/api/handlers/channels_test.go index 360f0da..c014dc4 100644 --- a/internal/api/handlers/channels_test.go +++ b/internal/api/handlers/channels_test.go @@ -7,6 +7,7 @@ import ( "context" "net/http" "net/http/httptest" + "reflect" "testing" "github.com/MeshCore-Beacon/beacon-server/internal/api" @@ -115,7 +116,7 @@ func TestListChannelMessages_InvalidCursor(t *testing.T) { func TestListChannels_OK(t *testing.T) { r := chi.NewRouter() r.Get("/channels", listChannels(stubReader{ - listChannels: func(_ context.Context, _ int32, _ []byte, _ string, _ int64) (api.Page[api.ChannelSummary], error) { + listChannels: func(_ context.Context, _ int32, _ []byte, _ []string, _ int64) (api.Page[api.ChannelSummary], error) { return api.Page[api.ChannelSummary]{Items: []api.ChannelSummary{{ID: 1, ChannelHash: "ab"}}}, nil }, })) @@ -127,6 +128,39 @@ func TestListChannels_OK(t *testing.T) { } } +func TestListChannels_IATAParsing(t *testing.T) { + cases := []struct { + name string + query string + want []string + }{ + {"single lowercased", "?iata=yow", []string{"YOW"}}, + {"multi csv", "?iatas=yow,%20yyz", []string{"YOW", "YYZ"}}, + {"none", "", nil}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + var got []string + r := chi.NewRouter() + r.Get("/channels", listChannels(stubReader{ + listChannels: func(_ context.Context, _ int32, _ []byte, iatas []string, _ int64) (api.Page[api.ChannelSummary], error) { + got = iatas + return api.Page[api.ChannelSummary]{}, nil + }, + })) + req := httptest.NewRequest(http.MethodGet, "/channels"+tc.query, nil) + w := httptest.NewRecorder() + r.ServeHTTP(w, req) + if w.Code != http.StatusOK { + t.Fatalf("expected 200, got %d", w.Code) + } + if !reflect.DeepEqual(got, tc.want) { + t.Errorf("expected iatas %v, got %v", tc.want, got) + } + }) + } +} + func TestGetChannel_OK(t *testing.T) { r := chi.NewRouter() r.Get("/channels/{channelID}", getChannel(stubReader{ diff --git a/internal/api/handlers/stub_reader_test.go b/internal/api/handlers/stub_reader_test.go index 948d57a..cdd575d 100644 --- a/internal/api/handlers/stub_reader_test.go +++ b/internal/api/handlers/stub_reader_test.go @@ -20,7 +20,7 @@ type stubReader struct { listRegions func(ctx context.Context) ([]api.RegionSummary, error) getRegion func(ctx context.Context, regionID int32) (*api.Region, error) getRegionBySlug func(ctx context.Context, slug string) (*api.Region, error) - listChannels func(ctx context.Context, limit int32, hash []byte, iata string, cursor int64) (api.Page[api.ChannelSummary], error) + listChannels func(ctx context.Context, limit int32, hash []byte, iatas []string, cursor int64) (api.Page[api.ChannelSummary], error) getChannel func(ctx context.Context, channelID int32) (*api.Channel, error) listChannelMessages func(ctx context.Context, channelID *int32, since time.Time, limit int32, iatas []string, scope string, cursor int64) (api.Page[api.ChannelMessage], error) listChannelMessagesByHash func(ctx context.Context, hash []byte, since time.Time, limit int32, iatas []string, scope string, cursor int64) (api.Page[api.ChannelMessage], error) @@ -96,9 +96,9 @@ func (s stubReader) GetRegionBySlug(ctx context.Context, slug string) (*api.Regi return nil, nil } -func (s stubReader) ListChannels(ctx context.Context, limit int32, hash []byte, iata string, cursor int64) (api.Page[api.ChannelSummary], error) { +func (s stubReader) ListChannels(ctx context.Context, limit int32, hash []byte, iatas []string, cursor int64) (api.Page[api.ChannelSummary], error) { if s.listChannels != nil { - return s.listChannels(ctx, limit, hash, iata, cursor) + return s.listChannels(ctx, limit, hash, iatas, cursor) } return api.Page[api.ChannelSummary]{}, nil } diff --git a/internal/api/reader.go b/internal/api/reader.go index ad95e58..9309f57 100644 --- a/internal/api/reader.go +++ b/internal/api/reader.go @@ -44,9 +44,10 @@ type Reader interface { // ListChannels returns a paginated list of channels ordered by last seen. // Includes both hashtag-derived and explicit key channels. - // Pass nil hash to skip hash filtering. Pass empty string iata to return all channels. + // Pass nil hash to skip hash filtering. Pass empty iatas to return all channels; + // IATAs must be uppercase. // cursor is last_seen epoch ms of the last item; pass 0 to start from the beginning. - ListChannels(ctx context.Context, limit int32, hash []byte, iata string, cursor int64) (Page[ChannelSummary], error) + ListChannels(ctx context.Context, limit int32, hash []byte, iatas []string, cursor int64) (Page[ChannelSummary], error) // GetChannel returns full detail for a single channel by its integer ID. // Returns nil, pgx.ErrNoRows if the channel is not found. diff --git a/internal/background/tasks.go b/internal/background/tasks.go index 6006cce..559a3dd 100644 --- a/internal/background/tasks.go +++ b/internal/background/tasks.go @@ -44,6 +44,9 @@ func CleanupTask(store *db.Store, telemetryRetention, packetRetention, interval if err := store.DeleteOldPackets(ctx, time.Now().Add(-packetRetention)); err != nil { return err } + if err := store.DeleteOldChannelIATAs(ctx, time.Now().Add(-packetRetention)); err != nil { + return err + } return nil }, } diff --git a/internal/cache/cache_test.go b/internal/cache/cache_test.go index 8ab4337..7ef6954 100644 --- a/internal/cache/cache_test.go +++ b/internal/cache/cache_test.go @@ -128,7 +128,7 @@ func (s *stubReader) GetCrossIATANeighbors(_ context.Context, _ uuid.UUID, _ str return nil, nil } -func (s *stubReader) ListChannels(_ context.Context, _ int32, _ []byte, _ string, _ int64) (api.Page[api.ChannelSummary], error) { +func (s *stubReader) ListChannels(_ context.Context, _ int32, _ []byte, _ []string, _ int64) (api.Page[api.ChannelSummary], error) { return api.Page[api.ChannelSummary]{}, nil } diff --git a/internal/cache/reader.go b/internal/cache/reader.go index 46d001e..4bdfd8c 100644 --- a/internal/cache/reader.go +++ b/internal/cache/reader.go @@ -354,8 +354,8 @@ func (cr *CachedReader) GetCrossIATANeighbors(ctx context.Context, nodeID uuid.U } // ListChannels implements [api.Reader]. -func (cr *CachedReader) ListChannels(ctx context.Context, limit int32, hash []byte, iata string, cursor int64) (api.Page[api.ChannelSummary], error) { - return cr.inner.ListChannels(ctx, limit, hash, iata, cursor) +func (cr *CachedReader) ListChannels(ctx context.Context, limit int32, hash []byte, iatas []string, cursor int64) (api.Page[api.ChannelSummary], error) { + return cr.inner.ListChannels(ctx, limit, hash, iatas, cursor) } // ListChannelMessages implements [api.Reader]. diff --git a/internal/ingest/ingest.go b/internal/ingest/ingest.go index 5bf7453..48310aa 100644 --- a/internal/ingest/ingest.go +++ b/internal/ingest/ingest.go @@ -151,6 +151,9 @@ type DB interface { // but can be safely ignored since unknown-key channels have no messages. UpsertChannelHashOnly(ctx context.Context, channelHash []byte) (int, error) + // UpsertChannelIATA upserts a channel_iatas row. + UpsertChannelIATA(ctx context.Context, channelHash []byte, iata string, heardAt time.Time) error + // GetPacketObservationCount returns the number of rows for the packet observations GetPacketObservationCount(ctx context.Context, packetHash []byte) (int64, error) diff --git a/internal/ingest/ingest_test.go b/internal/ingest/ingest_test.go index 15c54fc..81943d5 100644 --- a/internal/ingest/ingest_test.go +++ b/internal/ingest/ingest_test.go @@ -177,6 +177,8 @@ type stubDB struct { upsertNodeCalls int upsertChannelCalls int upsertChannelHashOnlyCalls int + upsertChannelIATACalls int + observationInserted bool } type setCapabilityCall struct { @@ -201,7 +203,7 @@ func (s *stubDB) UpsertPacket(_ context.Context, _ UpsertPacketParams) (bool, er } func (s *stubDB) SetPacketDecrypted(_ context.Context, _ []byte) error { return nil } func (s *stubDB) InsertObservation(_ context.Context, _ InsertObservationParams) (bool, error) { - return false, nil + return s.observationInserted, nil } func (s *stubDB) SetNodeDefaultScope(_ context.Context, _ uuid.UUID, _ int32) error { return nil } func (s *stubDB) UpsertNode(_ context.Context, _ UpsertNodeParams, _ RadioSettings) (uuid.UUID, error) { @@ -257,6 +259,11 @@ func (s *stubDB) UpsertChannelHashOnly(_ context.Context, _ []byte) (int, error) return 0, nil } +func (s *stubDB) UpsertChannelIATA(_ context.Context, _ []byte, _ string, _ time.Time) error { + s.upsertChannelIATACalls++ + return nil +} + func (s *stubDB) GetPacketObservationCount(_ context.Context, _ []byte) (int64, error) { return 0, nil } diff --git a/internal/ingest/packet.go b/internal/ingest/packet.go index bf64f29..3ff8d3b 100644 --- a/internal/ingest/packet.go +++ b/internal/ingest/packet.go @@ -783,6 +783,12 @@ func (w *Worker) handlePacket(ctx context.Context, iata, pubkeyHex string, raw [ } } + if channelHash != nil && inserted { + if err := w.db.UpsertChannelIATA(ctx, channelHash, iata, heardAt); err != nil { + log.Printf("ingest[%s]: db: upsert channel IATA failed from %s/%s: %v", w.cfg.BrokerName, iata, pubkeyHex, err) + } + } + // packet.PathHashes() reads packet.Path as hash-sized chunks, which is only true for // ordinary flood/direct-routed packets. TRACE repurposes packet.Path to carry one SNR // byte per hop instead, so for TRACE we resolve against the trace payload's own diff --git a/internal/ingest/side_effects_test.go b/internal/ingest/side_effects_test.go index acf0e56..a62a954 100644 --- a/internal/ingest/side_effects_test.go +++ b/internal/ingest/side_effects_test.go @@ -6,7 +6,10 @@ package ingest import ( "context" "crypto/ed25519" + "encoding/hex" + "encoding/json" "testing" + "time" "github.com/MeshCore-Beacon/beacon-server/internal/keystore" "github.com/meshcore-go/meshcore-go" @@ -133,3 +136,55 @@ func TestHandlePayloadTypeSideEffects_GrpTxt_UnknownKey_OnlyUpsertsHashOnlyChann t.Errorf("expected UpsertChannel NOT to be called when the key is unknown, got %d calls", db.upsertChannelCalls) } } + +// packetEnvelope wraps a packet in the minimal broker JSON that handlePacket expects. +func packetEnvelope(t *testing.T, packet *meshcore.Packet) []byte { + t.Helper() + raw, err := packet.ToBytes() + if err != nil { + t.Fatalf("packet to bytes: %v", err) + } + env, err := json.Marshal(map[string]string{ + "raw": hex.EncodeToString(raw), + "timestamp": time.Now().UTC().Format("2006-01-02T15:04:05.000000"), + }) + if err != nil { + t.Fatalf("marshal envelope: %v", err) + } + return env +} + +func TestHandlePacket_GrpTxt_UpsertsChannelIATA(t *testing.T) { + w, db := newTestWorker() + db.observationInserted = true + envelope := packetEnvelope(t, buildGrpTxtPacket(t, 0x1a, make([]byte, 16))) + + w.handlePacket(context.Background(), "YOW", "0102", envelope) + + if db.upsertChannelIATACalls != 1 { + t.Errorf("expected UpsertChannelIATA to be called once for a stored group text, got %d", db.upsertChannelIATACalls) + } +} + +func TestHandlePacket_GrpTxt_DedupObservation_SkipsChannelIATA(t *testing.T) { + w, db := newTestWorker() // stub reports the observation as a duplicate + envelope := packetEnvelope(t, buildGrpTxtPacket(t, 0x1a, make([]byte, 16))) + + w.handlePacket(context.Background(), "YOW", "0102", envelope) + + if db.upsertChannelIATACalls != 0 { + t.Errorf("expected UpsertChannelIATA NOT to be called for a duplicate observation, got %d calls", db.upsertChannelIATACalls) + } +} + +func TestHandlePacket_Advert_SkipsChannelIATA(t *testing.T) { + w, db := newTestWorker() + db.observationInserted = true + envelope := packetEnvelope(t, buildAdvertPacket(t, false)) + + w.handlePacket(context.Background(), "YOW", "0102", envelope) + + if db.upsertChannelIATACalls != 0 { + t.Errorf("expected UpsertChannelIATA NOT to be called for a non-channel packet, got %d calls", db.upsertChannelIATACalls) + } +} From fecb082a0c0579ee44674df9d0838c2218aaecfb Mon Sep 17 00:00:00 2001 From: MrAlders0n Date: Thu, 23 Jul 2026 13:08:18 -0400 Subject: [PATCH 2/8] fix(channels): review fixes for channel_iatas Refresh last_heard on duplicate observations too, capped at hourly, so steady traffic can't age a channel out of the filter while its packets stay retained (dedup key has no heard_at). Guard the seed join against the docker /dev/shm cap like the trace one, and drop the unused last_heard index column so the upserts stay HOT. --- db/migrations/011_channel_iatas.sql | 8 ++++++-- db/queries/queries.sql | 4 +++- db/sqlc/querier.go | 1 + db/sqlc/queries.sql.go | 4 +++- docs/docs.go | 2 +- docs/swagger.json | 2 +- docs/swagger.yaml | 2 +- internal/api/handlers/channels.go | 2 +- internal/ingest/packet.go | 3 ++- internal/ingest/side_effects_test.go | 6 +++--- 10 files changed, 22 insertions(+), 12 deletions(-) diff --git a/db/migrations/011_channel_iatas.sql b/db/migrations/011_channel_iatas.sql index 833a249..a89885d 100644 --- a/db/migrations/011_channel_iatas.sql +++ b/db/migrations/011_channel_iatas.sql @@ -8,12 +8,16 @@ CREATE TABLE channel_iatas ( PRIMARY KEY (channel_hash, iata) ); -CREATE INDEX idx_channel_iatas_iata ON channel_iatas(iata, last_heard DESC); +CREATE INDEX idx_channel_iatas_iata ON channel_iatas(iata); + +-- Seed from retained packets; parallelism off so the join spills to disk, not /dev/shm. +SET max_parallel_workers_per_gather = 0; --- Seed from retained packets so the filter works right away. INSERT INTO channel_iatas (channel_hash, iata, last_heard) SELECT p.channel_hash, po.iata, MAX(po.heard_at) FROM packets p JOIN packet_observations po ON po.packet_hash = p.packet_hash WHERE p.channel_hash IS NOT NULL GROUP BY p.channel_hash, po.iata; + +RESET max_parallel_workers_per_gather; diff --git a/db/queries/queries.sql b/db/queries/queries.sql index 57e824f..792d294 100644 --- a/db/queries/queries.sql +++ b/db/queries/queries.sql @@ -672,10 +672,12 @@ ON CONFLICT (channel_hash) WHERE key_fingerprint IS NULL DO UPDATE SET RETURNING id; -- name: UpsertChannelIATA :exec +-- Refreshes at most hourly so repeat hears don't churn the row. INSERT INTO channel_iatas (channel_hash, iata, last_heard) VALUES ($1, $2, $3) ON CONFLICT (channel_hash, iata) DO UPDATE SET - last_heard = GREATEST(channel_iatas.last_heard, EXCLUDED.last_heard); + last_heard = EXCLUDED.last_heard +WHERE EXCLUDED.last_heard > channel_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 9b05de4..ef36ddf 100644 --- a/db/sqlc/querier.go +++ b/db/sqlc/querier.go @@ -182,6 +182,7 @@ type Querier interface { // hash-only records (key unknown). Returns the channel row. UpsertChannel(ctx context.Context, arg UpsertChannelParams) (Channel, error) UpsertChannelHashOnly(ctx context.Context, channelHash []byte) (int32, error) + // Refreshes at most hourly so repeat hears don't churn the row. UpsertChannelIATA(ctx context.Context, arg UpsertChannelIATAParams) error // Copyright 2026 Beacon Contributors // SPDX-License-Identifier: agpl diff --git a/db/sqlc/queries.sql.go b/db/sqlc/queries.sql.go index e7080c9..1fe3ac6 100644 --- a/db/sqlc/queries.sql.go +++ b/db/sqlc/queries.sql.go @@ -3567,7 +3567,8 @@ const upsertChannelIATA = `-- name: UpsertChannelIATA :exec INSERT INTO channel_iatas (channel_hash, iata, last_heard) VALUES ($1, $2, $3) ON CONFLICT (channel_hash, iata) DO UPDATE SET - last_heard = GREATEST(channel_iatas.last_heard, EXCLUDED.last_heard) + last_heard = EXCLUDED.last_heard +WHERE EXCLUDED.last_heard > channel_iatas.last_heard + INTERVAL '1 hour' ` type UpsertChannelIATAParams struct { @@ -3576,6 +3577,7 @@ type UpsertChannelIATAParams struct { LastHeard pgtype.Timestamptz `json:"last_heard"` } +// Refreshes at most hourly so repeat hears don't churn the row. func (q *Queries) UpsertChannelIATA(ctx context.Context, arg UpsertChannelIATAParams) error { _, err := q.db.Exec(ctx, upsertChannelIATA, arg.ChannelHash, arg.Iata, arg.LastHeard) return err diff --git a/docs/docs.go b/docs/docs.go index 0f8b74d..4ce0893 100644 --- a/docs/docs.go +++ b/docs/docs.go @@ -62,7 +62,7 @@ const docTemplate = `{ }, { "type": "string", - "description": "Filter by IATA code (case-insensitive)", + "description": "Filter by IATA code", "name": "iata", "in": "query" }, diff --git a/docs/swagger.json b/docs/swagger.json index 984ff47..92312e1 100644 --- a/docs/swagger.json +++ b/docs/swagger.json @@ -60,7 +60,7 @@ }, { "type": "string", - "description": "Filter by IATA code (case-insensitive)", + "description": "Filter by IATA code", "name": "iata", "in": "query" }, diff --git a/docs/swagger.yaml b/docs/swagger.yaml index f846b68..98154ca 100644 --- a/docs/swagger.yaml +++ b/docs/swagger.yaml @@ -1098,7 +1098,7 @@ paths: in: query name: hash type: string - - description: Filter by IATA code (case-insensitive) + - description: Filter by IATA code in: query name: iata type: string diff --git a/internal/api/handlers/channels.go b/internal/api/handlers/channels.go index d1652f7..f394058 100644 --- a/internal/api/handlers/channels.go +++ b/internal/api/handlers/channels.go @@ -35,7 +35,7 @@ func ChannelsRouter(reader api.Reader) http.Handler { // @Tags Channels // @Produce json // @Param hash query string false "Single-byte channel hash (hex)" -// @Param iata query string false "Filter by IATA code (case-insensitive)" +// @Param iata query string false "Filter by IATA code" // @Param iatas query string false "Filter by IATA code(s), comma-separated e.g. YOW or YOW,YYZ" // @Param cursor query int false "last_seen epoch ms of last item for pagination" // @Param limit query int false "Max results (default 50)" diff --git a/internal/ingest/packet.go b/internal/ingest/packet.go index 3ff8d3b..741bb93 100644 --- a/internal/ingest/packet.go +++ b/internal/ingest/packet.go @@ -783,7 +783,8 @@ func (w *Worker) handlePacket(ctx context.Context, iata, pubkeyHex string, raw [ } } - if channelHash != nil && inserted { + // Runs on duplicate observations too; the upsert only writes when the row is >1h stale. + if channelHash != nil { if err := w.db.UpsertChannelIATA(ctx, channelHash, iata, heardAt); err != nil { log.Printf("ingest[%s]: db: upsert channel IATA failed from %s/%s: %v", w.cfg.BrokerName, iata, pubkeyHex, err) } diff --git a/internal/ingest/side_effects_test.go b/internal/ingest/side_effects_test.go index a62a954..077caef 100644 --- a/internal/ingest/side_effects_test.go +++ b/internal/ingest/side_effects_test.go @@ -166,14 +166,14 @@ func TestHandlePacket_GrpTxt_UpsertsChannelIATA(t *testing.T) { } } -func TestHandlePacket_GrpTxt_DedupObservation_SkipsChannelIATA(t *testing.T) { +func TestHandlePacket_GrpTxt_DedupObservation_StillUpsertsChannelIATA(t *testing.T) { w, db := newTestWorker() // stub reports the observation as a duplicate envelope := packetEnvelope(t, buildGrpTxtPacket(t, 0x1a, make([]byte, 16))) w.handlePacket(context.Background(), "YOW", "0102", envelope) - if db.upsertChannelIATACalls != 0 { - t.Errorf("expected UpsertChannelIATA NOT to be called for a duplicate observation, got %d calls", db.upsertChannelIATACalls) + if db.upsertChannelIATACalls != 1 { + t.Errorf("expected UpsertChannelIATA to run for a duplicate observation too, got %d calls", db.upsertChannelIATACalls) } } From 51acd409fb6305fbf2dba0f6bb1d1c733b933796 Mon Sep 17 00:00:00 2001 From: MrAlders0n Date: Thu, 23 Jul 2026 12:14:41 -0400 Subject: [PATCH 3/8] perf(traces): track trace activity per IATA in its own table The trace list joined every trace packet to ~50M observations to apply the IATA filter; the parallel hash join overran docker's 64MB /dev/shm and the filtered list errored. Keep a small trace_iatas table at ingest (like channel_iatas) and group packets by tag instead. Filtered stats are now per tag rather than per matching packet. --- db/migrations/012_trace_iatas.sql | 23 +++++++ db/queries/queries.sql | 74 ++++++++++++++-------- db/sqlc/mock/querier.go | 28 +++++++++ db/sqlc/models.go | 6 ++ db/sqlc/querier.go | 5 ++ db/sqlc/queries.sql.go | 92 ++++++++++++++++++++-------- db/traces.go | 12 ++++ internal/background/tasks.go | 3 + internal/ingest/ingest.go | 3 + internal/ingest/ingest_test.go | 6 ++ internal/ingest/packet.go | 6 ++ internal/ingest/side_effects_test.go | 27 ++++++++ 12 files changed, 233 insertions(+), 52 deletions(-) create mode 100644 db/migrations/012_trace_iatas.sql diff --git a/db/migrations/012_trace_iatas.sql b/db/migrations/012_trace_iatas.sql new file mode 100644 index 0000000..f26e05d --- /dev/null +++ b/db/migrations/012_trace_iatas.sql @@ -0,0 +1,23 @@ +-- Per-IATA trace activity so the trace filter stops joining every trace packet +-- to observations (the hash join was overrunning /dev/shm). + +CREATE TABLE trace_iatas ( + trace_tag BYTEA NOT NULL, + iata CHAR(3) NOT NULL REFERENCES iata_codes(iata) ON DELETE CASCADE, + last_heard TIMESTAMPTZ NOT NULL DEFAULT NOW(), + PRIMARY KEY (trace_tag, iata) +); + +CREATE INDEX idx_trace_iatas_iata ON trace_iatas(iata, last_heard DESC); + +-- Seed from retained packets; parallelism off so the join spills to disk, not /dev/shm. +SET max_parallel_workers_per_gather = 0; + +INSERT INTO trace_iatas (trace_tag, iata, last_heard) +SELECT p.trace_tag, po.iata, MAX(po.heard_at) +FROM packets p +JOIN packet_observations po ON po.packet_hash = p.packet_hash +WHERE p.trace_tag IS NOT NULL +GROUP BY p.trace_tag, po.iata; + +RESET max_parallel_workers_per_gather; diff --git a/db/queries/queries.sql b/db/queries/queries.sql index 792d294..698ab3a 100644 --- a/db/queries/queries.sql +++ b/db/queries/queries.sql @@ -498,6 +498,10 @@ DELETE FROM packets WHERE last_heard_at < $1; -- Keeps the channel IATA filter in step with packet retention. DELETE FROM channel_iatas WHERE last_heard < $1; +-- name: DeleteOldTraceIATAs :exec +-- Keeps the trace IATA filter in step with packet retention. +DELETE FROM trace_iatas WHERE last_heard < $1; + -- ============================================================ -- PACKET OBSERVATIONS -- ============================================================ @@ -679,6 +683,12 @@ ON CONFLICT (channel_hash, iata) DO UPDATE SET last_heard = EXCLUDED.last_heard WHERE EXCLUDED.last_heard > channel_iatas.last_heard + INTERVAL '1 hour'; +-- 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); + -- name: ListChannels :many -- Channels ordered by last seen, optionally filtered by hash and/or IATAs -- (membership via channel_iatas). NULL hash / empty array skip those filters. @@ -965,33 +975,45 @@ ON CONFLICT (region_id, iata) DO NOTHING; -- name: ListTraceTags :many -- Returns distinct trace tags with summary info, ordered by most recent first. +-- IATA membership comes from trace_iatas (joining observations here spilled the +-- hash join). Per-tag details filled in only for the returned page. +WITH tags AS ( + SELECT + p.trace_tag, + MIN(p.first_heard_at) AS first_heard_at, + MAX(p.last_heard_at) AS last_heard_at, + COUNT(*) AS packet_count, + MAX(p.parsed_payload->>'type') AS trace_type + FROM packets p + WHERE p.trace_tag IS NOT NULL + AND (COALESCE(cardinality($1::bpchar[]), 0) = 0 OR p.trace_tag IN ( + SELECT ti.trace_tag FROM trace_iatas ti WHERE ti.iata = ANY($1::bpchar[]))) + AND ($2::text = '' OR p.scope_id = (SELECT id FROM transport_scopes WHERE name = $2)) + AND ($3::timestamptz IS NULL OR p.first_heard_at >= $3) + AND ($4::timestamptz IS NULL OR p.first_heard_at <= $4) + AND ($5::timestamptz IS NULL OR p.last_heard_at < $5) + AND ($7::text = '' OR p.parsed_payload->>'type' = $7) + GROUP BY p.trace_tag + ORDER BY MAX(p.last_heard_at) DESC + LIMIT $6 +) SELECT - encode(p.trace_tag, 'hex') AS trace_tag, - MIN(p.first_heard_at)::timestamptz AS first_heard_at, - MAX(p.last_heard_at)::timestamptz AS last_heard_at, - COUNT(DISTINCT p.packet_hash) AS packet_count, - COUNT(DISTINCT po.iata) AS iata_count, - MAX(p.parsed_payload->>'type')::text AS trace_type, - best.parsed_payload AS best_payload -FROM packets p -LEFT JOIN packet_observations po ON po.packet_hash = p.packet_hash -LEFT JOIN LATERAL ( - SELECT parsed_payload - FROM packets p2 - WHERE p2.trace_tag = p.trace_tag - ORDER BY jsonb_array_length(p2.parsed_payload->'pathHashes') DESC - LIMIT 1 -) best ON true -WHERE p.trace_tag IS NOT NULL - AND (COALESCE(cardinality($1::bpchar[]), 0) = 0 OR po.iata = ANY($1::bpchar[])) - AND ($2::text = '' OR p.scope_id = (SELECT id FROM transport_scopes WHERE name = $2)) - AND ($3::timestamptz IS NULL OR p.first_heard_at >= $3) - AND ($4::timestamptz IS NULL OR p.first_heard_at <= $4) - AND ($5::timestamptz IS NULL OR p.last_heard_at < $5) - AND ($7::text = '' OR p.parsed_payload->>'type' = $7) -GROUP BY p.trace_tag, best.parsed_payload -ORDER BY MAX(p.last_heard_at) DESC -LIMIT $6; + encode(t.trace_tag, 'hex') AS trace_tag, + t.first_heard_at::timestamptz AS first_heard_at, + t.last_heard_at::timestamptz AS last_heard_at, + t.packet_count, + (SELECT COUNT(*) + FROM trace_iatas ti + WHERE ti.trace_tag = t.trace_tag + AND (COALESCE(cardinality($1::bpchar[]), 0) = 0 OR ti.iata = ANY($1::bpchar[]))) AS iata_count, + t.trace_type::text AS trace_type, + (SELECT p3.parsed_payload + FROM packets p3 + WHERE p3.trace_tag = t.trace_tag + ORDER BY jsonb_array_length(p3.parsed_payload->'pathHashes') DESC + LIMIT 1) AS best_payload +FROM tags t +ORDER BY t.last_heard_at DESC; -- ============================================================ -- ROUTES diff --git a/db/sqlc/mock/querier.go b/db/sqlc/mock/querier.go index 7ccdc26..22f0983 100644 --- a/db/sqlc/mock/querier.go +++ b/db/sqlc/mock/querier.go @@ -85,6 +85,20 @@ func (mr *MockQuerierMockRecorder) DeleteOldTelemetry(ctx, reportedAt any) *gomo return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "DeleteOldTelemetry", reflect.TypeOf((*MockQuerier)(nil).DeleteOldTelemetry), ctx, reportedAt) } +// DeleteOldTraceIATAs mocks base method. +func (m *MockQuerier) DeleteOldTraceIATAs(ctx context.Context, lastHeard pgtype.Timestamptz) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "DeleteOldTraceIATAs", ctx, lastHeard) + ret0, _ := ret[0].(error) + return ret0 +} + +// DeleteOldTraceIATAs indicates an expected call of DeleteOldTraceIATAs. +func (mr *MockQuerierMockRecorder) DeleteOldTraceIATAs(ctx, lastHeard any) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "DeleteOldTraceIATAs", reflect.TypeOf((*MockQuerier)(nil).DeleteOldTraceIATAs), ctx, lastHeard) +} + // GetChannelByID mocks base method. func (m *MockQuerier) GetChannelByID(ctx context.Context, id int32) (db.Channel, error) { m.ctrl.T.Helper() @@ -1427,6 +1441,20 @@ func (mr *MockQuerierMockRecorder) UpsertRegionIATA(ctx, arg any) *gomock.Call { return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "UpsertRegionIATA", reflect.TypeOf((*MockQuerier)(nil).UpsertRegionIATA), ctx, arg) } +// UpsertTraceIATA mocks base method. +func (m *MockQuerier) UpsertTraceIATA(ctx context.Context, arg db.UpsertTraceIATAParams) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "UpsertTraceIATA", ctx, arg) + ret0, _ := ret[0].(error) + return ret0 +} + +// UpsertTraceIATA indicates an expected call of UpsertTraceIATA. +func (mr *MockQuerierMockRecorder) UpsertTraceIATA(ctx, arg any) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "UpsertTraceIATA", reflect.TypeOf((*MockQuerier)(nil).UpsertTraceIATA), ctx, arg) +} + // UpsertTransportScope mocks base method. func (m *MockQuerier) UpsertTransportScope(ctx context.Context, arg db.UpsertTransportScopeParams) error { m.ctrl.T.Helper() diff --git a/db/sqlc/models.go b/db/sqlc/models.go index e31b389..6e88f25 100644 --- a/db/sqlc/models.go +++ b/db/sqlc/models.go @@ -271,6 +271,12 @@ type RegionIata struct { AddedAt pgtype.Timestamptz `json:"added_at"` } +type TraceIata struct { + TraceTag []byte `json:"trace_tag"` + Iata string `json:"iata"` + LastHeard pgtype.Timestamptz `json:"last_heard"` +} + type TransportScope struct { ID int32 `json:"id"` Name string `json:"name"` diff --git a/db/sqlc/querier.go b/db/sqlc/querier.go index ef36ddf..3b158bb 100644 --- a/db/sqlc/querier.go +++ b/db/sqlc/querier.go @@ -19,6 +19,8 @@ type Querier interface { DeleteOldPackets(ctx context.Context, lastHeardAt pgtype.Timestamptz) 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. + DeleteOldTraceIATAs(ctx context.Context, lastHeard pgtype.Timestamptz) error GetChannelByID(ctx context.Context, id int32) (Channel, error) // Returns neighbors of a node that are in a different IATA. GetCrossIATANeighbors(ctx context.Context, arg GetCrossIATANeighborsParams) ([]GetCrossIATANeighborsRow, error) @@ -142,6 +144,8 @@ type Querier interface { // TRACES // ============================================================ // Returns distinct trace tags with summary info, ordered by most recent first. + // IATA membership comes from trace_iatas (joining observations here spilled the + // hash join). Per-tag details filled in only for the returned page. ListTraceTags(ctx context.Context, arg ListTraceTagsParams) ([]ListTraceTagsRow, error) // Delete node_neighbors where the neighbor has departed from node_short_ids // for that IATA, or where its prefix_4 is now ambiguous. @@ -232,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 + UpsertTraceIATA(ctx context.Context, arg UpsertTraceIATAParams) error // ============================================================ // TRANSPORT CODES // ============================================================ diff --git a/db/sqlc/queries.sql.go b/db/sqlc/queries.sql.go index 1fe3ac6..748207a 100644 --- a/db/sqlc/queries.sql.go +++ b/db/sqlc/queries.sql.go @@ -43,6 +43,16 @@ func (q *Queries) DeleteOldTelemetry(ctx context.Context, reportedAt pgtype.Time return err } +const deleteOldTraceIATAs = `-- name: DeleteOldTraceIATAs :exec +DELETE FROM trace_iatas WHERE last_heard < $1 +` + +// Keeps the trace IATA filter in step with packet retention. +func (q *Queries) DeleteOldTraceIATAs(ctx context.Context, lastHeard pgtype.Timestamptz) error { + _, err := q.db.Exec(ctx, deleteOldTraceIATAs, lastHeard) + return err +} + const getChannelByID = `-- name: GetChannelByID :one SELECT id, channel_hash, key_fingerprint, name, hashtag, is_hashtag, is_public, key_known, first_seen, last_seen, message_count FROM channels WHERE id = $1 ` @@ -2901,33 +2911,43 @@ func (q *Queries) ListRegions(ctx context.Context) ([]ListRegionsRow, error) { const listTraceTags = `-- name: ListTraceTags :many +WITH tags AS ( + SELECT + p.trace_tag, + MIN(p.first_heard_at) AS first_heard_at, + MAX(p.last_heard_at) AS last_heard_at, + COUNT(*) AS packet_count, + MAX(p.parsed_payload->>'type') AS trace_type + FROM packets p + WHERE p.trace_tag IS NOT NULL + AND (COALESCE(cardinality($1::bpchar[]), 0) = 0 OR p.trace_tag IN ( + SELECT ti.trace_tag FROM trace_iatas ti WHERE ti.iata = ANY($1::bpchar[]))) + AND ($2::text = '' OR p.scope_id = (SELECT id FROM transport_scopes WHERE name = $2)) + AND ($3::timestamptz IS NULL OR p.first_heard_at >= $3) + AND ($4::timestamptz IS NULL OR p.first_heard_at <= $4) + AND ($5::timestamptz IS NULL OR p.last_heard_at < $5) + AND ($7::text = '' OR p.parsed_payload->>'type' = $7) + GROUP BY p.trace_tag + ORDER BY MAX(p.last_heard_at) DESC + LIMIT $6 +) SELECT - encode(p.trace_tag, 'hex') AS trace_tag, - MIN(p.first_heard_at)::timestamptz AS first_heard_at, - MAX(p.last_heard_at)::timestamptz AS last_heard_at, - COUNT(DISTINCT p.packet_hash) AS packet_count, - COUNT(DISTINCT po.iata) AS iata_count, - MAX(p.parsed_payload->>'type')::text AS trace_type, - best.parsed_payload AS best_payload -FROM packets p -LEFT JOIN packet_observations po ON po.packet_hash = p.packet_hash -LEFT JOIN LATERAL ( - SELECT parsed_payload - FROM packets p2 - WHERE p2.trace_tag = p.trace_tag - ORDER BY jsonb_array_length(p2.parsed_payload->'pathHashes') DESC - LIMIT 1 -) best ON true -WHERE p.trace_tag IS NOT NULL - AND (COALESCE(cardinality($1::bpchar[]), 0) = 0 OR po.iata = ANY($1::bpchar[])) - AND ($2::text = '' OR p.scope_id = (SELECT id FROM transport_scopes WHERE name = $2)) - AND ($3::timestamptz IS NULL OR p.first_heard_at >= $3) - AND ($4::timestamptz IS NULL OR p.first_heard_at <= $4) - AND ($5::timestamptz IS NULL OR p.last_heard_at < $5) - AND ($7::text = '' OR p.parsed_payload->>'type' = $7) -GROUP BY p.trace_tag, best.parsed_payload -ORDER BY MAX(p.last_heard_at) DESC -LIMIT $6 + encode(t.trace_tag, 'hex') AS trace_tag, + t.first_heard_at::timestamptz AS first_heard_at, + t.last_heard_at::timestamptz AS last_heard_at, + t.packet_count, + (SELECT COUNT(*) + FROM trace_iatas ti + WHERE ti.trace_tag = t.trace_tag + AND (COALESCE(cardinality($1::bpchar[]), 0) = 0 OR ti.iata = ANY($1::bpchar[]))) AS iata_count, + t.trace_type::text AS trace_type, + (SELECT p3.parsed_payload + FROM packets p3 + WHERE p3.trace_tag = t.trace_tag + ORDER BY jsonb_array_length(p3.parsed_payload->'pathHashes') DESC + LIMIT 1) AS best_payload +FROM tags t +ORDER BY t.last_heard_at DESC ` type ListTraceTagsParams struct { @@ -2954,6 +2974,8 @@ type ListTraceTagsRow struct { // TRACES // ============================================================ // Returns distinct trace tags with summary info, ordered by most recent first. +// IATA membership comes from trace_iatas (joining observations here spilled the +// hash join). Per-tag details filled in only for the returned page. func (q *Queries) ListTraceTags(ctx context.Context, arg ListTraceTagsParams) ([]ListTraceTagsRow, error) { rows, err := q.db.Query(ctx, listTraceTags, arg.Column1, @@ -4046,6 +4068,24 @@ func (q *Queries) UpsertRegionIATA(ctx context.Context, arg UpsertRegionIATAPara return err } +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) +` + +type UpsertTraceIATAParams struct { + TraceTag []byte `json:"trace_tag"` + Iata string `json:"iata"` + LastHeard pgtype.Timestamptz `json:"last_heard"` +} + +func (q *Queries) UpsertTraceIATA(ctx context.Context, arg UpsertTraceIATAParams) error { + _, err := q.db.Exec(ctx, upsertTraceIATA, arg.TraceTag, arg.Iata, arg.LastHeard) + return err +} + const upsertTransportScope = `-- name: UpsertTransportScope :exec INSERT INTO transport_scopes (name, display_name, transport_key, key_fingerprint) diff --git a/db/traces.go b/db/traces.go index 1b63c2f..2624564 100644 --- a/db/traces.go +++ b/db/traces.go @@ -20,6 +20,18 @@ type tracePayload struct { SNRValues []float32 `json:"snrValues"` } +func (s *Store) UpsertTraceIATA(ctx context.Context, traceTag []byte, iata string, heardAt time.Time) error { + return s.q.UpsertTraceIATA(ctx, sqlc.UpsertTraceIATAParams{ + TraceTag: traceTag, + Iata: iata, + LastHeard: pgtype.Timestamptz{Time: heardAt, Valid: true}, + }) +} + +func (s *Store) DeleteOldTraceIATAs(ctx context.Context, cutoff time.Time) error { + return s.q.DeleteOldTraceIATAs(ctx, pgtype.Timestamptz{Time: cutoff, Valid: true}) +} + func (s *Store) ListTraceTags(ctx context.Context, iatas []string, scope, traceType string, since, until time.Time, cursor time.Time, limit int32) ([]api.TraceTagSummary, error) { var sinceTS, untilTS, cursorTS pgtype.Timestamptz if !since.IsZero() { diff --git a/internal/background/tasks.go b/internal/background/tasks.go index 559a3dd..9855b80 100644 --- a/internal/background/tasks.go +++ b/internal/background/tasks.go @@ -47,6 +47,9 @@ func CleanupTask(store *db.Store, telemetryRetention, packetRetention, interval if err := store.DeleteOldChannelIATAs(ctx, time.Now().Add(-packetRetention)); err != nil { return err } + if err := store.DeleteOldTraceIATAs(ctx, time.Now().Add(-packetRetention)); err != nil { + return err + } return nil }, } diff --git a/internal/ingest/ingest.go b/internal/ingest/ingest.go index 48310aa..0fd5215 100644 --- a/internal/ingest/ingest.go +++ b/internal/ingest/ingest.go @@ -154,6 +154,9 @@ type DB interface { // UpsertChannelIATA upserts a channel_iatas row. UpsertChannelIATA(ctx context.Context, channelHash []byte, iata string, heardAt time.Time) error + // UpsertTraceIATA upserts a trace_iatas row. + UpsertTraceIATA(ctx context.Context, traceTag []byte, iata string, heardAt time.Time) error + // GetPacketObservationCount returns the number of rows for the packet observations GetPacketObservationCount(ctx context.Context, packetHash []byte) (int64, error) diff --git a/internal/ingest/ingest_test.go b/internal/ingest/ingest_test.go index 81943d5..fbde760 100644 --- a/internal/ingest/ingest_test.go +++ b/internal/ingest/ingest_test.go @@ -178,6 +178,7 @@ type stubDB struct { upsertChannelCalls int upsertChannelHashOnlyCalls int upsertChannelIATACalls int + upsertTraceIATACalls int observationInserted bool } @@ -264,6 +265,11 @@ func (s *stubDB) UpsertChannelIATA(_ context.Context, _ []byte, _ string, _ time return nil } +func (s *stubDB) UpsertTraceIATA(_ context.Context, _ []byte, _ string, _ time.Time) error { + s.upsertTraceIATACalls++ + return nil +} + func (s *stubDB) GetPacketObservationCount(_ context.Context, _ []byte) (int64, error) { return 0, nil } diff --git a/internal/ingest/packet.go b/internal/ingest/packet.go index 741bb93..9ac0ed8 100644 --- a/internal/ingest/packet.go +++ b/internal/ingest/packet.go @@ -790,6 +790,12 @@ func (w *Worker) handlePacket(ctx context.Context, iata, pubkeyHex string, raw [ } } + if traceTag != nil && inserted { + 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) + } + } + // packet.PathHashes() reads packet.Path as hash-sized chunks, which is only true for // ordinary flood/direct-routed packets. TRACE repurposes packet.Path to carry one SNR // byte per hop instead, so for TRACE we resolve against the trace payload's own diff --git a/internal/ingest/side_effects_test.go b/internal/ingest/side_effects_test.go index 077caef..25bdb2d 100644 --- a/internal/ingest/side_effects_test.go +++ b/internal/ingest/side_effects_test.go @@ -187,4 +187,31 @@ func TestHandlePacket_Advert_SkipsChannelIATA(t *testing.T) { if db.upsertChannelIATACalls != 0 { t.Errorf("expected UpsertChannelIATA NOT to be called for a non-channel packet, got %d calls", db.upsertChannelIATACalls) } + if db.upsertTraceIATACalls != 0 { + t.Errorf("expected UpsertTraceIATA NOT to be called for a non-trace packet, got %d calls", db.upsertTraceIATACalls) + } +} + +func buildTracePacket(t *testing.T) *meshcore.Packet { + t.Helper() + payload, err := (&meshcore.Trace{Tag: 0xdeadbeef, AuthCode: 1}).ToBytes() + if err != nil { + t.Fatalf("trace to bytes: %v", err) + } + return &meshcore.Packet{ + Header: meshcore.MakeHeader(meshcore.RouteTypeFlood, meshcore.PayloadTypeTrace, 0), + Payload: payload, + } +} + +func TestHandlePacket_Trace_UpsertsTraceIATA(t *testing.T) { + w, db := newTestWorker() + db.observationInserted = true + envelope := packetEnvelope(t, buildTracePacket(t)) + + w.handlePacket(context.Background(), "YOW", "0102", envelope) + + if db.upsertTraceIATACalls != 1 { + t.Errorf("expected UpsertTraceIATA to be called once for a stored trace, got %d", db.upsertTraceIATACalls) + } } From 67cfb3f8ed0baf51bb297741d822c3c1e8d8e82d Mon Sep 17 00:00:00 2001 From: MrAlders0n Date: Thu, 23 Jul 2026 13:10:28 -0400 Subject: [PATCH 4/8] 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) } From b2557d8eb71d0a7f9cccc901b4d629fbd071c21e Mon Sep 17 00:00:00 2001 From: MrAlders0n Date: Thu, 23 Jul 2026 12:15:41 -0400 Subject: [PATCH 5/8] perf(scopes): stop cross-joining packets in scope stats GetScopeStats left-joined packets, observer_scopes and nodes at once, multiplying rows into the millions before COUNT(DISTINCT) deduped them (~10s per call). Count each table on its own and index packets(scope_id). --- db/migrations/013_packets_scope_index.sql | 3 +++ db/queries/queries.sql | 12 +++++------- db/sqlc/querier.go | 2 ++ db/sqlc/queries.sql.go | 12 +++++------- 4 files changed, 15 insertions(+), 14 deletions(-) create mode 100644 db/migrations/013_packets_scope_index.sql diff --git a/db/migrations/013_packets_scope_index.sql b/db/migrations/013_packets_scope_index.sql new file mode 100644 index 0000000..73029cf --- /dev/null +++ b/db/migrations/013_packets_scope_index.sql @@ -0,0 +1,3 @@ +-- Scope stats count off scope_id, which had no index. + +CREATE INDEX idx_packets_scope ON packets(scope_id) WHERE scope_id IS NOT NULL; diff --git a/db/queries/queries.sql b/db/queries/queries.sql index 4f1c265..bb8008b 100644 --- a/db/queries/queries.sql +++ b/db/queries/queries.sql @@ -917,16 +917,14 @@ WHERE ($1::text = '' OR preset = $1::text) ORDER BY preset, iata, source_type; -- name: GetScopeStats :many +-- Count each table on its own; the old cross-join blew up to millions of rows +-- before COUNT(DISTINCT) (~10s). SELECT ts.name, - COUNT(DISTINCT p.packet_hash) AS packet_count, - COUNT(DISTINCT os.observer_id) AS observer_count, - COUNT(DISTINCT n.id) AS node_count + (SELECT COUNT(*) FROM packets p WHERE p.scope_id = ts.id) AS packet_count, + (SELECT COUNT(*) FROM observer_scopes os WHERE os.scope_id = ts.id) AS observer_count, + (SELECT COUNT(*) FROM nodes n WHERE n.default_scope_id = ts.id) AS node_count FROM transport_scopes ts -LEFT JOIN packets p ON p.scope_id = ts.id -LEFT JOIN observer_scopes os ON os.scope_id = ts.id -LEFT JOIN nodes n ON n.default_scope_id = ts.id -GROUP BY ts.name ORDER BY ts.name; -- ============================================================ diff --git a/db/sqlc/querier.go b/db/sqlc/querier.go index fa93a88..efb1ba2 100644 --- a/db/sqlc/querier.go +++ b/db/sqlc/querier.go @@ -50,6 +50,8 @@ type Querier interface { GetRegionIATAs(ctx context.Context, regionID int32) ([]string, error) GetScopeByName(ctx context.Context, name string) (GetScopeByNameRow, error) GetScopeNames(ctx context.Context) ([]string, error) + // Count each table on its own; the old cross-join blew up to millions of rows + // before COUNT(DISTINCT) (~10s). GetScopeStats(ctx context.Context) ([]GetScopeStatsRow, error) GetScopesByIATAs(ctx context.Context, dollar_1 []string) ([]GetScopesByIATAsRow, error) // Returns node counts grouped by type, optionally filtered by IATA. diff --git a/db/sqlc/queries.sql.go b/db/sqlc/queries.sql.go index 1698cbe..787f5cd 100644 --- a/db/sqlc/queries.sql.go +++ b/db/sqlc/queries.sql.go @@ -1035,14 +1035,10 @@ func (q *Queries) GetScopeNames(ctx context.Context) ([]string, error) { const getScopeStats = `-- name: GetScopeStats :many SELECT ts.name, - COUNT(DISTINCT p.packet_hash) AS packet_count, - COUNT(DISTINCT os.observer_id) AS observer_count, - COUNT(DISTINCT n.id) AS node_count + (SELECT COUNT(*) FROM packets p WHERE p.scope_id = ts.id) AS packet_count, + (SELECT COUNT(*) FROM observer_scopes os WHERE os.scope_id = ts.id) AS observer_count, + (SELECT COUNT(*) FROM nodes n WHERE n.default_scope_id = ts.id) AS node_count FROM transport_scopes ts -LEFT JOIN packets p ON p.scope_id = ts.id -LEFT JOIN observer_scopes os ON os.scope_id = ts.id -LEFT JOIN nodes n ON n.default_scope_id = ts.id -GROUP BY ts.name ORDER BY ts.name ` @@ -1053,6 +1049,8 @@ type GetScopeStatsRow struct { NodeCount int64 `json:"node_count"` } +// Count each table on its own; the old cross-join blew up to millions of rows +// before COUNT(DISTINCT) (~10s). func (q *Queries) GetScopeStats(ctx context.Context) ([]GetScopeStatsRow, error) { rows, err := q.db.Query(ctx, getScopeStats) if err != nil { From b870ba1e897b3eb755a2d849e9a892b75d78faec Mon Sep 17 00:00:00 2001 From: MrAlders0n Date: Thu, 23 Jul 2026 12:13:20 -0400 Subject: [PATCH 6/8] perf(routes): index known_routes(last_seen) The routes list orders by last_seen with no index, so every page seq-scanned and sorted ~150k rows. Same fix as the packets list got in 008. --- db/migrations/014_known_routes_last_seen_index.sql | 3 +++ 1 file changed, 3 insertions(+) create mode 100644 db/migrations/014_known_routes_last_seen_index.sql diff --git a/db/migrations/014_known_routes_last_seen_index.sql b/db/migrations/014_known_routes_last_seen_index.sql new file mode 100644 index 0000000..faa5808 --- /dev/null +++ b/db/migrations/014_known_routes_last_seen_index.sql @@ -0,0 +1,3 @@ +-- Routes list orders by last_seen; index it (was seq-scanning ~150k rows per page). + +CREATE INDEX idx_known_routes_last_seen ON known_routes(last_seen DESC); From efb845b6769a0ae41c9d44aa2f87656e882f220c Mon Sep 17 00:00:00 2001 From: MrAlders0n Date: Thu, 23 Jul 2026 22:13:29 -0400 Subject: [PATCH 7/8] track payload_type on observations and serve the breakdown from a matview --- .../015_observations_payload_type.sql | 17 +++++++ db/migrations/016_mv_payload_breakdown.sql | 15 ++++++ db/packets.go | 1 + db/queries/queries.sql | 22 ++++---- db/sqlc/mock/querier.go | 22 ++++++-- db/sqlc/models.go | 7 +++ db/sqlc/querier.go | 5 +- db/sqlc/queries.sql.go | 50 +++++++++++-------- db/stats.go | 22 ++++---- internal/background/tasks.go | 3 ++ internal/ingest/packet.go | 2 + 11 files changed, 119 insertions(+), 47 deletions(-) create mode 100644 db/migrations/015_observations_payload_type.sql create mode 100644 db/migrations/016_mv_payload_breakdown.sql diff --git a/db/migrations/015_observations_payload_type.sql b/db/migrations/015_observations_payload_type.sql new file mode 100644 index 0000000..e3bb79b --- /dev/null +++ b/db/migrations/015_observations_payload_type.sql @@ -0,0 +1,17 @@ +-- Copy payload_type onto the observation (like source_broker) so the breakdown +-- drops its join back to packets, which was overrunning /dev/shm and 500ing. + +ALTER TABLE packet_observations ADD COLUMN payload_type SMALLINT; + +-- Backfill the last 7 days only (all the breakdown reads); parallelism off so +-- the join spills to disk, not /dev/shm. +SET max_parallel_workers_per_gather = 0; + +UPDATE packet_observations po +SET payload_type = p.payload_type +FROM packets p +WHERE p.packet_hash = po.packet_hash + AND po.heard_at > NOW() - INTERVAL '7 days' + AND po.payload_type IS NULL; + +RESET max_parallel_workers_per_gather; diff --git a/db/migrations/016_mv_payload_breakdown.sql b/db/migrations/016_mv_payload_breakdown.sql new file mode 100644 index 0000000..0710015 --- /dev/null +++ b/db/migrations/016_mv_payload_breakdown.sql @@ -0,0 +1,15 @@ +-- Precomputed payload breakdown so the request reads a small table instead of +-- scanning a week of observations. + +CREATE MATERIALIZED VIEW mv_payload_breakdown_by_iata AS +SELECT + iata, + payload_type, + COUNT(*) AS count +FROM packet_observations +WHERE heard_at > NOW() - INTERVAL '7 days' + AND payload_type IS NOT NULL +GROUP BY iata, payload_type; + +CREATE UNIQUE INDEX idx_mv_payload_breakdown + ON mv_payload_breakdown_by_iata(iata, payload_type); diff --git a/db/packets.go b/db/packets.go index 66a0be6..68bd34c 100644 --- a/db/packets.go +++ b/db/packets.go @@ -505,6 +505,7 @@ func (s *Store) InsertObservation(ctx context.Context, o ingest.InsertObservatio BandwidthKhz: &o.BandwidthKHz, CodingRate: &o.CodingRate, SourceBroker: &o.SourceBroker, + PayloadType: &o.PayloadType, } row, err := s.q.InsertObservation(ctx, params) if errors.Is(err, pgx.ErrNoRows) { diff --git a/db/queries/queries.sql b/db/queries/queries.sql index bb8008b..b203e4c 100644 --- a/db/queries/queries.sql +++ b/db/queries/queries.sql @@ -523,9 +523,10 @@ INSERT INTO packet_observations ( spread_factor, bandwidth_khz, coding_rate, - source_broker + source_broker, + payload_type ) VALUES ( - $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16 + $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17 ) ON CONFLICT (packet_hash, observer_id) DO NOTHING RETURNING *; @@ -821,15 +822,13 @@ ORDER BY observation_count DESC LIMIT $2; -- name: GetStatsPayloadBreakdown :many --- Returns observation counts grouped by payload type for the given window and IATA. +-- Payload-type counts for the IATA, from the precomputed view. SELECT - p.payload_type, - COUNT(*) AS count -FROM packet_observations po -JOIN packets p ON p.packet_hash = po.packet_hash -WHERE po.heard_at > $1 - AND (COALESCE(cardinality($2::bpchar[]), 0) = 0 OR po.iata = ANY($2::bpchar[])) -GROUP BY p.payload_type + payload_type, + SUM(count)::bigint AS count +FROM mv_payload_breakdown_by_iata +WHERE (COALESCE(cardinality($1::bpchar[]), 0) = 0 OR iata = ANY($1::bpchar[])) +GROUP BY payload_type ORDER BY count DESC; -- name: GetStatsNodeTypes :many @@ -1141,6 +1140,9 @@ REFRESH MATERIALIZED VIEW CONCURRENTLY mv_hourly_iata_stats; -- name: RefreshTopNodes :exec REFRESH MATERIALIZED VIEW CONCURRENTLY mv_top_nodes_by_iata; +-- name: RefreshPayloadBreakdown :exec +REFRESH MATERIALIZED VIEW CONCURRENTLY mv_payload_breakdown_by_iata; + -- name: RefreshRadioPresets :exec REFRESH MATERIALIZED VIEW CONCURRENTLY mv_radio_presets; diff --git a/db/sqlc/mock/querier.go b/db/sqlc/mock/querier.go index 22f0983..aa8d1da 100644 --- a/db/sqlc/mock/querier.go +++ b/db/sqlc/mock/querier.go @@ -550,18 +550,18 @@ func (mr *MockQuerierMockRecorder) GetStatsOverview(ctx, dollar_1 any) *gomock.C } // GetStatsPayloadBreakdown mocks base method. -func (m *MockQuerier) GetStatsPayloadBreakdown(ctx context.Context, arg db.GetStatsPayloadBreakdownParams) ([]db.GetStatsPayloadBreakdownRow, error) { +func (m *MockQuerier) GetStatsPayloadBreakdown(ctx context.Context, dollar_1 []string) ([]db.GetStatsPayloadBreakdownRow, error) { m.ctrl.T.Helper() - ret := m.ctrl.Call(m, "GetStatsPayloadBreakdown", ctx, arg) + ret := m.ctrl.Call(m, "GetStatsPayloadBreakdown", ctx, dollar_1) ret0, _ := ret[0].([]db.GetStatsPayloadBreakdownRow) ret1, _ := ret[1].(error) return ret0, ret1 } // GetStatsPayloadBreakdown indicates an expected call of GetStatsPayloadBreakdown. -func (mr *MockQuerierMockRecorder) GetStatsPayloadBreakdown(ctx, arg any) *gomock.Call { +func (mr *MockQuerierMockRecorder) GetStatsPayloadBreakdown(ctx, dollar_1 any) *gomock.Call { mr.mock.ctrl.T.Helper() - return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "GetStatsPayloadBreakdown", reflect.TypeOf((*MockQuerier)(nil).GetStatsPayloadBreakdown), ctx, arg) + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "GetStatsPayloadBreakdown", reflect.TypeOf((*MockQuerier)(nil).GetStatsPayloadBreakdown), ctx, dollar_1) } // GetStatsTopAdvertisers mocks base method. @@ -995,6 +995,20 @@ func (mr *MockQuerierMockRecorder) RefreshHourlyStats(ctx any) *gomock.Call { return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "RefreshHourlyStats", reflect.TypeOf((*MockQuerier)(nil).RefreshHourlyStats), ctx) } +// RefreshPayloadBreakdown mocks base method. +func (m *MockQuerier) RefreshPayloadBreakdown(ctx context.Context) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "RefreshPayloadBreakdown", ctx) + ret0, _ := ret[0].(error) + return ret0 +} + +// RefreshPayloadBreakdown indicates an expected call of RefreshPayloadBreakdown. +func (mr *MockQuerierMockRecorder) RefreshPayloadBreakdown(ctx any) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "RefreshPayloadBreakdown", reflect.TypeOf((*MockQuerier)(nil).RefreshPayloadBreakdown), ctx) +} + // RefreshRadioPresets mocks base method. func (m *MockQuerier) RefreshRadioPresets(ctx context.Context) error { m.ctrl.T.Helper() diff --git a/db/sqlc/models.go b/db/sqlc/models.go index 6e88f25..388fc85 100644 --- a/db/sqlc/models.go +++ b/db/sqlc/models.go @@ -74,6 +74,12 @@ type MvHourlyIataStat struct { ActiveObservers int64 `json:"active_observers"` } +type MvPayloadBreakdownByIatum struct { + Iata string `json:"iata"` + PayloadType *int16 `json:"payload_type"` + Count int64 `json:"count"` +} + type MvRadioPreset struct { Preset string `json:"preset"` Iata string `json:"iata"` @@ -250,6 +256,7 @@ type PacketObservation struct { BandwidthKhz *float32 `json:"bandwidth_khz"` CodingRate *int16 `json:"coding_rate"` SourceBroker *string `json:"source_broker"` + PayloadType *int16 `json:"payload_type"` } type Region struct { diff --git a/db/sqlc/querier.go b/db/sqlc/querier.go index efb1ba2..bb1712a 100644 --- a/db/sqlc/querier.go +++ b/db/sqlc/querier.go @@ -60,8 +60,8 @@ type Querier interface { // STATS // ============================================================ GetStatsOverview(ctx context.Context, dollar_1 []string) (GetStatsOverviewRow, error) - // Returns observation counts grouped by payload type for the given window and IATA. - GetStatsPayloadBreakdown(ctx context.Context, arg GetStatsPayloadBreakdownParams) ([]GetStatsPayloadBreakdownRow, error) + // Payload-type counts for the IATA, from the precomputed view. + GetStatsPayloadBreakdown(ctx context.Context, dollar_1 []string) ([]GetStatsPayloadBreakdownRow, error) // Returns the top N nodes by distinct ADVERT packet count for the given window and IATA. // COUNT(DISTINCT p.packet_hash) rather than COUNT(*): the same advert broadcast is commonly // heard by more than one observer, and each hearing is its own packet_observations row -- @@ -156,6 +156,7 @@ type Querier interface { // that IATA, or where any hop's prefix_4 is now ambiguous (matches >1 node). ReconfirmRoutes(ctx context.Context) error RefreshHourlyStats(ctx context.Context) error + RefreshPayloadBreakdown(ctx context.Context) error RefreshRadioPresets(ctx context.Context) error RefreshTopNodes(ctx context.Context) error // ============================================================ diff --git a/db/sqlc/queries.sql.go b/db/sqlc/queries.sql.go index 787f5cd..6871173 100644 --- a/db/sqlc/queries.sql.go +++ b/db/sqlc/queries.sql.go @@ -1197,29 +1197,22 @@ func (q *Queries) GetStatsOverview(ctx context.Context, dollar_1 []string) (GetS const getStatsPayloadBreakdown = `-- name: GetStatsPayloadBreakdown :many SELECT - p.payload_type, - COUNT(*) AS count -FROM packet_observations po -JOIN packets p ON p.packet_hash = po.packet_hash -WHERE po.heard_at > $1 - AND (COALESCE(cardinality($2::bpchar[]), 0) = 0 OR po.iata = ANY($2::bpchar[])) -GROUP BY p.payload_type + payload_type, + SUM(count)::bigint AS count +FROM mv_payload_breakdown_by_iata +WHERE (COALESCE(cardinality($1::bpchar[]), 0) = 0 OR iata = ANY($1::bpchar[])) +GROUP BY payload_type ORDER BY count DESC ` -type GetStatsPayloadBreakdownParams struct { - HeardAt pgtype.Timestamptz `json:"heard_at"` - Column2 []string `json:"column_2"` -} - type GetStatsPayloadBreakdownRow struct { - PayloadType int16 `json:"payload_type"` - Count int64 `json:"count"` + PayloadType *int16 `json:"payload_type"` + Count int64 `json:"count"` } -// Returns observation counts grouped by payload type for the given window and IATA. -func (q *Queries) GetStatsPayloadBreakdown(ctx context.Context, arg GetStatsPayloadBreakdownParams) ([]GetStatsPayloadBreakdownRow, error) { - rows, err := q.db.Query(ctx, getStatsPayloadBreakdown, arg.HeardAt, arg.Column2) +// Payload-type counts for the IATA, from the precomputed view. +func (q *Queries) GetStatsPayloadBreakdown(ctx context.Context, dollar_1 []string) ([]GetStatsPayloadBreakdownRow, error) { + rows, err := q.db.Query(ctx, getStatsPayloadBreakdown, dollar_1) if err != nil { return nil, err } @@ -1551,12 +1544,13 @@ INSERT INTO packet_observations ( spread_factor, bandwidth_khz, coding_rate, - source_broker + source_broker, + payload_type ) VALUES ( - $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16 + $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17 ) ON CONFLICT (packet_hash, observer_id) DO NOTHING -RETURNING id, packet_hash, observer_id, iata, heard_at, path_length_byte, hash_size, hop_count, path_bytes, rssi, snr, propagation_time_ms, radio_freq_mhz, spread_factor, bandwidth_khz, coding_rate, source_broker +RETURNING id, packet_hash, observer_id, iata, heard_at, path_length_byte, hash_size, hop_count, path_bytes, rssi, snr, propagation_time_ms, radio_freq_mhz, spread_factor, bandwidth_khz, coding_rate, source_broker, payload_type ` type InsertObservationParams struct { @@ -1576,6 +1570,7 @@ type InsertObservationParams struct { BandwidthKhz *float32 `json:"bandwidth_khz"` CodingRate *int16 `json:"coding_rate"` SourceBroker *string `json:"source_broker"` + PayloadType *int16 `json:"payload_type"` } // ============================================================ @@ -1599,6 +1594,7 @@ func (q *Queries) InsertObservation(ctx context.Context, arg InsertObservationPa arg.BandwidthKhz, arg.CodingRate, arg.SourceBroker, + arg.PayloadType, ) var i PacketObservation err := row.Scan( @@ -1619,6 +1615,7 @@ func (q *Queries) InsertObservation(ctx context.Context, arg InsertObservationPa &i.BandwidthKhz, &i.CodingRate, &i.SourceBroker, + &i.PayloadType, ) return i, err } @@ -2293,7 +2290,7 @@ func (q *Queries) ListNodes(ctx context.Context, arg ListNodesParams) ([]ListNod } const listObservationsForPacket = `-- name: ListObservationsForPacket :many -SELECT po.id, po.packet_hash, po.observer_id, po.iata, po.heard_at, po.path_length_byte, po.hash_size, po.hop_count, po.path_bytes, po.rssi, po.snr, po.propagation_time_ms, po.radio_freq_mhz, po.spread_factor, po.bandwidth_khz, po.coding_rate, po.source_broker, o.display_name AS observer_name +SELECT po.id, po.packet_hash, po.observer_id, po.iata, po.heard_at, po.path_length_byte, po.hash_size, po.hop_count, po.path_bytes, po.rssi, po.snr, po.propagation_time_ms, po.radio_freq_mhz, po.spread_factor, po.bandwidth_khz, po.coding_rate, po.source_broker, po.payload_type, o.display_name AS observer_name FROM packet_observations po LEFT JOIN observers o ON o.id = po.observer_id WHERE po.packet_hash = $1 @@ -2318,6 +2315,7 @@ type ListObservationsForPacketRow struct { BandwidthKhz *float32 `json:"bandwidth_khz"` CodingRate *int16 `json:"coding_rate"` SourceBroker *string `json:"source_broker"` + PayloadType *int16 `json:"payload_type"` ObserverName *string `json:"observer_name"` } @@ -2348,6 +2346,7 @@ func (q *Queries) ListObservationsForPacket(ctx context.Context, packetHash []by &i.BandwidthKhz, &i.CodingRate, &i.SourceBroker, + &i.PayloadType, &i.ObserverName, ); err != nil { return nil, err @@ -3073,6 +3072,15 @@ func (q *Queries) RefreshHourlyStats(ctx context.Context) error { return err } +const refreshPayloadBreakdown = `-- name: RefreshPayloadBreakdown :exec +REFRESH MATERIALIZED VIEW CONCURRENTLY mv_payload_breakdown_by_iata +` + +func (q *Queries) RefreshPayloadBreakdown(ctx context.Context) error { + _, err := q.db.Exec(ctx, refreshPayloadBreakdown) + return err +} + const refreshRadioPresets = `-- name: RefreshRadioPresets :exec REFRESH MATERIALIZED VIEW CONCURRENTLY mv_radio_presets ` diff --git a/db/stats.go b/db/stats.go index 930601b..a152ad6 100644 --- a/db/stats.go +++ b/db/stats.go @@ -51,22 +51,20 @@ func (s *Store) GetStatsObservations(ctx context.Context, iatas []string, since return points, nil } -func (s *Store) GetStatsPayloadBreakdown(ctx context.Context, iatas []string, since time.Time) ([]api.PayloadBreakdownItem, error) { - if since.IsZero() { - since = time.Now().Add(-24 * time.Hour) - } - rows, err := s.q.GetStatsPayloadBreakdown(ctx, sqlc.GetStatsPayloadBreakdownParams{ - HeardAt: pgtype.Timestamptz{Time: since, Valid: true}, - Column2: iatas, - }) +// since is unused: the view is a rolling 7-day window refreshed on a timer. +func (s *Store) GetStatsPayloadBreakdown(ctx context.Context, iatas []string, _ time.Time) ([]api.PayloadBreakdownItem, error) { + rows, err := s.q.GetStatsPayloadBreakdown(ctx, iatas) if err != nil { return nil, err } items := make([]api.PayloadBreakdownItem, 0, len(rows)) for _, v := range rows { + if v.PayloadType == nil { + continue + } items = append(items, api.PayloadBreakdownItem{ - PayloadType: v.PayloadType, - PayloadTypeName: api.PayloadTypeName(v.PayloadType), + PayloadType: *v.PayloadType, + PayloadTypeName: api.PayloadTypeName(*v.PayloadType), Count: v.Count, }) } @@ -244,6 +242,10 @@ func (s *Store) RefreshTopNodes(ctx context.Context) error { return s.q.RefreshTopNodes(ctx) } +func (s *Store) RefreshPayloadBreakdown(ctx context.Context) error { + return s.q.RefreshPayloadBreakdown(ctx) +} + func (s *Store) RefreshRadioPresets(ctx context.Context) error { return s.q.RefreshRadioPresets(ctx) } diff --git a/internal/background/tasks.go b/internal/background/tasks.go index 8824304..1966c6f 100644 --- a/internal/background/tasks.go +++ b/internal/background/tasks.go @@ -24,6 +24,9 @@ func ViewRefreshTask(store *db.Store, interval time.Duration) Task { if err := store.RefreshTopNodes(ctx); err != nil { log.Printf("background[view_refresh]: top nodes: %v", err) } + if err := store.RefreshPayloadBreakdown(ctx); err != nil { + log.Printf("background[view_refresh]: payload breakdown: %v", err) + } if err := store.RefreshRadioPresets(ctx); err != nil { log.Printf("background[view_refresh]: radio presets: %v", err) } diff --git a/internal/ingest/packet.go b/internal/ingest/packet.go index b51c99e..eb90db1 100644 --- a/internal/ingest/packet.go +++ b/internal/ingest/packet.go @@ -54,6 +54,7 @@ type InsertObservationParams struct { BandwidthKHz float32 CodingRate int16 SourceBroker string + PayloadType int16 } // RadioSettings holds the radio configuration for an observer, populated from @@ -770,6 +771,7 @@ func (w *Worker) handlePacket(ctx context.Context, iata, pubkeyHex string, raw [ BandwidthKHz: radio.BWKHz, CodingRate: radio.CR, SourceBroker: w.cfg.BrokerName, + PayloadType: int16(packet.PayloadType()), } inserted, err := w.db.InsertObservation(ctx, oParams) if err != nil { From db40e6b3e214e4b5ca11944177121110cfa65355 Mon Sep 17 00:00:00 2001 From: MrAlders0n Date: Thu, 23 Jul 2026 22:14:40 -0400 Subject: [PATCH 8/8] serve top observers from a matview instead of scanning per request --- db/migrations/017_mv_top_observers.sql | 18 +++++++++ db/queries/queries.sql | 29 +++++++------- db/sqlc/mock/querier.go | 14 +++++++ db/sqlc/models.go | 9 +++++ db/sqlc/querier.go | 4 +- db/sqlc/queries.sql.go | 53 ++++++++++++++------------ db/stats.go | 13 ++++--- internal/background/tasks.go | 3 ++ 8 files changed, 98 insertions(+), 45 deletions(-) create mode 100644 db/migrations/017_mv_top_observers.sql diff --git a/db/migrations/017_mv_top_observers.sql b/db/migrations/017_mv_top_observers.sql new file mode 100644 index 0000000..c1e1c6b --- /dev/null +++ b/db/migrations/017_mv_top_observers.sql @@ -0,0 +1,18 @@ +-- Precomputed top observers, same shape as mv_top_nodes_by_iata (was a ~2s +-- per-request scan of a week of observations). + +CREATE MATERIALIZED VIEW mv_top_observers_by_iata AS +SELECT + po.iata, + po.observer_id, + o.display_name, + o.observer_type, + COUNT(*) AS observation_count, + MAX(po.heard_at) AS last_heard +FROM packet_observations po +JOIN observers o ON o.id = po.observer_id +WHERE po.heard_at > NOW() - INTERVAL '7 days' +GROUP BY po.iata, po.observer_id, o.display_name, o.observer_type; + +CREATE UNIQUE INDEX idx_mv_top_observers + ON mv_top_observers_by_iata(iata, observer_id); diff --git a/db/queries/queries.sql b/db/queries/queries.sql index b203e4c..baa29e5 100644 --- a/db/queries/queries.sql +++ b/db/queries/queries.sql @@ -843,22 +843,18 @@ GROUP BY n.node_type ORDER BY count DESC; -- name: GetStatsTopObservers :many --- Returns the top N observers by observation count for the given window and IATA. +-- Top N observers for the IATA, from the precomputed view. Counts sum across +-- matched IATAs; iata is a representative one for display. SELECT - o.id, - o.display_name, - o.observer_type, - COUNT(*) AS observation_count, - COALESCE(( - SELECT po2.iata FROM packet_observations po2 - WHERE po2.observer_id = o.id - ORDER BY po2.heard_at DESC LIMIT 1 - ), '') AS iata -FROM packet_observations po -JOIN observers o ON o.id = po.observer_id -WHERE po.heard_at > $1 - AND (COALESCE(cardinality($2::bpchar[]), 0) = 0 OR po.iata = ANY($2::bpchar[])) -GROUP BY o.id + observer_id AS id, + display_name, + observer_type, + COALESCE(SUM(observation_count), 0)::bigint AS observation_count, + COALESCE(MAX(iata), '')::bpchar AS iata +FROM mv_top_observers_by_iata +WHERE last_heard > $1 + AND (COALESCE(cardinality($2::bpchar[]), 0) = 0 OR iata = ANY($2::bpchar[])) +GROUP BY observer_id, display_name, observer_type ORDER BY observation_count DESC LIMIT $3; @@ -1140,6 +1136,9 @@ REFRESH MATERIALIZED VIEW CONCURRENTLY mv_hourly_iata_stats; -- name: RefreshTopNodes :exec REFRESH MATERIALIZED VIEW CONCURRENTLY mv_top_nodes_by_iata; +-- name: RefreshTopObservers :exec +REFRESH MATERIALIZED VIEW CONCURRENTLY mv_top_observers_by_iata; + -- name: RefreshPayloadBreakdown :exec REFRESH MATERIALIZED VIEW CONCURRENTLY mv_payload_breakdown_by_iata; diff --git a/db/sqlc/mock/querier.go b/db/sqlc/mock/querier.go index aa8d1da..0e217ba 100644 --- a/db/sqlc/mock/querier.go +++ b/db/sqlc/mock/querier.go @@ -1037,6 +1037,20 @@ func (mr *MockQuerierMockRecorder) RefreshTopNodes(ctx any) *gomock.Call { return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "RefreshTopNodes", reflect.TypeOf((*MockQuerier)(nil).RefreshTopNodes), ctx) } +// RefreshTopObservers mocks base method. +func (m *MockQuerier) RefreshTopObservers(ctx context.Context) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "RefreshTopObservers", ctx) + ret0, _ := ret[0].(error) + return ret0 +} + +// RefreshTopObservers indicates an expected call of RefreshTopObservers. +func (mr *MockQuerierMockRecorder) RefreshTopObservers(ctx any) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "RefreshTopObservers", reflect.TypeOf((*MockQuerier)(nil).RefreshTopObservers), ctx) +} + // ResolvePathHashesP1 mocks base method. func (m *MockQuerier) ResolvePathHashesP1(ctx context.Context, arg db.ResolvePathHashesP1Params) ([]db.ResolvePathHashesP1Row, error) { m.ctrl.T.Helper() diff --git a/db/sqlc/models.go b/db/sqlc/models.go index 388fc85..60c3030 100644 --- a/db/sqlc/models.go +++ b/db/sqlc/models.go @@ -96,6 +96,15 @@ type MvTopNodesByIatum struct { LastHeard pgtype.Timestamptz `json:"last_heard"` } +type MvTopObserversByIatum struct { + Iata string `json:"iata"` + ObserverID uuid.UUID `json:"observer_id"` + DisplayName *string `json:"display_name"` + ObserverType *string `json:"observer_type"` + ObservationCount int64 `json:"observation_count"` + LastHeard interface{} `json:"last_heard"` +} + type Node struct { ID uuid.UUID `json:"id"` PublicKey []byte `json:"public_key"` diff --git a/db/sqlc/querier.go b/db/sqlc/querier.go index bb1712a..fda8940 100644 --- a/db/sqlc/querier.go +++ b/db/sqlc/querier.go @@ -67,7 +67,8 @@ type Querier interface { // heard by more than one observer, and each hearing is its own packet_observations row -- // this counts adverts sent, not adverts heard. GetStatsTopAdvertisers(ctx context.Context, arg GetStatsTopAdvertisersParams) ([]GetStatsTopAdvertisersRow, error) - // Returns the top N observers by observation count for the given window and IATA. + // Top N observers for the IATA, from the precomputed view. Counts sum across + // matched IATAs; iata is a representative one for display. GetStatsTopObservers(ctx context.Context, arg GetStatsTopObserversParams) ([]GetStatsTopObserversRow, error) // Returns the top N companion names by decrypted channel message count for the given // window and IATA. Grouped by sender_name as decrypted from the message itself, not by @@ -159,6 +160,7 @@ type Querier interface { RefreshPayloadBreakdown(ctx context.Context) error RefreshRadioPresets(ctx context.Context) error RefreshTopNodes(ctx context.Context) error + RefreshTopObservers(ctx context.Context) error // ============================================================ // HELPERS // ============================================================ diff --git a/db/sqlc/queries.sql.go b/db/sqlc/queries.sql.go index 6871173..87fb880 100644 --- a/db/sqlc/queries.sql.go +++ b/db/sqlc/queries.sql.go @@ -1303,41 +1303,37 @@ func (q *Queries) GetStatsTopAdvertisers(ctx context.Context, arg GetStatsTopAdv const getStatsTopObservers = `-- name: GetStatsTopObservers :many SELECT - o.id, - o.display_name, - o.observer_type, - COUNT(*) AS observation_count, - COALESCE(( - SELECT po2.iata FROM packet_observations po2 - WHERE po2.observer_id = o.id - ORDER BY po2.heard_at DESC LIMIT 1 - ), '') AS iata -FROM packet_observations po -JOIN observers o ON o.id = po.observer_id -WHERE po.heard_at > $1 - AND (COALESCE(cardinality($2::bpchar[]), 0) = 0 OR po.iata = ANY($2::bpchar[])) -GROUP BY o.id + observer_id AS id, + display_name, + observer_type, + COALESCE(SUM(observation_count), 0)::bigint AS observation_count, + COALESCE(MAX(iata), '')::bpchar AS iata +FROM mv_top_observers_by_iata +WHERE last_heard > $1 + AND (COALESCE(cardinality($2::bpchar[]), 0) = 0 OR iata = ANY($2::bpchar[])) +GROUP BY observer_id, display_name, observer_type ORDER BY observation_count DESC LIMIT $3 ` type GetStatsTopObserversParams struct { - HeardAt pgtype.Timestamptz `json:"heard_at"` - Column2 []string `json:"column_2"` - Limit int32 `json:"limit"` + LastHeard interface{} `json:"last_heard"` + Column2 []string `json:"column_2"` + Limit int32 `json:"limit"` } type GetStatsTopObserversRow struct { - ID uuid.UUID `json:"id"` - DisplayName *string `json:"display_name"` - ObserverType *string `json:"observer_type"` - ObservationCount int64 `json:"observation_count"` - Iata interface{} `json:"iata"` + ID uuid.UUID `json:"id"` + DisplayName *string `json:"display_name"` + ObserverType *string `json:"observer_type"` + ObservationCount int64 `json:"observation_count"` + Iata string `json:"iata"` } -// Returns the top N observers by observation count for the given window and IATA. +// Top N observers for the IATA, from the precomputed view. Counts sum across +// matched IATAs; iata is a representative one for display. func (q *Queries) GetStatsTopObservers(ctx context.Context, arg GetStatsTopObserversParams) ([]GetStatsTopObserversRow, error) { - rows, err := q.db.Query(ctx, getStatsTopObservers, arg.HeardAt, arg.Column2, arg.Limit) + rows, err := q.db.Query(ctx, getStatsTopObservers, arg.LastHeard, arg.Column2, arg.Limit) if err != nil { return nil, err } @@ -3099,6 +3095,15 @@ func (q *Queries) RefreshTopNodes(ctx context.Context) error { return err } +const refreshTopObservers = `-- name: RefreshTopObservers :exec +REFRESH MATERIALIZED VIEW CONCURRENTLY mv_top_observers_by_iata +` + +func (q *Queries) RefreshTopObservers(ctx context.Context) error { + _, err := q.db.Exec(ctx, refreshTopObservers) + return err +} + const resolvePathHashesP1 = `-- name: ResolvePathHashesP1 :many diff --git a/db/stats.go b/db/stats.go index a152ad6..f2e5d1e 100644 --- a/db/stats.go +++ b/db/stats.go @@ -103,21 +103,20 @@ func (s *Store) GetStatsTopObservers(ctx context.Context, iatas []string, since since = time.Now().Add(-24 * time.Hour) } rows, err := s.q.GetStatsTopObservers(ctx, sqlc.GetStatsTopObserversParams{ - HeardAt: pgtype.Timestamptz{Time: since, Valid: true}, - Column2: iatas, - Limit: limit, + LastHeard: pgtype.Timestamptz{Time: since, Valid: true}, + Column2: iatas, + Limit: limit, }) if err != nil { return nil, err } items := make([]api.TopObserver, 0, len(rows)) for _, v := range rows { - iata, _ := v.Iata.(string) items = append(items, api.TopObserver{ ObserverID: v.ID, DisplayName: v.DisplayName, ObserverType: v.ObserverType, - IATA: iata, + IATA: v.Iata, ObservationCount: v.ObservationCount, }) } @@ -246,6 +245,10 @@ func (s *Store) RefreshPayloadBreakdown(ctx context.Context) error { return s.q.RefreshPayloadBreakdown(ctx) } +func (s *Store) RefreshTopObservers(ctx context.Context) error { + return s.q.RefreshTopObservers(ctx) +} + func (s *Store) RefreshRadioPresets(ctx context.Context) error { return s.q.RefreshRadioPresets(ctx) } diff --git a/internal/background/tasks.go b/internal/background/tasks.go index 1966c6f..b9a4d1a 100644 --- a/internal/background/tasks.go +++ b/internal/background/tasks.go @@ -24,6 +24,9 @@ func ViewRefreshTask(store *db.Store, interval time.Duration) Task { if err := store.RefreshTopNodes(ctx); err != nil { log.Printf("background[view_refresh]: top nodes: %v", err) } + if err := store.RefreshTopObservers(ctx); err != nil { + log.Printf("background[view_refresh]: top observers: %v", err) + } if err := store.RefreshPayloadBreakdown(ctx); err != nil { log.Printf("background[view_refresh]: payload breakdown: %v", err) }