Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 17 additions & 5 deletions db/channels.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
42 changes: 31 additions & 11 deletions db/channels_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
Expand Down Expand Up @@ -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)
}
Expand All @@ -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)
Expand Down
23 changes: 23 additions & 0 deletions db/migrations/011_channel_iatas.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
-- 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);

-- Seed from retained packets; parallelism off so the join spills to disk, not /dev/shm.
SET max_parallel_workers_per_gather = 0;

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;
23 changes: 23 additions & 0 deletions db/migrations/012_trace_iatas.sql
Original file line number Diff line number Diff line change
@@ -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);

-- 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;
3 changes: 3 additions & 0 deletions db/migrations/013_packets_scope_index.sql
Original file line number Diff line number Diff line change
@@ -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;
3 changes: 3 additions & 0 deletions db/migrations/014_known_routes_last_seen_index.sql
Original file line number Diff line number Diff line change
@@ -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);
121 changes: 76 additions & 45 deletions db/queries/queries.sql
Original file line number Diff line number Diff line change
Expand Up @@ -494,6 +494,14 @@ 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;

-- name: DeleteOldTraceIATAs :exec
-- Keeps the trace IATA filter in step with packet retention.
DELETE FROM trace_iatas WHERE last_heard < $1;

-- ============================================================
-- PACKET OBSERVATIONS
-- ============================================================
Expand Down Expand Up @@ -667,22 +675,35 @@ ON CONFLICT (channel_hash) WHERE key_fingerprint IS NULL DO UPDATE SET
last_seen = NOW()
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 = EXCLUDED.last_heard
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 = EXCLUDED.last_heard
WHERE EXCLUDED.last_heard > trace_iatas.last_heard + INTERVAL '1 hour';

-- 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;
Expand Down Expand Up @@ -896,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;

-- ============================================================
Expand Down Expand Up @@ -956,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
Expand Down
Loading
Loading