// Code generated by sqlc. DO NOT EDIT. // versions: // sqlc v1.31.1 // source: queries.sql package db import ( "context" "github.com/google/uuid" "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 deleteOldNodes = `-- name: DeleteOldNodes :exec DELETE FROM nodes WHERE last_seen < $1 AND id NOT IN (SELECT owner_node_id FROM observer_owners WHERE owner_node_id IS NOT NULL) ` // Deletes nodes not seen since the given cutoff. node_iatas and node_neighbors cascade- // delete via FK. Excludes nodes referenced by observer_owners.owner_node_id -- that FK has // no ON DELETE action, so deleting one directly would fail the whole statement anyway, and // an operator manually recorded ownership for that node, so leave it alone even if stale. // known_routes.node_ids is a plain UUID[] with no FK; a deleted node's id can be left // dangling in old routes there, but ReconfirmTask already prunes stale/ambiguous routes // periodically and will clean those up on its own schedule. func (q *Queries) DeleteOldNodes(ctx context.Context, lastSeen pgtype.Timestamptz) error { _, err := q.db.Exec(ctx, deleteOldNodes, lastSeen) return err } const deleteOldPackets = `-- name: DeleteOldPackets :exec DELETE FROM packets WHERE last_heard_at < $1 ` // Deletes packets and their observations older than the given cutoff. // packet_observations cascade-delete via FK. func (q *Queries) DeleteOldPackets(ctx context.Context, lastHeardAt pgtype.Timestamptz) error { _, err := q.db.Exec(ctx, deleteOldPackets, lastHeardAt) return err } const deleteOldRoutes = `-- name: DeleteOldRoutes :exec DELETE FROM known_routes WHERE last_seen < $1 OR (observation_count < $2 AND last_seen < $3) ` type DeleteOldRoutesParams struct { LastSeen pgtype.Timestamptz `json:"last_seen"` ObservationCount int64 `json:"observation_count"` LastSeen_2 pgtype.Timestamptz `json:"last_seen_2"` } // Deletes routes not observed since the retention cutoff ($1), and rarely-observed // routes (observation_count < $2) not observed since the grace cutoff ($3). func (q *Queries) DeleteOldRoutes(ctx context.Context, arg DeleteOldRoutesParams) error { _, err := q.db.Exec(ctx, deleteOldRoutes, arg.LastSeen, arg.ObservationCount, arg.LastSeen_2) return err } const deleteOldTelemetry = `-- name: DeleteOldTelemetry :exec DELETE FROM observer_telemetry WHERE reported_at < $1 ` // Deletes telemetry rows older than the given cutoff. Called by the cleanup goroutine. func (q *Queries) DeleteOldTelemetry(ctx context.Context, reportedAt pgtype.Timestamptz) error { _, err := q.db.Exec(ctx, deleteOldTelemetry, reportedAt) 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 ` func (q *Queries) GetChannelByID(ctx context.Context, id int32) (Channel, error) { row := q.db.QueryRow(ctx, getChannelByID, id) var i Channel err := row.Scan( &i.ID, &i.ChannelHash, &i.KeyFingerprint, &i.Name, &i.Hashtag, &i.IsHashtag, &i.IsPublic, &i.KeyKnown, &i.FirstSeen, &i.LastSeen, &i.MessageCount, ) return i, err } const getCrossIATANeighbors = `-- name: GetCrossIATANeighbors :many SELECT n.id, n.name, n.node_type, n.latitude, n.longitude, nn.iata AS neighbor_iata, nn.observation_count, nn.last_seen, nn.snr FROM node_neighbors nn JOIN nodes n ON n.id = nn.neighbor_id WHERE nn.node_id = $1 AND nn.iata != $2 ORDER BY nn.last_seen DESC ` type GetCrossIATANeighborsParams struct { NodeID uuid.UUID `json:"node_id"` Iata string `json:"iata"` } type GetCrossIATANeighborsRow struct { ID uuid.UUID `json:"id"` Name *string `json:"name"` NodeType int16 `json:"node_type"` Latitude *float64 `json:"latitude"` Longitude *float64 `json:"longitude"` NeighborIata string `json:"neighbor_iata"` ObservationCount int64 `json:"observation_count"` LastSeen pgtype.Timestamptz `json:"last_seen"` Snr *float32 `json:"snr"` } // Returns neighbors of a node that are in a different IATA. func (q *Queries) GetCrossIATANeighbors(ctx context.Context, arg GetCrossIATANeighborsParams) ([]GetCrossIATANeighborsRow, error) { rows, err := q.db.Query(ctx, getCrossIATANeighbors, arg.NodeID, arg.Iata) if err != nil { return nil, err } defer rows.Close() items := []GetCrossIATANeighborsRow{} for rows.Next() { var i GetCrossIATANeighborsRow if err := rows.Scan( &i.ID, &i.Name, &i.NodeType, &i.Latitude, &i.Longitude, &i.NeighborIata, &i.ObservationCount, &i.LastSeen, &i.Snr, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const getHourlyStats = `-- name: GetHourlyStats :many SELECT iata, hour, observation_count, unique_packets, active_observers FROM mv_hourly_iata_stats WHERE (COALESCE(cardinality($1::bpchar[]), 0) = 0 OR iata = ANY($1::bpchar[])) AND hour >= NOW() - $2::interval ORDER BY iata, hour ` type GetHourlyStatsParams struct { Column1 []string `json:"column_1"` Column2 pgtype.Interval `json:"column_2"` } func (q *Queries) GetHourlyStats(ctx context.Context, arg GetHourlyStatsParams) ([]MvHourlyIataStat, error) { rows, err := q.db.Query(ctx, getHourlyStats, arg.Column1, arg.Column2) if err != nil { return nil, err } defer rows.Close() items := []MvHourlyIataStat{} for rows.Next() { var i MvHourlyIataStat if err := rows.Scan( &i.Iata, &i.Hour, &i.ObservationCount, &i.UniquePackets, &i.ActiveObservers, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const getIATA = `-- name: GetIATA :one SELECT iata, display_name, approx_lat, approx_lng, added_at, border FROM iata_codes WHERE iata = $1 ` func (q *Queries) GetIATA(ctx context.Context, iata string) (IataCode, error) { row := q.db.QueryRow(ctx, getIATA, iata) var i IataCode err := row.Scan( &i.Iata, &i.DisplayName, &i.ApproxLat, &i.ApproxLng, &i.AddedAt, &i.Border, ) return i, err } const getIATABorder = `-- name: GetIATABorder :one SELECT border FROM iata_codes WHERE iata = $1 ` // border is NULL when the IATA exists but has no border configured; a // missing row (unknown IATA) is sql.ErrNoRows, same not-found distinction // GetIATA already makes. func (q *Queries) GetIATABorder(ctx context.Context, iata string) ([]byte, error) { row := q.db.QueryRow(ctx, getIATABorder, iata) var border []byte err := row.Scan(&border) return border, err } const getKnownRoutesByNode = `-- name: GetKnownRoutesByNode :many SELECT id, node_ids, hash_prefix, iata, hop_count, first_seen, last_seen, observation_count FROM known_routes WHERE iata = $1 AND $2::uuid = ANY(node_ids) ORDER BY hop_count ASC, last_seen DESC ` type GetKnownRoutesByNodeParams struct { Iata string `json:"iata"` Column2 uuid.UUID `json:"column_2"` } type GetKnownRoutesByNodeRow struct { ID int64 `json:"id"` NodeIds []uuid.UUID `json:"node_ids"` HashPrefix [][]byte `json:"hash_prefix"` Iata string `json:"iata"` HopCount int32 `json:"hop_count"` FirstSeen pgtype.Timestamptz `json:"first_seen"` LastSeen pgtype.Timestamptz `json:"last_seen"` ObservationCount int64 `json:"observation_count"` } func (q *Queries) GetKnownRoutesByNode(ctx context.Context, arg GetKnownRoutesByNodeParams) ([]GetKnownRoutesByNodeRow, error) { rows, err := q.db.Query(ctx, getKnownRoutesByNode, arg.Iata, arg.Column2) if err != nil { return nil, err } defer rows.Close() items := []GetKnownRoutesByNodeRow{} for rows.Next() { var i GetKnownRoutesByNodeRow if err := rows.Scan( &i.ID, &i.NodeIds, &i.HashPrefix, &i.Iata, &i.HopCount, &i.FirstSeen, &i.LastSeen, &i.ObservationCount, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const getNodeByID = `-- name: GetNodeByID :one SELECT n.id, n.public_key, n.node_type, n.name, n.latitude, n.longitude, n.location_source, n.last_advert_at, n.supports_multibyte_paths, n.supports_multibyte_traces, n.default_scope_id, n.min_firmware_version, n.first_seen, n.last_seen, n.radio_freq_mhz, n.radio_sf, n.radio_bw_khz, n.metadata, n.device_clock_drift_seconds, ts.name AS default_scope_name, EXISTS (SELECT 1 FROM observers o WHERE o.public_key = n.public_key) AS is_observer, (SELECT o.id FROM observers o WHERE o.public_key = n.public_key LIMIT 1) AS observer_id, (SELECT json_agg(json_build_object('iata', ni.iata, 'lastHeard', (extract(epoch from ni.last_heard) * 1000)::bigint) ORDER BY ni.last_heard DESC) FROM node_iatas ni WHERE ni.node_id = n.id) AS iatas, (SELECT COUNT(DISTINCT nn.neighbor_id) FROM node_neighbors nn WHERE nn.node_id = n.id)::bigint AS known_neighbor_count FROM nodes n LEFT JOIN transport_scopes ts ON ts.id = n.default_scope_id WHERE n.id = $1 ` type GetNodeByIDRow struct { ID uuid.UUID `json:"id"` PublicKey []byte `json:"public_key"` NodeType int16 `json:"node_type"` Name *string `json:"name"` Latitude *float64 `json:"latitude"` Longitude *float64 `json:"longitude"` LocationSource *string `json:"location_source"` LastAdvertAt pgtype.Timestamptz `json:"last_advert_at"` SupportsMultibytePaths bool `json:"supports_multibyte_paths"` SupportsMultibyteTraces bool `json:"supports_multibyte_traces"` DefaultScopeID *int32 `json:"default_scope_id"` MinFirmwareVersion *string `json:"min_firmware_version"` FirstSeen pgtype.Timestamptz `json:"first_seen"` LastSeen pgtype.Timestamptz `json:"last_seen"` RadioFreqMhz *float32 `json:"radio_freq_mhz"` RadioSf *int16 `json:"radio_sf"` RadioBwKhz *float32 `json:"radio_bw_khz"` Metadata []byte `json:"metadata"` DeviceClockDriftSeconds *int32 `json:"device_clock_drift_seconds"` DefaultScopeName *string `json:"default_scope_name"` IsObserver bool `json:"is_observer"` ObserverID uuid.UUID `json:"observer_id"` Iatas []byte `json:"iatas"` KnownNeighborCount int64 `json:"known_neighbor_count"` } func (q *Queries) GetNodeByID(ctx context.Context, id uuid.UUID) (GetNodeByIDRow, error) { row := q.db.QueryRow(ctx, getNodeByID, id) var i GetNodeByIDRow err := row.Scan( &i.ID, &i.PublicKey, &i.NodeType, &i.Name, &i.Latitude, &i.Longitude, &i.LocationSource, &i.LastAdvertAt, &i.SupportsMultibytePaths, &i.SupportsMultibyteTraces, &i.DefaultScopeID, &i.MinFirmwareVersion, &i.FirstSeen, &i.LastSeen, &i.RadioFreqMhz, &i.RadioSf, &i.RadioBwKhz, &i.Metadata, &i.DeviceClockDriftSeconds, &i.DefaultScopeName, &i.IsObserver, &i.ObserverID, &i.Iatas, &i.KnownNeighborCount, ) return i, err } const getNodeByPubkey = `-- name: GetNodeByPubkey :one SELECT id FROM nodes WHERE public_key = $1 ` func (q *Queries) GetNodeByPubkey(ctx context.Context, publicKey []byte) (uuid.UUID, error) { row := q.db.QueryRow(ctx, getNodeByPubkey, publicKey) var id uuid.UUID err := row.Scan(&id) return id, err } const getNodeNeighbors = `-- name: GetNodeNeighbors :many SELECT n.id, n.public_key, n.name, n.node_type, n.latitude, n.longitude, nn.iata, nn.observation_count, nn.first_seen, nn.last_seen, nn.snr FROM node_neighbors nn JOIN nodes n ON n.id = nn.neighbor_id WHERE nn.node_id = $1 ORDER BY nn.last_seen DESC ` type GetNodeNeighborsRow struct { ID uuid.UUID `json:"id"` PublicKey []byte `json:"public_key"` Name *string `json:"name"` NodeType int16 `json:"node_type"` Latitude *float64 `json:"latitude"` Longitude *float64 `json:"longitude"` Iata string `json:"iata"` ObservationCount int64 `json:"observation_count"` FirstSeen pgtype.Timestamptz `json:"first_seen"` LastSeen pgtype.Timestamptz `json:"last_seen"` Snr *float32 `json:"snr"` } // Returns the neighbors of a node with details, ordered by most recently seen. func (q *Queries) GetNodeNeighbors(ctx context.Context, nodeID uuid.UUID) ([]GetNodeNeighborsRow, error) { rows, err := q.db.Query(ctx, getNodeNeighbors, nodeID) if err != nil { return nil, err } defer rows.Close() items := []GetNodeNeighborsRow{} for rows.Next() { var i GetNodeNeighborsRow if err := rows.Scan( &i.ID, &i.PublicKey, &i.Name, &i.NodeType, &i.Latitude, &i.Longitude, &i.Iata, &i.ObservationCount, &i.FirstSeen, &i.LastSeen, &i.Snr, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const getNodesByIDs = `-- name: GetNodesByIDs :many SELECT id, public_key, name, latitude, longitude FROM nodes WHERE id = ANY($1::uuid[]) ` type GetNodesByIDsRow struct { ID uuid.UUID `json:"id"` PublicKey []byte `json:"public_key"` Name *string `json:"name"` Latitude *float64 `json:"latitude"` Longitude *float64 `json:"longitude"` } func (q *Queries) GetNodesByIDs(ctx context.Context, dollar_1 []uuid.UUID) ([]GetNodesByIDsRow, error) { rows, err := q.db.Query(ctx, getNodesByIDs, dollar_1) if err != nil { return nil, err } defer rows.Close() items := []GetNodesByIDsRow{} for rows.Next() { var i GetNodesByIDsRow if err := rows.Scan( &i.ID, &i.PublicKey, &i.Name, &i.Latitude, &i.Longitude, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const getObserverBrokers = `-- name: GetObserverBrokers :many SELECT broker_name, last_seen, last_packet_at FROM observer_brokers WHERE observer_id = $1 ORDER BY last_seen DESC ` type GetObserverBrokersRow struct { BrokerName string `json:"broker_name"` LastSeen pgtype.Timestamptz `json:"last_seen"` LastPacketAt pgtype.Timestamptz `json:"last_packet_at"` } func (q *Queries) GetObserverBrokers(ctx context.Context, observerID uuid.UUID) ([]GetObserverBrokersRow, error) { rows, err := q.db.Query(ctx, getObserverBrokers, observerID) if err != nil { return nil, err } defer rows.Close() items := []GetObserverBrokersRow{} for rows.Next() { var i GetObserverBrokersRow if err := rows.Scan(&i.BrokerName, &i.LastSeen, &i.LastPacketAt); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const getObserverByID = `-- name: GetObserverByID :one SELECT id, public_key, display_name, observer_type, software_version, hardware_model, firmware_version, firmware_build, radio_freq_mhz, radio_sf, radio_bw_khz, radio_cr, battery_level, uptime_seconds, status_metadata, last_status_at, first_seen, last_seen, observation_count, metadata, region_scope FROM observers WHERE id = $1 ` func (q *Queries) GetObserverByID(ctx context.Context, id uuid.UUID) (Observer, error) { row := q.db.QueryRow(ctx, getObserverByID, id) var i Observer err := row.Scan( &i.ID, &i.PublicKey, &i.DisplayName, &i.ObserverType, &i.SoftwareVersion, &i.HardwareModel, &i.FirmwareVersion, &i.FirmwareBuild, &i.RadioFreqMhz, &i.RadioSf, &i.RadioBwKhz, &i.RadioCr, &i.BatteryLevel, &i.UptimeSeconds, &i.StatusMetadata, &i.LastStatusAt, &i.FirstSeen, &i.LastSeen, &i.ObservationCount, &i.Metadata, &i.RegionScope, ) return i, err } const getObserverByPubkey = `-- name: GetObserverByPubkey :one SELECT id, public_key, display_name, observer_type, software_version, hardware_model, firmware_version, firmware_build, radio_freq_mhz, radio_sf, radio_bw_khz, radio_cr, battery_level, uptime_seconds, status_metadata, last_status_at, first_seen, last_seen, observation_count, metadata, region_scope FROM observers WHERE public_key = $1 ` func (q *Queries) GetObserverByPubkey(ctx context.Context, publicKey []byte) (Observer, error) { row := q.db.QueryRow(ctx, getObserverByPubkey, publicKey) var i Observer err := row.Scan( &i.ID, &i.PublicKey, &i.DisplayName, &i.ObserverType, &i.SoftwareVersion, &i.HardwareModel, &i.FirmwareVersion, &i.FirmwareBuild, &i.RadioFreqMhz, &i.RadioSf, &i.RadioBwKhz, &i.RadioCr, &i.BatteryLevel, &i.UptimeSeconds, &i.StatusMetadata, &i.LastStatusAt, &i.FirstSeen, &i.LastSeen, &i.ObservationCount, &i.Metadata, &i.RegionScope, ) return i, err } const getObserverLastIATA = `-- name: GetObserverLastIATA :one SELECT iata FROM packet_observations WHERE observer_id = $1 ORDER BY heard_at DESC LIMIT 1 ` func (q *Queries) GetObserverLastIATA(ctx context.Context, observerID uuid.UUID) (string, error) { row := q.db.QueryRow(ctx, getObserverLastIATA, observerID) var iata string err := row.Scan(&iata) return iata, err } const getObserverRadio = `-- name: GetObserverRadio :one SELECT radio_freq_mhz, radio_bw_khz, radio_sf, radio_cr FROM observers WHERE id = $1 ` type GetObserverRadioRow struct { RadioFreqMhz *float32 `json:"radio_freq_mhz"` RadioBwKhz *float32 `json:"radio_bw_khz"` RadioSf *int16 `json:"radio_sf"` RadioCr *int16 `json:"radio_cr"` } func (q *Queries) GetObserverRadio(ctx context.Context, id uuid.UUID) (GetObserverRadioRow, error) { row := q.db.QueryRow(ctx, getObserverRadio, id) var i GetObserverRadioRow err := row.Scan( &i.RadioFreqMhz, &i.RadioBwKhz, &i.RadioSf, &i.RadioCr, ) return i, err } const getObserverScopes = `-- name: GetObserverScopes :many SELECT ts.name FROM observer_scopes os JOIN transport_scopes ts ON ts.id = os.scope_id WHERE os.observer_id = $1 ORDER BY ts.name ` func (q *Queries) GetObserverScopes(ctx context.Context, observerID uuid.UUID) ([]string, error) { rows, err := q.db.Query(ctx, getObserverScopes, observerID) if err != nil { return nil, err } defer rows.Close() items := []string{} for rows.Next() { var name string if err := rows.Scan(&name); err != nil { return nil, err } items = append(items, name) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const getObserverTelemetry = `-- name: GetObserverTelemetry :many SELECT id, reported_at, battery_voltage_mv, airtime_tx_pct, airtime_rx_pct, noise_floor_db, uptime_seconds, queue_length, debug_flags, receive_errors FROM observer_telemetry WHERE observer_id = $1 AND ($2::timestamptz IS NULL OR reported_at >= $2) AND ($3::timestamptz IS NULL OR reported_at <= $3) AND ($4 = 0 OR id > $4) ORDER BY reported_at ASC ` type GetObserverTelemetryParams struct { ObserverID uuid.UUID `json:"observer_id"` Column2 pgtype.Timestamptz `json:"column_2"` Column3 pgtype.Timestamptz `json:"column_3"` Column4 interface{} `json:"column_4"` } type GetObserverTelemetryRow struct { ID int64 `json:"id"` ReportedAt pgtype.Timestamptz `json:"reported_at"` BatteryVoltageMv *int32 `json:"battery_voltage_mv"` AirtimeTxPct *float32 `json:"airtime_tx_pct"` AirtimeRxPct *float32 `json:"airtime_rx_pct"` NoiseFloorDb *float32 `json:"noise_floor_db"` UptimeSeconds *int64 `json:"uptime_seconds"` QueueLength *int32 `json:"queue_length"` DebugFlags *int32 `json:"debug_flags"` ReceiveErrors *int32 `json:"receive_errors"` } func (q *Queries) GetObserverTelemetry(ctx context.Context, arg GetObserverTelemetryParams) ([]GetObserverTelemetryRow, error) { rows, err := q.db.Query(ctx, getObserverTelemetry, arg.ObserverID, arg.Column2, arg.Column3, arg.Column4, ) if err != nil { return nil, err } defer rows.Close() items := []GetObserverTelemetryRow{} for rows.Next() { var i GetObserverTelemetryRow if err := rows.Scan( &i.ID, &i.ReportedAt, &i.BatteryVoltageMv, &i.AirtimeTxPct, &i.AirtimeRxPct, &i.NoiseFloorDb, &i.UptimeSeconds, &i.QueueLength, &i.DebugFlags, &i.ReceiveErrors, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const getObserverTelemetryBucketed = `-- name: GetObserverTelemetryBucketed :many SELECT (date_trunc('day', reported_at) + (EXTRACT(HOUR FROM reported_at)::int / $4::int) * ($4::int * interval '1 hour'))::timestamptz AS bucket, AVG(battery_voltage_mv)::int AS battery_voltage_mv, GREATEST(MAX(airtime_tx_pct) - MIN(airtime_tx_pct), 0)::real AS airtime_tx_pct, GREATEST(MAX(airtime_rx_pct) - MIN(airtime_rx_pct), 0)::real AS airtime_rx_pct, AVG(noise_floor_db)::real AS noise_floor_db, MAX(uptime_seconds)::bigint AS uptime_seconds, AVG(queue_length)::int AS queue_length, GREATEST(MAX(receive_errors) - MIN(receive_errors), 0)::int AS receive_errors FROM observer_telemetry WHERE observer_id = $1 AND ($2::timestamptz IS NULL OR reported_at >= $2) AND ($3::timestamptz IS NULL OR reported_at <= $3) GROUP BY bucket ORDER BY bucket ASC ` type GetObserverTelemetryBucketedParams struct { ObserverID uuid.UUID `json:"observer_id"` Column2 pgtype.Timestamptz `json:"column_2"` Column3 pgtype.Timestamptz `json:"column_3"` Column4 int32 `json:"column_4"` } type GetObserverTelemetryBucketedRow struct { Bucket pgtype.Timestamptz `json:"bucket"` BatteryVoltageMv int32 `json:"battery_voltage_mv"` AirtimeTxPct float32 `json:"airtime_tx_pct"` AirtimeRxPct float32 `json:"airtime_rx_pct"` NoiseFloorDb float32 `json:"noise_floor_db"` UptimeSeconds int64 `json:"uptime_seconds"` QueueLength int32 `json:"queue_length"` ReceiveErrors int32 `json:"receive_errors"` } func (q *Queries) GetObserverTelemetryBucketed(ctx context.Context, arg GetObserverTelemetryBucketedParams) ([]GetObserverTelemetryBucketedRow, error) { rows, err := q.db.Query(ctx, getObserverTelemetryBucketed, arg.ObserverID, arg.Column2, arg.Column3, arg.Column4, ) if err != nil { return nil, err } defer rows.Close() items := []GetObserverTelemetryBucketedRow{} for rows.Next() { var i GetObserverTelemetryBucketedRow if err := rows.Scan( &i.Bucket, &i.BatteryVoltageMv, &i.AirtimeTxPct, &i.AirtimeRxPct, &i.NoiseFloorDb, &i.UptimeSeconds, &i.QueueLength, &i.ReceiveErrors, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const getPacketByHash = `-- name: GetPacketByHash :one SELECT p.packet_hash, p.payload_type, p.payload_version, p.route_type, p.transport_codes_present, p.region_code, p.sub_region_code, p.scope_id, p.origin_pubkey, p.raw_payload, p.raw_header, p.parsed_payload, p.decrypted, p.channel_hash, p.trace_tag, p.first_heard_at, p.last_heard_at, ts.name AS scope_name, cm.sender_name AS cm_sender_name, cm.content AS cm_content, cm.sent_at AS cm_sent_at FROM packets p LEFT JOIN transport_scopes ts ON ts.id = p.scope_id LEFT JOIN channel_messages cm ON cm.packet_hash = p.packet_hash WHERE p.packet_hash = $1 ` type GetPacketByHashRow struct { PacketHash []byte `json:"packet_hash"` PayloadType int16 `json:"payload_type"` PayloadVersion int16 `json:"payload_version"` RouteType int16 `json:"route_type"` TransportCodesPresent *bool `json:"transport_codes_present"` RegionCode *int32 `json:"region_code"` SubRegionCode *int32 `json:"sub_region_code"` ScopeID *int32 `json:"scope_id"` OriginPubkey []byte `json:"origin_pubkey"` RawPayload []byte `json:"raw_payload"` RawHeader []byte `json:"raw_header"` ParsedPayload []byte `json:"parsed_payload"` Decrypted *bool `json:"decrypted"` ChannelHash []byte `json:"channel_hash"` TraceTag []byte `json:"trace_tag"` FirstHeardAt pgtype.Timestamptz `json:"first_heard_at"` LastHeardAt pgtype.Timestamptz `json:"last_heard_at"` ScopeName *string `json:"scope_name"` CmSenderName *string `json:"cm_sender_name"` CmContent *string `json:"cm_content"` CmSentAt pgtype.Timestamptz `json:"cm_sent_at"` } func (q *Queries) GetPacketByHash(ctx context.Context, packetHash []byte) (GetPacketByHashRow, error) { row := q.db.QueryRow(ctx, getPacketByHash, packetHash) var i GetPacketByHashRow err := row.Scan( &i.PacketHash, &i.PayloadType, &i.PayloadVersion, &i.RouteType, &i.TransportCodesPresent, &i.RegionCode, &i.SubRegionCode, &i.ScopeID, &i.OriginPubkey, &i.RawPayload, &i.RawHeader, &i.ParsedPayload, &i.Decrypted, &i.ChannelHash, &i.TraceTag, &i.FirstHeardAt, &i.LastHeardAt, &i.ScopeName, &i.CmSenderName, &i.CmContent, &i.CmSentAt, ) return i, err } const getPacketObservationCount = `-- name: GetPacketObservationCount :one SELECT COUNT(*) FROM packet_observations WHERE packet_hash = $1 ` func (q *Queries) GetPacketObservationCount(ctx context.Context, packetHash []byte) (int64, error) { row := q.db.QueryRow(ctx, getPacketObservationCount, packetHash) var count int64 err := row.Scan(&count) return count, err } const getPacketsByTraceTag = `-- name: GetPacketsByTraceTag :many SELECT encode(p.packet_hash, 'hex') AS packet_hash_hex, p.route_type, p.first_heard_at, p.last_heard_at, p.parsed_payload, p.scope_id, ts.name AS scope_name FROM packets p LEFT JOIN transport_scopes ts ON ts.id = p.scope_id WHERE p.trace_tag = decode($1, 'hex') ORDER BY p.first_heard_at ASC ` type GetPacketsByTraceTagRow struct { PacketHashHex string `json:"packet_hash_hex"` RouteType int16 `json:"route_type"` FirstHeardAt pgtype.Timestamptz `json:"first_heard_at"` LastHeardAt pgtype.Timestamptz `json:"last_heard_at"` ParsedPayload []byte `json:"parsed_payload"` ScopeID *int32 `json:"scope_id"` ScopeName *string `json:"scope_name"` } // Returns all packets for a given trace tag with observations. func (q *Queries) GetPacketsByTraceTag(ctx context.Context, decode string) ([]GetPacketsByTraceTagRow, error) { rows, err := q.db.Query(ctx, getPacketsByTraceTag, decode) if err != nil { return nil, err } defer rows.Close() items := []GetPacketsByTraceTagRow{} for rows.Next() { var i GetPacketsByTraceTagRow if err := rows.Scan( &i.PacketHashHex, &i.RouteType, &i.FirstHeardAt, &i.LastHeardAt, &i.ParsedPayload, &i.ScopeID, &i.ScopeName, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const getRadioPresets = `-- name: GetRadioPresets :many SELECT preset, iata, source_type, count FROM mv_radio_presets WHERE ($1::text = '' OR preset = $1::text) AND (COALESCE(cardinality($2::bpchar[]), 0) = 0 OR iata = ANY($2::bpchar[])) ORDER BY preset, iata, source_type ` type GetRadioPresetsParams struct { Column1 string `json:"column_1"` Column2 []string `json:"column_2"` } func (q *Queries) GetRadioPresets(ctx context.Context, arg GetRadioPresetsParams) ([]MvRadioPreset, error) { rows, err := q.db.Query(ctx, getRadioPresets, arg.Column1, arg.Column2) if err != nil { return nil, err } defer rows.Close() items := []MvRadioPreset{} for rows.Next() { var i MvRadioPreset if err := rows.Scan( &i.Preset, &i.Iata, &i.SourceType, &i.Count, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const getRegion = `-- name: GetRegion :one SELECT id, slug, name, description, center_lat, center_lng, zoom_level FROM regions WHERE id = $1 ` type GetRegionRow struct { ID int32 `json:"id"` Slug string `json:"slug"` Name string `json:"name"` Description *string `json:"description"` CenterLat *float64 `json:"center_lat"` CenterLng *float64 `json:"center_lng"` ZoomLevel *int32 `json:"zoom_level"` } func (q *Queries) GetRegion(ctx context.Context, id int32) (GetRegionRow, error) { row := q.db.QueryRow(ctx, getRegion, id) var i GetRegionRow err := row.Scan( &i.ID, &i.Slug, &i.Name, &i.Description, &i.CenterLat, &i.CenterLng, &i.ZoomLevel, ) return i, err } const getRegionBySlug = `-- name: GetRegionBySlug :one SELECT id, slug, name, description, center_lat, center_lng, zoom_level FROM regions WHERE slug = $1 ` type GetRegionBySlugRow struct { ID int32 `json:"id"` Slug string `json:"slug"` Name string `json:"name"` Description *string `json:"description"` CenterLat *float64 `json:"center_lat"` CenterLng *float64 `json:"center_lng"` ZoomLevel *int32 `json:"zoom_level"` } func (q *Queries) GetRegionBySlug(ctx context.Context, slug string) (GetRegionBySlugRow, error) { row := q.db.QueryRow(ctx, getRegionBySlug, slug) var i GetRegionBySlugRow err := row.Scan( &i.ID, &i.Slug, &i.Name, &i.Description, &i.CenterLat, &i.CenterLng, &i.ZoomLevel, ) return i, err } const getRegionIATAs = `-- name: GetRegionIATAs :many SELECT iata FROM region_iatas WHERE region_id = $1 ORDER BY iata ` func (q *Queries) GetRegionIATAs(ctx context.Context, regionID int32) ([]string, error) { rows, err := q.db.Query(ctx, getRegionIATAs, regionID) if err != nil { return nil, err } defer rows.Close() items := []string{} for rows.Next() { var iata string if err := rows.Scan(&iata); err != nil { return nil, err } items = append(items, iata) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const getScopeByName = `-- name: GetScopeByName :one 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, COUNT(DISTINCT po.iata) AS iata_count, array_remove(array_agg(DISTINCT po.iata ORDER BY po.iata), NULL)::text[] AS iatas 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 observers o ON o.id = os.observer_id LEFT JOIN packet_observations po ON po.observer_id = o.id LEFT JOIN nodes n ON n.default_scope_id = ts.id WHERE ts.name = $1 GROUP BY ts.name ` type GetScopeByNameRow struct { Name string `json:"name"` PacketCount int64 `json:"packet_count"` ObserverCount int64 `json:"observer_count"` NodeCount int64 `json:"node_count"` IataCount int64 `json:"iata_count"` Iatas []string `json:"iatas"` } func (q *Queries) GetScopeByName(ctx context.Context, name string) (GetScopeByNameRow, error) { row := q.db.QueryRow(ctx, getScopeByName, name) var i GetScopeByNameRow err := row.Scan( &i.Name, &i.PacketCount, &i.ObserverCount, &i.NodeCount, &i.IataCount, &i.Iatas, ) return i, err } const getScopeNames = `-- name: GetScopeNames :many SELECT name FROM transport_scopes ORDER BY name ` func (q *Queries) GetScopeNames(ctx context.Context) ([]string, error) { rows, err := q.db.Query(ctx, getScopeNames) if err != nil { return nil, err } defer rows.Close() items := []string{} for rows.Next() { var name string if err := rows.Scan(&name); err != nil { return nil, err } items = append(items, name) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const getScopeStats = `-- name: GetScopeStats :many SELECT ts.name, (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 ORDER BY ts.name ` type GetScopeStatsRow struct { Name string `json:"name"` PacketCount int64 `json:"packet_count"` ObserverCount int64 `json:"observer_count"` 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 { return nil, err } defer rows.Close() items := []GetScopeStatsRow{} for rows.Next() { var i GetScopeStatsRow if err := rows.Scan( &i.Name, &i.PacketCount, &i.ObserverCount, &i.NodeCount, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const getScopesByIATAs = `-- name: GetScopesByIATAs :many SELECT ts.name, COUNT(DISTINCT os.observer_id) AS observer_count, COUNT(DISTINCT n.id) AS node_count, COUNT(DISTINCT po.iata) AS iata_count FROM transport_scopes ts LEFT JOIN observer_scopes os ON os.scope_id = ts.id LEFT JOIN observers o ON o.id = os.observer_id LEFT JOIN packet_observations po ON po.observer_id = o.id LEFT JOIN nodes n ON n.default_scope_id = ts.id WHERE (COALESCE(cardinality($1::bpchar[]), 0) = 0 OR po.iata = ANY($1::bpchar[])) GROUP BY ts.name ORDER BY ts.name ` type GetScopesByIATAsRow struct { Name string `json:"name"` ObserverCount int64 `json:"observer_count"` NodeCount int64 `json:"node_count"` IataCount int64 `json:"iata_count"` } func (q *Queries) GetScopesByIATAs(ctx context.Context, dollar_1 []string) ([]GetScopesByIATAsRow, error) { rows, err := q.db.Query(ctx, getScopesByIATAs, dollar_1) if err != nil { return nil, err } defer rows.Close() items := []GetScopesByIATAsRow{} for rows.Next() { var i GetScopesByIATAsRow if err := rows.Scan( &i.Name, &i.ObserverCount, &i.NodeCount, &i.IataCount, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const getStatsClockDrift = `-- name: GetStatsClockDrift :many SELECT n.id, n.name, n.node_type, n.device_clock_drift_seconds, n.last_advert_at, json_agg(json_build_object('iata', ni.iata, 'lastHeard', (extract(epoch from ni.last_heard) * 1000)::bigint) ORDER BY ni.last_heard DESC) FILTER (WHERE ni.iata IS NOT NULL) AS iatas FROM nodes n LEFT JOIN node_iatas ni ON ni.node_id = n.id WHERE n.node_type IN (2, 3) AND n.device_clock_drift_seconds IS NOT NULL AND ABS(n.device_clock_drift_seconds) > $1::int AND (COALESCE(cardinality($2::bpchar[]), 0) = 0 OR n.id IN (SELECT node_id FROM node_iatas WHERE iata = ANY($2::bpchar[]))) GROUP BY n.id, n.name, n.node_type, n.device_clock_drift_seconds, n.last_advert_at ORDER BY ABS(n.device_clock_drift_seconds) DESC LIMIT $3 ` type GetStatsClockDriftParams struct { Column1 int32 `json:"column_1"` Column2 []string `json:"column_2"` Limit int32 `json:"limit"` } type GetStatsClockDriftRow struct { ID uuid.UUID `json:"id"` Name *string `json:"name"` NodeType int16 `json:"node_type"` DeviceClockDriftSeconds *int32 `json:"device_clock_drift_seconds"` LastAdvertAt pgtype.Timestamptz `json:"last_advert_at"` Iatas []byte `json:"iatas"` } // Repeaters/room servers (node_type 2/3) whose current advert-derived clock drift exceeds // the given threshold in magnitude, worst first. Not time-windowed -- reflects each node's // latest measured drift, not an aggregate over a period. func (q *Queries) GetStatsClockDrift(ctx context.Context, arg GetStatsClockDriftParams) ([]GetStatsClockDriftRow, error) { rows, err := q.db.Query(ctx, getStatsClockDrift, arg.Column1, arg.Column2, arg.Limit) if err != nil { return nil, err } defer rows.Close() items := []GetStatsClockDriftRow{} for rows.Next() { var i GetStatsClockDriftRow if err := rows.Scan( &i.ID, &i.Name, &i.NodeType, &i.DeviceClockDriftSeconds, &i.LastAdvertAt, &i.Iatas, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const getStatsNodeTypes = `-- name: GetStatsNodeTypes :many SELECT n.node_type, COUNT(DISTINCT n.id)::bigint AS count FROM nodes n LEFT JOIN node_iatas ni ON ni.node_id = n.id WHERE (COALESCE(cardinality($1::bpchar[]), 0) = 0 OR ni.iata = ANY($1::bpchar[])) GROUP BY n.node_type ORDER BY count DESC ` type GetStatsNodeTypesRow struct { NodeType int16 `json:"node_type"` Count int64 `json:"count"` } // Returns node counts grouped by type, optionally filtered by IATA. func (q *Queries) GetStatsNodeTypes(ctx context.Context, dollar_1 []string) ([]GetStatsNodeTypesRow, error) { rows, err := q.db.Query(ctx, getStatsNodeTypes, dollar_1) if err != nil { return nil, err } defer rows.Close() items := []GetStatsNodeTypesRow{} for rows.Next() { var i GetStatsNodeTypesRow if err := rows.Scan(&i.NodeType, &i.Count); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const getStatsOverview = `-- name: GetStatsOverview :one SELECT COUNT(DISTINCT po.packet_hash) AS total_packets, COUNT(*) AS total_observations, COUNT(DISTINCT po.observer_id) AS active_observers, COUNT(DISTINCT po.iata) AS active_iatas FROM packet_observations po WHERE po.heard_at > NOW() - INTERVAL '24 hours' AND (COALESCE(cardinality($1::bpchar[]), 0) = 0 OR po.iata = ANY($1::bpchar[])) ` type GetStatsOverviewRow struct { TotalPackets int64 `json:"total_packets"` TotalObservations int64 `json:"total_observations"` ActiveObservers int64 `json:"active_observers"` ActiveIatas int64 `json:"active_iatas"` } // ============================================================ // STATS // ============================================================ func (q *Queries) GetStatsOverview(ctx context.Context, dollar_1 []string) (GetStatsOverviewRow, error) { row := q.db.QueryRow(ctx, getStatsOverview, dollar_1) var i GetStatsOverviewRow err := row.Scan( &i.TotalPackets, &i.TotalObservations, &i.ActiveObservers, &i.ActiveIatas, ) return i, err } const getStatsPayloadBreakdown = `-- name: GetStatsPayloadBreakdown :many SELECT 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[])) AND bucket >= NOW() - $2::interval GROUP BY payload_type ORDER BY count DESC ` type GetStatsPayloadBreakdownParams struct { Column1 []string `json:"column_1"` Column2 pgtype.Interval `json:"column_2"` } type GetStatsPayloadBreakdownRow struct { PayloadType *int16 `json:"payload_type"` Count int64 `json:"count"` } // Payload-type counts for the IATA within the window, summed from the // precomputed hourly buckets. func (q *Queries) GetStatsPayloadBreakdown(ctx context.Context, arg GetStatsPayloadBreakdownParams) ([]GetStatsPayloadBreakdownRow, error) { rows, err := q.db.Query(ctx, getStatsPayloadBreakdown, arg.Column1, arg.Column2) if err != nil { return nil, err } defer rows.Close() items := []GetStatsPayloadBreakdownRow{} for rows.Next() { var i GetStatsPayloadBreakdownRow if err := rows.Scan(&i.PayloadType, &i.Count); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const getStatsTopAdvertisers = `-- name: GetStatsTopAdvertisers :many SELECT node_id AS id, name, node_type, COALESCE(SUM(advert_count), 0)::bigint AS advert_count, COALESCE(SUM(flood_advert_count), 0)::bigint AS flood_advert_count, COALESCE(SUM(direct_advert_count), 0)::bigint AS direct_advert_count, MAX(last_heard)::timestamptz AS last_heard, COALESCE(MAX(iata), '')::bpchar AS iata FROM mv_top_advertisers_by_iata WHERE bucket >= NOW() - $1::interval AND (COALESCE(cardinality($2::bpchar[]), 0) = 0 OR iata = ANY($2::bpchar[])) GROUP BY node_id, name, node_type ORDER BY advert_count DESC LIMIT $3 ` type GetStatsTopAdvertisersParams struct { Column1 pgtype.Interval `json:"column_1"` Column2 []string `json:"column_2"` Limit int32 `json:"limit"` } type GetStatsTopAdvertisersRow struct { ID uuid.UUID `json:"id"` Name *string `json:"name"` NodeType int16 `json:"node_type"` AdvertCount int64 `json:"advert_count"` FloodAdvertCount int64 `json:"flood_advert_count"` DirectAdvertCount int64 `json:"direct_advert_count"` LastHeard pgtype.Timestamptz `json:"last_heard"` Iata string `json:"iata"` } // Top N advertisers in the window, summed from the hourly buckets. func (q *Queries) GetStatsTopAdvertisers(ctx context.Context, arg GetStatsTopAdvertisersParams) ([]GetStatsTopAdvertisersRow, error) { rows, err := q.db.Query(ctx, getStatsTopAdvertisers, arg.Column1, arg.Column2, arg.Limit) if err != nil { return nil, err } defer rows.Close() items := []GetStatsTopAdvertisersRow{} for rows.Next() { var i GetStatsTopAdvertisersRow if err := rows.Scan( &i.ID, &i.Name, &i.NodeType, &i.AdvertCount, &i.FloodAdvertCount, &i.DirectAdvertCount, &i.LastHeard, &i.Iata, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const getStatsTopObservers = `-- name: GetStatsTopObservers :many SELECT 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 bucket >= NOW() - $1::interval 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 { Column1 pgtype.Interval `json:"column_1"` 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 string `json:"iata"` } // Top N observers for the IATA within the window, summed from the precomputed // hourly buckets. Counts sum across matched IATAs; iata is a representative one. func (q *Queries) GetStatsTopObservers(ctx context.Context, arg GetStatsTopObserversParams) ([]GetStatsTopObserversRow, error) { rows, err := q.db.Query(ctx, getStatsTopObservers, arg.Column1, arg.Column2, arg.Limit) if err != nil { return nil, err } defer rows.Close() items := []GetStatsTopObserversRow{} for rows.Next() { var i GetStatsTopObserversRow if err := rows.Scan( &i.ID, &i.DisplayName, &i.ObserverType, &i.ObservationCount, &i.Iata, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const getStatsTopTalkers = `-- name: GetStatsTopTalkers :many SELECT sender_name, COALESCE(SUM(message_count), 0)::bigint AS message_count, MAX(last_sent)::timestamptz AS last_sent FROM mv_top_talkers_by_iata WHERE bucket >= NOW() - $1::interval AND (COALESCE(cardinality($2::bpchar[]), 0) = 0 OR iata = ANY($2::bpchar[])) GROUP BY sender_name ORDER BY message_count DESC LIMIT $3 ` type GetStatsTopTalkersParams struct { Column1 pgtype.Interval `json:"column_1"` Column2 []string `json:"column_2"` Limit int32 `json:"limit"` } type GetStatsTopTalkersRow struct { SenderName *string `json:"sender_name"` MessageCount int64 `json:"message_count"` LastSent pgtype.Timestamptz `json:"last_sent"` } // Top N talkers (by decrypted sender_name) in the window, summed from the hourly buckets. func (q *Queries) GetStatsTopTalkers(ctx context.Context, arg GetStatsTopTalkersParams) ([]GetStatsTopTalkersRow, error) { rows, err := q.db.Query(ctx, getStatsTopTalkers, arg.Column1, arg.Column2, arg.Limit) if err != nil { return nil, err } defer rows.Close() items := []GetStatsTopTalkersRow{} for rows.Next() { var i GetStatsTopTalkersRow if err := rows.Scan(&i.SenderName, &i.MessageCount, &i.LastSent); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const getTopNodes = `-- name: GetTopNodes :many SELECT iata, node_id, name, node_type, observation_count, last_heard FROM mv_top_nodes_by_iata WHERE (COALESCE(cardinality($1::bpchar[]), 0) = 0 OR iata = ANY($1::bpchar[])) ORDER BY observation_count DESC LIMIT $2 ` type GetTopNodesParams struct { Column1 []string `json:"column_1"` Limit int32 `json:"limit"` } func (q *Queries) GetTopNodes(ctx context.Context, arg GetTopNodesParams) ([]MvTopNodesByIatum, error) { rows, err := q.db.Query(ctx, getTopNodes, arg.Column1, arg.Limit) if err != nil { return nil, err } defer rows.Close() items := []MvTopNodesByIatum{} for rows.Next() { var i MvTopNodesByIatum if err := rows.Scan( &i.Iata, &i.NodeID, &i.Name, &i.NodeType, &i.ObservationCount, &i.LastHeard, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const getTransportScopeByName = `-- name: GetTransportScopeByName :one SELECT id FROM transport_scopes WHERE name = $1 ` func (q *Queries) GetTransportScopeByName(ctx context.Context, name string) (int32, error) { row := q.db.QueryRow(ctx, getTransportScopeByName, name) var id int32 err := row.Scan(&id) return id, err } const getTransportScopes = `-- name: GetTransportScopes :many SELECT name, transport_key, key_fingerprint FROM transport_scopes ORDER BY name ` type GetTransportScopesRow struct { Name string `json:"name"` TransportKey []byte `json:"transport_key"` KeyFingerprint []byte `json:"key_fingerprint"` } func (q *Queries) GetTransportScopes(ctx context.Context) ([]GetTransportScopesRow, error) { rows, err := q.db.Query(ctx, getTransportScopes) if err != nil { return nil, err } defer rows.Close() items := []GetTransportScopesRow{} for rows.Next() { var i GetTransportScopesRow if err := rows.Scan(&i.Name, &i.TransportKey, &i.KeyFingerprint); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const insertChannelMessage = `-- name: InsertChannelMessage :one INSERT INTO channel_messages (channel_id, packet_hash, sender_name, content, sent_at) VALUES ($1, $2, $3, $4, $5) ON CONFLICT (packet_hash) DO NOTHING RETURNING id ` type InsertChannelMessageParams struct { ChannelID int32 `json:"channel_id"` PacketHash []byte `json:"packet_hash"` SenderName *string `json:"sender_name"` Content *string `json:"content"` SentAt pgtype.Timestamptz `json:"sent_at"` } // ============================================================ // CHANNEL MESSAGES // ============================================================ func (q *Queries) InsertChannelMessage(ctx context.Context, arg InsertChannelMessageParams) (int64, error) { row := q.db.QueryRow(ctx, insertChannelMessage, arg.ChannelID, arg.PacketHash, arg.SenderName, arg.Content, arg.SentAt, ) var id int64 err := row.Scan(&id) return id, err } const insertObservation = `-- name: InsertObservation :one INSERT INTO packet_observations ( 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 ) VALUES ( $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, payload_type ` type InsertObservationParams struct { PacketHash []byte `json:"packet_hash"` ObserverID uuid.UUID `json:"observer_id"` Iata string `json:"iata"` HeardAt pgtype.Timestamptz `json:"heard_at"` PathLengthByte int16 `json:"path_length_byte"` HashSize int16 `json:"hash_size"` HopCount int16 `json:"hop_count"` PathBytes []byte `json:"path_bytes"` Rssi *int16 `json:"rssi"` Snr *float32 `json:"snr"` PropagationTimeMs *int32 `json:"propagation_time_ms"` RadioFreqMhz *float32 `json:"radio_freq_mhz"` SpreadFactor *int16 `json:"spread_factor"` BandwidthKhz *float32 `json:"bandwidth_khz"` CodingRate *int16 `json:"coding_rate"` SourceBroker *string `json:"source_broker"` PayloadType *int16 `json:"payload_type"` } // ============================================================ // PACKET OBSERVATIONS // ============================================================ func (q *Queries) InsertObservation(ctx context.Context, arg InsertObservationParams) (PacketObservation, error) { row := q.db.QueryRow(ctx, insertObservation, arg.PacketHash, arg.ObserverID, arg.Iata, arg.HeardAt, arg.PathLengthByte, arg.HashSize, arg.HopCount, arg.PathBytes, arg.Rssi, arg.Snr, arg.PropagationTimeMs, arg.RadioFreqMhz, arg.SpreadFactor, arg.BandwidthKhz, arg.CodingRate, arg.SourceBroker, arg.PayloadType, ) var i PacketObservation err := row.Scan( &i.ID, &i.PacketHash, &i.ObserverID, &i.Iata, &i.HeardAt, &i.PathLengthByte, &i.HashSize, &i.HopCount, &i.PathBytes, &i.Rssi, &i.Snr, &i.PropagationTimeMs, &i.RadioFreqMhz, &i.SpreadFactor, &i.BandwidthKhz, &i.CodingRate, &i.SourceBroker, &i.PayloadType, ) return i, err } const insertObserverTelemetry = `-- name: InsertObserverTelemetry :exec INSERT INTO observer_telemetry ( observer_id, reported_at, battery_voltage_mv, airtime_tx_pct, airtime_rx_pct, noise_floor_db, uptime_seconds, queue_length, debug_flags, receive_errors ) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10) ON CONFLICT (observer_id, reported_at) DO NOTHING ` type InsertObserverTelemetryParams struct { ObserverID uuid.UUID `json:"observer_id"` ReportedAt pgtype.Timestamptz `json:"reported_at"` BatteryVoltageMv *int32 `json:"battery_voltage_mv"` AirtimeTxPct *float32 `json:"airtime_tx_pct"` AirtimeRxPct *float32 `json:"airtime_rx_pct"` NoiseFloorDb *float32 `json:"noise_floor_db"` UptimeSeconds *int64 `json:"uptime_seconds"` QueueLength *int32 `json:"queue_length"` DebugFlags *int32 `json:"debug_flags"` ReceiveErrors *int32 `json:"receive_errors"` } // Inserts a telemetry snapshot for an observer. The reported_at timestamp should // be truncated to the configured resolution before calling to ensure deduplication. func (q *Queries) InsertObserverTelemetry(ctx context.Context, arg InsertObserverTelemetryParams) error { _, err := q.db.Exec(ctx, insertObserverTelemetry, arg.ObserverID, arg.ReportedAt, arg.BatteryVoltageMv, arg.AirtimeTxPct, arg.AirtimeRxPct, arg.NoiseFloorDb, arg.UptimeSeconds, arg.QueueLength, arg.DebugFlags, arg.ReceiveErrors, ) return err } const listAllChannelMessages = `-- name: ListAllChannelMessages :many SELECT DISTINCT ON (cm.id) cm.id, cm.channel_id, cm.packet_hash, cm.sender_name, cm.sender_pubkey, cm.content, cm.sent_at, encode(cm.packet_hash, 'hex') as packet_hash_hex, c.channel_hash, (SELECT COUNT(*) FROM packet_observations po2 WHERE po2.packet_hash = cm.packet_hash) AS observation_count FROM channel_messages cm JOIN channels c ON c.id = cm.channel_id JOIN packet_observations po ON po.packet_hash = cm.packet_hash JOIN packets p ON p.packet_hash = cm.packet_hash LEFT JOIN transport_scopes ts ON ts.id = p.scope_id WHERE ($1::timestamptz IS NULL OR cm.sent_at >= $1) AND (COALESCE(cardinality($2::bpchar[]), 0) = 0 OR po.iata = ANY($2::bpchar[])) AND ($3::text = '' OR ts.name = $3::text) AND ($4 = 0 OR cm.id < $4) ORDER BY cm.id DESC LIMIT $5 ` type ListAllChannelMessagesParams struct { Column1 pgtype.Timestamptz `json:"column_1"` Column2 []string `json:"column_2"` Column3 string `json:"column_3"` Column4 interface{} `json:"column_4"` Limit int32 `json:"limit"` } type ListAllChannelMessagesRow struct { ID int64 `json:"id"` ChannelID int32 `json:"channel_id"` PacketHash []byte `json:"packet_hash"` SenderName *string `json:"sender_name"` SenderPubkey []byte `json:"sender_pubkey"` Content *string `json:"content"` SentAt pgtype.Timestamptz `json:"sent_at"` PacketHashHex string `json:"packet_hash_hex"` ChannelHash []byte `json:"channel_hash"` ObservationCount int64 `json:"observation_count"` } // Returns all messages across all channels with optional time, IATA, scope and cursor filters. // Pass empty string for iata or scope to skip those filters. // Pass cursor=0 to start from the beginning. func (q *Queries) ListAllChannelMessages(ctx context.Context, arg ListAllChannelMessagesParams) ([]ListAllChannelMessagesRow, error) { rows, err := q.db.Query(ctx, listAllChannelMessages, arg.Column1, arg.Column2, arg.Column3, arg.Column4, arg.Limit, ) if err != nil { return nil, err } defer rows.Close() items := []ListAllChannelMessagesRow{} for rows.Next() { var i ListAllChannelMessagesRow if err := rows.Scan( &i.ID, &i.ChannelID, &i.PacketHash, &i.SenderName, &i.SenderPubkey, &i.Content, &i.SentAt, &i.PacketHashHex, &i.ChannelHash, &i.ObservationCount, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const listChannelMessages = `-- name: ListChannelMessages :many SELECT DISTINCT ON (cm.id) cm.id, cm.channel_id, cm.packet_hash, cm.sender_name, cm.sender_pubkey, cm.content, cm.sent_at, encode(cm.packet_hash, 'hex') as packet_hash_hex, c.channel_hash, (SELECT COUNT(*) FROM packet_observations po2 WHERE po2.packet_hash = cm.packet_hash) AS observation_count FROM channel_messages cm JOIN channels c ON c.id = cm.channel_id JOIN packet_observations po ON po.packet_hash = cm.packet_hash JOIN packets p ON p.packet_hash = cm.packet_hash LEFT JOIN transport_scopes ts ON ts.id = p.scope_id WHERE cm.channel_id = $1 AND ($2::timestamptz IS NULL OR cm.sent_at >= $2) AND (COALESCE(cardinality($3::bpchar[]), 0) = 0 OR po.iata = ANY($3::bpchar[])) AND ($4::text = '' OR ts.name = $4::text) AND ($5::bigint = 0 OR cm.id < $5::bigint) ORDER BY cm.id DESC LIMIT $6 ` type ListChannelMessagesParams struct { ChannelID int32 `json:"channel_id"` Column2 pgtype.Timestamptz `json:"column_2"` Column3 []string `json:"column_3"` Column4 string `json:"column_4"` Column5 int64 `json:"column_5"` Limit int32 `json:"limit"` } type ListChannelMessagesRow struct { ID int64 `json:"id"` ChannelID int32 `json:"channel_id"` PacketHash []byte `json:"packet_hash"` SenderName *string `json:"sender_name"` SenderPubkey []byte `json:"sender_pubkey"` Content *string `json:"content"` SentAt pgtype.Timestamptz `json:"sent_at"` PacketHashHex string `json:"packet_hash_hex"` ChannelHash []byte `json:"channel_hash"` ObservationCount int64 `json:"observation_count"` } // Returns messages for a channel identified by integer ID. // Pass a zero/null timestamp for since to return all messages up to limit. // Pass empty string for iata to skip IATA filtering. // Pass cursor=0 to start from the beginning. func (q *Queries) ListChannelMessages(ctx context.Context, arg ListChannelMessagesParams) ([]ListChannelMessagesRow, error) { rows, err := q.db.Query(ctx, listChannelMessages, arg.ChannelID, arg.Column2, arg.Column3, arg.Column4, arg.Column5, arg.Limit, ) if err != nil { return nil, err } defer rows.Close() items := []ListChannelMessagesRow{} for rows.Next() { var i ListChannelMessagesRow if err := rows.Scan( &i.ID, &i.ChannelID, &i.PacketHash, &i.SenderName, &i.SenderPubkey, &i.Content, &i.SentAt, &i.PacketHashHex, &i.ChannelHash, &i.ObservationCount, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const listChannelMessagesByHash = `-- name: ListChannelMessagesByHash :many SELECT DISTINCT ON (cm.id) cm.id, cm.channel_id, cm.packet_hash, cm.sender_name, cm.sender_pubkey, cm.content, cm.sent_at, c.channel_hash, (SELECT COUNT(*) FROM packet_observations po2 WHERE po2.packet_hash = cm.packet_hash) AS observation_count FROM channel_messages cm JOIN channels c ON c.id = cm.channel_id JOIN packet_observations po ON po.packet_hash = cm.packet_hash JOIN packets p ON p.packet_hash = cm.packet_hash LEFT JOIN transport_scopes ts ON ts.id = p.scope_id WHERE c.channel_hash = $1 AND ($2::timestamptz IS NULL OR cm.sent_at >= $2) AND (COALESCE(cardinality($3::bpchar[]), 0) = 0 OR po.iata = ANY($3::bpchar[])) AND ($4::text = '' OR ts.name = $4::text) AND ($5::bigint = 0 OR cm.id < $5::bigint) ORDER BY cm.id DESC LIMIT $6 ` type ListChannelMessagesByHashParams struct { ChannelHash []byte `json:"channel_hash"` Column2 pgtype.Timestamptz `json:"column_2"` Column3 []string `json:"column_3"` Column4 string `json:"column_4"` Column5 int64 `json:"column_5"` Limit int32 `json:"limit"` } type ListChannelMessagesByHashRow struct { ID int64 `json:"id"` ChannelID int32 `json:"channel_id"` PacketHash []byte `json:"packet_hash"` SenderName *string `json:"sender_name"` SenderPubkey []byte `json:"sender_pubkey"` Content *string `json:"content"` SentAt pgtype.Timestamptz `json:"sent_at"` ChannelHash []byte `json:"channel_hash"` ObservationCount int64 `json:"observation_count"` } // Returns messages for all channels matching a hash byte. // May return messages from multiple channels if the hash collides across different keys. // Pass empty string for iata or scope to skip those filters. // Pass cursor=0 to start from the beginning. func (q *Queries) ListChannelMessagesByHash(ctx context.Context, arg ListChannelMessagesByHashParams) ([]ListChannelMessagesByHashRow, error) { rows, err := q.db.Query(ctx, listChannelMessagesByHash, arg.ChannelHash, arg.Column2, arg.Column3, arg.Column4, arg.Column5, arg.Limit, ) if err != nil { return nil, err } defer rows.Close() items := []ListChannelMessagesByHashRow{} for rows.Next() { var i ListChannelMessagesByHashRow if err := rows.Scan( &i.ID, &i.ChannelID, &i.PacketHash, &i.SenderName, &i.SenderPubkey, &i.Content, &i.SentAt, &i.ChannelHash, &i.ObservationCount, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const listChannels = `-- name: ListChannels :many 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 (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 LIMIT $4 ` type ListChannelsParams struct { ChannelHash []byte `json:"channel_hash"` Iatas []string `json:"iatas"` CursorTs pgtype.Timestamptz `json:"cursor_ts"` PageLimit int32 `json:"page_limit"` } // 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.ChannelHash, arg.Iatas, arg.CursorTs, arg.PageLimit, ) if err != nil { return nil, err } defer rows.Close() items := []Channel{} for rows.Next() { var i Channel if err := rows.Scan( &i.ID, &i.ChannelHash, &i.KeyFingerprint, &i.Name, &i.Hashtag, &i.IsHashtag, &i.IsPublic, &i.KeyKnown, &i.FirstSeen, &i.LastSeen, &i.MessageCount, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const listIATAs = `-- name: ListIATAs :many SELECT iata, display_name, approx_lat, approx_lng, added_at, border FROM iata_codes ORDER BY iata ` func (q *Queries) ListIATAs(ctx context.Context) ([]IataCode, error) { rows, err := q.db.Query(ctx, listIATAs) if err != nil { return nil, err } defer rows.Close() items := []IataCode{} for rows.Next() { var i IataCode if err := rows.Scan( &i.Iata, &i.DisplayName, &i.ApproxLat, &i.ApproxLng, &i.AddedAt, &i.Border, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const listKnownRoutes = `-- name: ListKnownRoutes :many SELECT id, node_ids, hash_prefix, iata, hop_count, first_seen, last_seen, observation_count FROM known_routes WHERE ($1 = '' OR iata = $1) AND ($2 = 0 OR hop_count = $2) AND ($3::timestamptz IS NULL OR last_seen < $3) ORDER BY last_seen DESC LIMIT $4 ` type ListKnownRoutesParams struct { Column1 interface{} `json:"column_1"` Column2 interface{} `json:"column_2"` Column3 pgtype.Timestamptz `json:"column_3"` Limit int32 `json:"limit"` } type ListKnownRoutesRow struct { ID int64 `json:"id"` NodeIds []uuid.UUID `json:"node_ids"` HashPrefix [][]byte `json:"hash_prefix"` Iata string `json:"iata"` HopCount int32 `json:"hop_count"` FirstSeen pgtype.Timestamptz `json:"first_seen"` LastSeen pgtype.Timestamptz `json:"last_seen"` ObservationCount int64 `json:"observation_count"` } func (q *Queries) ListKnownRoutes(ctx context.Context, arg ListKnownRoutesParams) ([]ListKnownRoutesRow, error) { rows, err := q.db.Query(ctx, listKnownRoutes, arg.Column1, arg.Column2, arg.Column3, arg.Limit, ) if err != nil { return nil, err } defer rows.Close() items := []ListKnownRoutesRow{} for rows.Next() { var i ListKnownRoutesRow if err := rows.Scan( &i.ID, &i.NodeIds, &i.HashPrefix, &i.Iata, &i.HopCount, &i.FirstSeen, &i.LastSeen, &i.ObservationCount, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const listMessagesAfterID = `-- name: ListMessagesAfterID :many SELECT DISTINCT ON (cm.id) cm.id, cm.channel_id, cm.packet_hash, cm.sender_name, cm.sender_pubkey, cm.content, cm.sent_at, encode(cm.packet_hash, 'hex') as packet_hash_hex, c.channel_hash, (SELECT COUNT(*) FROM packet_observations po2 WHERE po2.packet_hash = cm.packet_hash) AS observation_count FROM channel_messages cm JOIN channels c ON c.id = cm.channel_id JOIN packet_observations po ON po.packet_hash = cm.packet_hash JOIN packets p ON p.packet_hash = cm.packet_hash LEFT JOIN transport_scopes ts ON ts.id = p.scope_id WHERE cm.id > $1 AND (COALESCE(cardinality($2::bpchar[]), 0) = 0 OR po.iata = ANY($2::bpchar[])) AND ($3::text = '' OR ts.name = $3::text) ORDER BY cm.id ASC LIMIT $4 ` type ListMessagesAfterIDParams struct { ID int64 `json:"id"` Column2 []string `json:"column_2"` Column3 string `json:"column_3"` Limit int32 `json:"limit"` } type ListMessagesAfterIDRow struct { ID int64 `json:"id"` ChannelID int32 `json:"channel_id"` PacketHash []byte `json:"packet_hash"` SenderName *string `json:"sender_name"` SenderPubkey []byte `json:"sender_pubkey"` Content *string `json:"content"` SentAt pgtype.Timestamptz `json:"sent_at"` PacketHashHex string `json:"packet_hash_hex"` ChannelHash []byte `json:"channel_hash"` ObservationCount int64 `json:"observation_count"` } // Returns messages after the given message ID, ordered oldest first. // Used for WS reconnect backfill. func (q *Queries) ListMessagesAfterID(ctx context.Context, arg ListMessagesAfterIDParams) ([]ListMessagesAfterIDRow, error) { rows, err := q.db.Query(ctx, listMessagesAfterID, arg.ID, arg.Column2, arg.Column3, arg.Limit, ) if err != nil { return nil, err } defer rows.Close() items := []ListMessagesAfterIDRow{} for rows.Next() { var i ListMessagesAfterIDRow if err := rows.Scan( &i.ID, &i.ChannelID, &i.PacketHash, &i.SenderName, &i.SenderPubkey, &i.Content, &i.SentAt, &i.PacketHashHex, &i.ChannelHash, &i.ObservationCount, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const listNodeObservations = `-- name: ListNodeObservations :many SELECT po.id, encode(po.packet_hash, 'hex') AS packet_hash_hex, p.payload_type, po.iata, po.heard_at, po.rssi, po.snr, po.hop_count FROM packet_observations po JOIN packets p ON p.packet_hash = po.packet_hash JOIN nodes n ON n.public_key = p.origin_pubkey WHERE n.id = $1 AND ($2 = 0 OR po.id < $2) ORDER BY po.id DESC LIMIT $3 ` type ListNodeObservationsParams struct { ID uuid.UUID `json:"id"` Column2 interface{} `json:"column_2"` Limit int32 `json:"limit"` } type ListNodeObservationsRow struct { ID int64 `json:"id"` PacketHashHex string `json:"packet_hash_hex"` PayloadType int16 `json:"payload_type"` Iata string `json:"iata"` HeardAt pgtype.Timestamptz `json:"heard_at"` Rssi *int16 `json:"rssi"` Snr *float32 `json:"snr"` HopCount int16 `json:"hop_count"` } func (q *Queries) ListNodeObservations(ctx context.Context, arg ListNodeObservationsParams) ([]ListNodeObservationsRow, error) { rows, err := q.db.Query(ctx, listNodeObservations, arg.ID, arg.Column2, arg.Limit) if err != nil { return nil, err } defer rows.Close() items := []ListNodeObservationsRow{} for rows.Next() { var i ListNodeObservationsRow if err := rows.Scan( &i.ID, &i.PacketHashHex, &i.PayloadType, &i.Iata, &i.HeardAt, &i.Rssi, &i.Snr, &i.HopCount, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const listNodes = `-- name: ListNodes :many SELECT n.id, n.public_key, n.node_type, n.name, n.latitude, n.longitude, n.last_seen, n.radio_freq_mhz, n.radio_sf, n.radio_bw_khz, ts.name AS default_scope_name, json_agg(json_build_object('iata', ni.iata, 'lastHeard', (extract(epoch from ni.last_heard) * 1000)::bigint) ORDER BY ni.last_heard DESC) FILTER (WHERE ni.iata IS NOT NULL) AS iatas, EXISTS (SELECT 1 FROM observers o WHERE o.public_key = n.public_key) AS is_observer, (SELECT o.id FROM observers o WHERE o.public_key = n.public_key LIMIT 1) AS observer_id, (SELECT COUNT(DISTINCT nn.neighbor_id) FROM node_neighbors nn WHERE nn.node_id = n.id)::bigint AS known_neighbor_count, -- CASE short-circuits: the array_agg subquery only runs when $10 is true, -- so requests that don't ask for neighbor IDs don't pay for it. (CASE WHEN $10::bool THEN (SELECT COALESCE(array_agg(DISTINCT nn.neighbor_id), '{}'::uuid[]) FROM node_neighbors nn WHERE nn.node_id = n.id) ELSE NULL END)::uuid[] AS neighbor_ids FROM nodes n LEFT JOIN node_iatas ni ON ni.node_id = n.id LEFT JOIN transport_scopes ts ON ts.id = n.default_scope_id WHERE ($1 = 0 OR n.node_type = $1) AND (COALESCE(cardinality($2::bpchar[]), 0) = 0 OR n.id IN (SELECT node_id FROM node_iatas WHERE iata = ANY($2::bpchar[]))) AND ( $3::text = 'any' OR ($3::text = 'true' AND n.supports_multibyte_paths = TRUE) OR ($3::text = 'false' AND n.supports_multibyte_paths = FALSE) ) AND ( $4::text = 'any' OR ($4::text = 'true' AND n.supports_multibyte_traces = TRUE) OR ($4::text = 'false' AND n.supports_multibyte_traces = FALSE) ) AND ($5::bytea IS NULL OR n.public_key = $5) AND ($6 = '' OR n.name ILIKE '%' || $6 || '%') AND ($7::timestamptz IS NULL OR n.last_seen < $7) AND ($9::text = '' OR ts.name = $9::text) AND ($11::text = '' OR encode(n.public_key, 'hex') ILIKE $11 || '%') GROUP BY n.id, ts.name ORDER BY n.last_seen DESC LIMIT $8 ` type ListNodesParams struct { Column1 interface{} `json:"column_1"` Column2 []string `json:"column_2"` Column3 string `json:"column_3"` Column4 string `json:"column_4"` Column5 []byte `json:"column_5"` Column6 interface{} `json:"column_6"` Column7 pgtype.Timestamptz `json:"column_7"` Limit int32 `json:"limit"` Column9 string `json:"column_9"` Column10 bool `json:"column_10"` Column11 string `json:"column_11"` } type ListNodesRow struct { ID uuid.UUID `json:"id"` PublicKey []byte `json:"public_key"` NodeType int16 `json:"node_type"` Name *string `json:"name"` Latitude *float64 `json:"latitude"` Longitude *float64 `json:"longitude"` LastSeen pgtype.Timestamptz `json:"last_seen"` RadioFreqMhz *float32 `json:"radio_freq_mhz"` RadioSf *int16 `json:"radio_sf"` RadioBwKhz *float32 `json:"radio_bw_khz"` DefaultScopeName *string `json:"default_scope_name"` Iatas []byte `json:"iatas"` IsObserver bool `json:"is_observer"` ObserverID uuid.UUID `json:"observer_id"` KnownNeighborCount int64 `json:"known_neighbor_count"` NeighborIds []uuid.UUID `json:"neighbor_ids"` } func (q *Queries) ListNodes(ctx context.Context, arg ListNodesParams) ([]ListNodesRow, error) { rows, err := q.db.Query(ctx, listNodes, arg.Column1, arg.Column2, arg.Column3, arg.Column4, arg.Column5, arg.Column6, arg.Column7, arg.Limit, arg.Column9, arg.Column10, arg.Column11, ) if err != nil { return nil, err } defer rows.Close() items := []ListNodesRow{} for rows.Next() { var i ListNodesRow if err := rows.Scan( &i.ID, &i.PublicKey, &i.NodeType, &i.Name, &i.Latitude, &i.Longitude, &i.LastSeen, &i.RadioFreqMhz, &i.RadioSf, &i.RadioBwKhz, &i.DefaultScopeName, &i.Iatas, &i.IsObserver, &i.ObserverID, &i.KnownNeighborCount, &i.NeighborIds, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } 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, 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 ORDER BY po.heard_at ASC ` type ListObservationsForPacketRow struct { ID int64 `json:"id"` PacketHash []byte `json:"packet_hash"` ObserverID uuid.UUID `json:"observer_id"` Iata string `json:"iata"` HeardAt pgtype.Timestamptz `json:"heard_at"` PathLengthByte int16 `json:"path_length_byte"` HashSize int16 `json:"hash_size"` HopCount int16 `json:"hop_count"` PathBytes []byte `json:"path_bytes"` Rssi *int16 `json:"rssi"` Snr *float32 `json:"snr"` PropagationTimeMs *int32 `json:"propagation_time_ms"` RadioFreqMhz *float32 `json:"radio_freq_mhz"` SpreadFactor *int16 `json:"spread_factor"` 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"` } func (q *Queries) ListObservationsForPacket(ctx context.Context, packetHash []byte) ([]ListObservationsForPacketRow, error) { rows, err := q.db.Query(ctx, listObservationsForPacket, packetHash) if err != nil { return nil, err } defer rows.Close() items := []ListObservationsForPacketRow{} for rows.Next() { var i ListObservationsForPacketRow if err := rows.Scan( &i.ID, &i.PacketHash, &i.ObserverID, &i.Iata, &i.HeardAt, &i.PathLengthByte, &i.HashSize, &i.HopCount, &i.PathBytes, &i.Rssi, &i.Snr, &i.PropagationTimeMs, &i.RadioFreqMhz, &i.SpreadFactor, &i.BandwidthKhz, &i.CodingRate, &i.SourceBroker, &i.PayloadType, &i.ObserverName, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const listObserverAdverts = `-- name: ListObserverAdverts :many SELECT po.id, encode(po.packet_hash, 'hex') AS packet_hash_hex, p.payload_type, po.iata, po.heard_at, po.rssi, po.snr, po.hop_count, n.name AS node_name, encode(p.origin_pubkey, 'hex') AS node_public_key FROM packet_observations po JOIN packets p ON p.packet_hash = po.packet_hash LEFT JOIN nodes n ON n.public_key = p.origin_pubkey WHERE po.observer_id = $1 AND p.payload_type = 4 AND ($2 = 0 OR po.id > $2) ORDER BY po.id ASC LIMIT $3 ` type ListObserverAdvertsParams struct { ObserverID uuid.UUID `json:"observer_id"` Column2 interface{} `json:"column_2"` Limit int32 `json:"limit"` } type ListObserverAdvertsRow struct { ID int64 `json:"id"` PacketHashHex string `json:"packet_hash_hex"` PayloadType int16 `json:"payload_type"` Iata string `json:"iata"` HeardAt pgtype.Timestamptz `json:"heard_at"` Rssi *int16 `json:"rssi"` Snr *float32 `json:"snr"` HopCount int16 `json:"hop_count"` NodeName *string `json:"node_name"` NodePublicKey string `json:"node_public_key"` } // Returns advert packets (payload_type=4) heard by a specific observer. // Pass cursor=0 to start from the beginning, or the last seen id for pagination. func (q *Queries) ListObserverAdverts(ctx context.Context, arg ListObserverAdvertsParams) ([]ListObserverAdvertsRow, error) { rows, err := q.db.Query(ctx, listObserverAdverts, arg.ObserverID, arg.Column2, arg.Limit) if err != nil { return nil, err } defer rows.Close() items := []ListObserverAdvertsRow{} for rows.Next() { var i ListObserverAdvertsRow if err := rows.Scan( &i.ID, &i.PacketHashHex, &i.PayloadType, &i.Iata, &i.HeardAt, &i.Rssi, &i.Snr, &i.HopCount, &i.NodeName, &i.NodePublicKey, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const listObservers = `-- name: ListObservers :many SELECT o.id, o.display_name, o.observer_type, o.last_status_at, o.radio_freq_mhz, o.radio_sf, o.radio_bw_khz, array_remove(array_agg(DISTINCT ts.name ORDER BY ts.name), NULL)::text[] AS scopes, COALESCE(CASE WHEN GREATEST(COALESCE(o.last_status_at, o.last_seen), o.last_seen) > NOW() - INTERVAL '5 minutes' THEN 'online' ELSE 'offline' END, 'offline')::text AS status, COALESCE(( SELECT po.iata FROM packet_observations po WHERE po.observer_id = o.id ORDER BY po.heard_at DESC LIMIT 1 ), '')::text AS iata FROM observers o LEFT JOIN observer_brokers ob ON ob.observer_id = o.id LEFT JOIN observer_scopes os ON os.observer_id = o.id LEFT JOIN transport_scopes ts ON ts.id = os.scope_id WHERE (COALESCE(cardinality($1::bpchar[]), 0) = 0 OR ( SELECT po.iata FROM packet_observations po WHERE po.observer_id = o.id ORDER BY po.heard_at DESC LIMIT 1 ) = ANY($1::bpchar[])) AND ($2 = '' OR o.observer_type = $2) AND ($3 = '' OR ob.broker_name = $3) AND ($4 = '' OR CASE WHEN GREATEST(COALESCE(o.last_status_at, o.last_seen), o.last_seen) > NOW() - INTERVAL '5 minutes' THEN 'online' ELSE 'offline' END = $4) AND ($5 = '' OR o.display_name ILIKE '%' || $5 || '%') AND ($6::timestamptz IS NULL OR o.last_seen < $6) AND ($8::text = '' OR EXISTS ( SELECT 1 FROM observer_scopes os2 JOIN transport_scopes ts2 ON ts2.id = os2.scope_id WHERE os2.observer_id = o.id AND ts2.name = $8::text )) GROUP BY o.id ORDER BY o.last_seen DESC LIMIT $7 ` type ListObserversParams struct { Column1 []string `json:"column_1"` Column2 interface{} `json:"column_2"` Column3 interface{} `json:"column_3"` Column4 interface{} `json:"column_4"` Column5 interface{} `json:"column_5"` Column6 pgtype.Timestamptz `json:"column_6"` Limit int32 `json:"limit"` Column8 string `json:"column_8"` } type ListObserversRow struct { ID uuid.UUID `json:"id"` DisplayName *string `json:"display_name"` ObserverType *string `json:"observer_type"` LastStatusAt pgtype.Timestamptz `json:"last_status_at"` RadioFreqMhz *float32 `json:"radio_freq_mhz"` RadioSf *int16 `json:"radio_sf"` RadioBwKhz *float32 `json:"radio_bw_khz"` Scopes []string `json:"scopes"` Status string `json:"status"` Iata string `json:"iata"` } // Pass cursor=0 to start from the beginning, or the last seen observer's rownum for pagination. // Note: observers use UUID PKs so we order by last_seen and use a keyset on last_seen+id. func (q *Queries) ListObservers(ctx context.Context, arg ListObserversParams) ([]ListObserversRow, error) { rows, err := q.db.Query(ctx, listObservers, arg.Column1, arg.Column2, arg.Column3, arg.Column4, arg.Column5, arg.Column6, arg.Limit, arg.Column8, ) if err != nil { return nil, err } defer rows.Close() items := []ListObserversRow{} for rows.Next() { var i ListObserversRow if err := rows.Scan( &i.ID, &i.DisplayName, &i.ObserverType, &i.LastStatusAt, &i.RadioFreqMhz, &i.RadioSf, &i.RadioBwKhz, &i.Scopes, &i.Status, &i.Iata, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const listPackets = `-- name: ListPackets :many SELECT p.packet_hash, p.payload_type, p.route_type, p.first_heard_at, p.last_heard_at, p.scope_id, ts.name AS scope_name, (SELECT COUNT(*) FROM packet_observations po2 WHERE po2.packet_hash = p.packet_hash) AS observation_count, po.observer_id AS latest_observer_id, o.display_name AS latest_observer_name, po.iata AS latest_observer_iata, po.path_length_byte AS latest_observer_path_length_byte, po.hash_size AS latest_observer_hash_size, po.hop_count AS latest_observer_hop_count, po.path_bytes AS latest_observer_path_bytes FROM packets p LEFT JOIN LATERAL ( SELECT observer_id, iata, path_length_byte, hash_size, hop_count, path_bytes FROM packet_observations WHERE packet_hash = p.packet_hash ORDER BY heard_at DESC LIMIT 1 ) po ON true LEFT JOIN observers o ON o.id = po.observer_id LEFT JOIN transport_scopes ts ON ts.id = p.scope_id WHERE (COALESCE(cardinality($1::smallint[]), 0) = 0 OR p.payload_type = ANY($1::smallint[])) AND (COALESCE(cardinality($2::smallint[]), 0) = 0 OR p.route_type = ANY($2::smallint[])) 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 (COALESCE(cardinality($7::text[]), 0) = 0 OR ts.name = ANY($7::text[])) ORDER BY p.last_heard_at DESC LIMIT $6 ` type ListPacketsParams struct { Column1 []int16 `json:"column_1"` Column2 []int16 `json:"column_2"` Column3 pgtype.Timestamptz `json:"column_3"` Column4 pgtype.Timestamptz `json:"column_4"` Column5 pgtype.Timestamptz `json:"column_5"` Limit int32 `json:"limit"` Column7 []string `json:"column_7"` } type ListPacketsRow struct { PacketHash []byte `json:"packet_hash"` PayloadType int16 `json:"payload_type"` RouteType int16 `json:"route_type"` FirstHeardAt pgtype.Timestamptz `json:"first_heard_at"` LastHeardAt pgtype.Timestamptz `json:"last_heard_at"` ScopeID *int32 `json:"scope_id"` ScopeName *string `json:"scope_name"` ObservationCount int64 `json:"observation_count"` LatestObserverID uuid.UUID `json:"latest_observer_id"` LatestObserverName *string `json:"latest_observer_name"` LatestObserverIata string `json:"latest_observer_iata"` LatestObserverPathLengthByte int16 `json:"latest_observer_path_length_byte"` LatestObserverHashSize int16 `json:"latest_observer_hash_size"` LatestObserverHopCount int16 `json:"latest_observer_hop_count"` LatestObserverPathBytes []byte `json:"latest_observer_path_bytes"` } // Returns packets with the latest observation rolled in for display. // Pass cursor=0 to start from the beginning. IATA-filtered requests are // served by ListPacketsByIATAs instead. func (q *Queries) ListPackets(ctx context.Context, arg ListPacketsParams) ([]ListPacketsRow, error) { rows, err := q.db.Query(ctx, listPackets, arg.Column1, arg.Column2, arg.Column3, arg.Column4, arg.Column5, arg.Limit, arg.Column7, ) if err != nil { return nil, err } defer rows.Close() items := []ListPacketsRow{} for rows.Next() { var i ListPacketsRow if err := rows.Scan( &i.PacketHash, &i.PayloadType, &i.RouteType, &i.FirstHeardAt, &i.LastHeardAt, &i.ScopeID, &i.ScopeName, &i.ObservationCount, &i.LatestObserverID, &i.LatestObserverName, &i.LatestObserverIata, &i.LatestObserverPathLengthByte, &i.LatestObserverHashSize, &i.LatestObserverHopCount, &i.LatestObserverPathBytes, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const listPacketsAfterID = `-- name: ListPacketsAfterID :many SELECT p.packet_hash, p.payload_type, p.route_type, p.first_heard_at, p.last_heard_at, (SELECT COUNT(*) FROM packet_observations po2 WHERE po2.packet_hash = p.packet_hash) AS observation_count, po.observer_id AS latest_observer_id, o.display_name AS latest_observer_name, po.iata AS latest_observer_iata, po.path_length_byte AS latest_observer_path_length_byte, po.hash_size AS latest_observer_hash_size, po.hop_count AS latest_observer_hop_count, po.path_bytes AS latest_observer_path_bytes, ts.name AS scope_name FROM packets p JOIN packet_observations po ON po.packet_hash = p.packet_hash LEFT JOIN observers o ON o.id = po.observer_id LEFT JOIN transport_scopes ts ON ts.id = p.scope_id WHERE po.id > $1 AND ($2::smallint = -1 OR p.payload_type = $2::smallint) AND ($3::smallint = -1 OR p.route_type = $3::smallint) AND (COALESCE(cardinality($4::bpchar[]), 0) = 0 OR po.iata = ANY($4::bpchar[])) AND ($5::text = '' OR ts.name = $5::text) ORDER BY po.id ASC LIMIT $6 ` type ListPacketsAfterIDParams struct { ID int64 `json:"id"` Column2 int16 `json:"column_2"` Column3 int16 `json:"column_3"` Column4 []string `json:"column_4"` Column5 string `json:"column_5"` Limit int32 `json:"limit"` } type ListPacketsAfterIDRow struct { PacketHash []byte `json:"packet_hash"` PayloadType int16 `json:"payload_type"` RouteType int16 `json:"route_type"` FirstHeardAt pgtype.Timestamptz `json:"first_heard_at"` LastHeardAt pgtype.Timestamptz `json:"last_heard_at"` ObservationCount int64 `json:"observation_count"` LatestObserverID uuid.UUID `json:"latest_observer_id"` LatestObserverName *string `json:"latest_observer_name"` LatestObserverIata string `json:"latest_observer_iata"` LatestObserverPathLengthByte int16 `json:"latest_observer_path_length_byte"` LatestObserverHashSize int16 `json:"latest_observer_hash_size"` LatestObserverHopCount int16 `json:"latest_observer_hop_count"` LatestObserverPathBytes []byte `json:"latest_observer_path_bytes"` ScopeName *string `json:"scope_name"` } // Returns packets with observations after the given observation ID, ordered oldest first. // Used for WS reconnect backfill. Pass afterObservationId=0 to start from the beginning. func (q *Queries) ListPacketsAfterID(ctx context.Context, arg ListPacketsAfterIDParams) ([]ListPacketsAfterIDRow, error) { rows, err := q.db.Query(ctx, listPacketsAfterID, arg.ID, arg.Column2, arg.Column3, arg.Column4, arg.Column5, arg.Limit, ) if err != nil { return nil, err } defer rows.Close() items := []ListPacketsAfterIDRow{} for rows.Next() { var i ListPacketsAfterIDRow if err := rows.Scan( &i.PacketHash, &i.PayloadType, &i.RouteType, &i.FirstHeardAt, &i.LastHeardAt, &i.ObservationCount, &i.LatestObserverID, &i.LatestObserverName, &i.LatestObserverIata, &i.LatestObserverPathLengthByte, &i.LatestObserverHashSize, &i.LatestObserverHopCount, &i.LatestObserverPathBytes, &i.ScopeName, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const listPacketsByIATAs = `-- name: ListPacketsByIATAs :many SELECT p.packet_hash, p.payload_type, p.route_type, p.first_heard_at, p.last_heard_at, p.scope_id, ts.name AS scope_name, sh.site_heard_at, (SELECT COUNT(*) FROM packet_observations po2 WHERE po2.packet_hash = p.packet_hash) AS observation_count, po.observer_id AS latest_observer_id, o.display_name AS latest_observer_name, po.iata AS latest_observer_iata, po.path_length_byte AS latest_observer_path_length_byte, po.hash_size AS latest_observer_hash_size, po.hop_count AS latest_observer_hop_count, po.path_bytes AS latest_observer_path_bytes FROM ( SELECT hits.packet_hash, MAX(hits.heard_at)::timestamptz AS site_heard_at FROM unnest($1::bpchar[]) AS req(iata) CROSS JOIN LATERAL ( SELECT po3.packet_hash, po3.heard_at FROM packet_observations po3 JOIN packets p2 ON p2.packet_hash = po3.packet_hash WHERE po3.iata = req.iata AND po3.heard_at < COALESCE($2::timestamptz, 'infinity'::timestamptz) AND (COALESCE(cardinality($3::smallint[]), 0) = 0 OR p2.payload_type = ANY($3::smallint[])) AND (COALESCE(cardinality($4::smallint[]), 0) = 0 OR p2.route_type = ANY($4::smallint[])) AND ($5::timestamptz IS NULL OR p2.first_heard_at >= $5) AND ($6::timestamptz IS NULL OR p2.first_heard_at <= $6) AND (COALESCE(cardinality($7::text[]), 0) = 0 OR EXISTS ( SELECT 1 FROM transport_scopes ts2 WHERE ts2.id = p2.scope_id AND ts2.name = ANY($7::text[]))) ORDER BY po3.heard_at DESC LIMIT $8 ) hits GROUP BY hits.packet_hash HAVING ($2::timestamptz IS NULL OR NOT EXISTS ( SELECT 1 FROM packet_observations px WHERE px.packet_hash = hits.packet_hash AND px.iata = ANY($1::bpchar[]) AND px.heard_at >= $2)) ORDER BY site_heard_at DESC LIMIT $9 ) sh JOIN packets p ON p.packet_hash = sh.packet_hash LEFT JOIN LATERAL ( SELECT observer_id, iata, path_length_byte, hash_size, hop_count, path_bytes FROM packet_observations WHERE packet_hash = p.packet_hash ORDER BY heard_at DESC LIMIT 1 ) po ON true LEFT JOIN observers o ON o.id = po.observer_id LEFT JOIN transport_scopes ts ON ts.id = p.scope_id ORDER BY sh.site_heard_at DESC ` type ListPacketsByIATAsParams struct { Iatas []string `json:"iatas"` CursorTs pgtype.Timestamptz `json:"cursor_ts"` PayloadTypes []int16 `json:"payload_types"` RouteTypes []int16 `json:"route_types"` SinceTs pgtype.Timestamptz `json:"since_ts"` UntilTs pgtype.Timestamptz `json:"until_ts"` ScopeNames []string `json:"scope_names"` ScanDepth int32 `json:"scan_depth"` PageLimit int32 `json:"page_limit"` } type ListPacketsByIATAsRow struct { PacketHash []byte `json:"packet_hash"` PayloadType int16 `json:"payload_type"` RouteType int16 `json:"route_type"` FirstHeardAt pgtype.Timestamptz `json:"first_heard_at"` LastHeardAt pgtype.Timestamptz `json:"last_heard_at"` ScopeID *int32 `json:"scope_id"` ScopeName *string `json:"scope_name"` SiteHeardAt pgtype.Timestamptz `json:"site_heard_at"` ObservationCount int64 `json:"observation_count"` LatestObserverID uuid.UUID `json:"latest_observer_id"` LatestObserverName *string `json:"latest_observer_name"` LatestObserverIata string `json:"latest_observer_iata"` LatestObserverPathLengthByte int16 `json:"latest_observer_path_length_byte"` LatestObserverHashSize int16 `json:"latest_observer_hash_size"` LatestObserverHopCount int16 `json:"latest_observer_hop_count"` LatestObserverPathBytes []byte `json:"latest_observer_path_bytes"` } // IATA-filtered packet list, driven from idx_observations_iata_heard. // Walking packets newest-first and probing for the site probes ~589k packets // to fill a page for a quiet site; walking the site's own observation log is // proportional to the page size instead. Results are ordered by when the // requested sites heard the packet (site-local recency) and the cursor // follows that ordering. scan_depth is a multiple of the page size to // absorb per-observer duplicates; if duplication exceeds it across a // page, pagination ends early (hasMore=false) rather than returning a // short page, even though deeper matches exist. func (q *Queries) ListPacketsByIATAs(ctx context.Context, arg ListPacketsByIATAsParams) ([]ListPacketsByIATAsRow, error) { rows, err := q.db.Query(ctx, listPacketsByIATAs, arg.Iatas, arg.CursorTs, arg.PayloadTypes, arg.RouteTypes, arg.SinceTs, arg.UntilTs, arg.ScopeNames, arg.ScanDepth, arg.PageLimit, ) if err != nil { return nil, err } defer rows.Close() items := []ListPacketsByIATAsRow{} for rows.Next() { var i ListPacketsByIATAsRow if err := rows.Scan( &i.PacketHash, &i.PayloadType, &i.RouteType, &i.FirstHeardAt, &i.LastHeardAt, &i.ScopeID, &i.ScopeName, &i.SiteHeardAt, &i.ObservationCount, &i.LatestObserverID, &i.LatestObserverName, &i.LatestObserverIata, &i.LatestObserverPathLengthByte, &i.LatestObserverHashSize, &i.LatestObserverHopCount, &i.LatestObserverPathBytes, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const listRegions = `-- name: ListRegions :many SELECT id, slug, name FROM regions ORDER BY display_order, name ` type ListRegionsRow struct { ID int32 `json:"id"` Slug string `json:"slug"` Name string `json:"name"` } // ============================================================ // REGIONS // ============================================================ func (q *Queries) ListRegions(ctx context.Context) ([]ListRegionsRow, error) { rows, err := q.db.Query(ctx, listRegions) if err != nil { return nil, err } defer rows.Close() items := []ListRegionsRow{} for rows.Next() { var i ListRegionsRow if err := rows.Scan(&i.ID, &i.Slug, &i.Name); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } 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(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 { Column1 []string `json:"column_1"` Column2 string `json:"column_2"` Column3 pgtype.Timestamptz `json:"column_3"` Column4 pgtype.Timestamptz `json:"column_4"` Column5 pgtype.Timestamptz `json:"column_5"` Limit int32 `json:"limit"` Column7 string `json:"column_7"` } type ListTraceTagsRow struct { TraceTag string `json:"trace_tag"` FirstHeardAt pgtype.Timestamptz `json:"first_heard_at"` LastHeardAt pgtype.Timestamptz `json:"last_heard_at"` PacketCount int64 `json:"packet_count"` IataCount int64 `json:"iata_count"` TraceType string `json:"trace_type"` BestPayload []byte `json:"best_payload"` } // ============================================================ // 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, arg.Column2, arg.Column3, arg.Column4, arg.Column5, arg.Limit, arg.Column7, ) if err != nil { return nil, err } defer rows.Close() items := []ListTraceTagsRow{} for rows.Next() { var i ListTraceTagsRow if err := rows.Scan( &i.TraceTag, &i.FirstHeardAt, &i.LastHeardAt, &i.PacketCount, &i.IataCount, &i.TraceType, &i.BestPayload, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const listUndecryptedGroupTextPackets = `-- name: ListUndecryptedGroupTextPackets :many SELECT packet_hash, raw_payload FROM packets WHERE payload_type = 5 AND decrypted IS NOT TRUE ` type ListUndecryptedGroupTextPacketsRow struct { PacketHash []byte `json:"packet_hash"` RawPayload []byte `json:"raw_payload"` } // Returns GRP_TXT packets (payload_type=5) never successfully decrypted. Used at boot to // retry decryption against the current keystore for packets whose channel key was only added // to the config after they'd already been ingested -- see // internal/ingest.BackfillChannelMessages. func (q *Queries) ListUndecryptedGroupTextPackets(ctx context.Context) ([]ListUndecryptedGroupTextPacketsRow, error) { rows, err := q.db.Query(ctx, listUndecryptedGroupTextPackets) if err != nil { return nil, err } defer rows.Close() items := []ListUndecryptedGroupTextPacketsRow{} for rows.Next() { var i ListUndecryptedGroupTextPacketsRow if err := rows.Scan(&i.PacketHash, &i.RawPayload); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const reconfirmNeighbors = `-- name: ReconfirmNeighbors :exec DELETE FROM node_neighbors nn WHERE NOT EXISTS ( SELECT 1 FROM node_short_ids ns WHERE ns.node_id = nn.neighbor_id AND ns.iata = nn.iata ) OR ( SELECT COUNT(*) FROM node_short_ids ns WHERE ns.iata = nn.iata AND ns.prefix_4 = ( SELECT prefix_4 FROM node_short_ids WHERE node_id = nn.neighbor_id AND iata = nn.iata ) ) > 1 ` // Delete node_neighbors where the neighbor has departed from node_short_ids // for that IATA, or where its prefix_4 is now ambiguous. func (q *Queries) ReconfirmNeighbors(ctx context.Context) error { _, err := q.db.Exec(ctx, reconfirmNeighbors) return err } const reconfirmRoutes = `-- name: ReconfirmRoutes :exec WITH batch AS ( SELECT iata, path_key, node_ids, hash_prefix FROM known_routes ORDER BY last_reconfirmed_at LIMIT $1 ), amb AS MATERIALIZED ( SELECT iata, 1 AS len, prefix_1 AS p FROM node_short_ids GROUP BY iata, prefix_1 HAVING COUNT(*) > 1 UNION ALL SELECT iata, 2, prefix_2 FROM node_short_ids GROUP BY iata, prefix_2 HAVING COUNT(*) > 1 UNION ALL SELECT iata, 3, prefix_3 FROM node_short_ids GROUP BY iata, prefix_3 HAVING COUNT(*) > 1 UNION ALL SELECT iata, 4, prefix_4 FROM node_short_ids GROUP BY iata, prefix_4 HAVING COUNT(*) > 1 ), dead AS ( SELECT b.iata, b.path_key FROM batch b WHERE EXISTS ( SELECT 1 FROM unnest(b.node_ids) AS hop_node_id WHERE NOT EXISTS ( SELECT 1 FROM node_short_ids ns WHERE ns.node_id = hop_node_id AND ns.iata = b.iata ) ) UNION SELECT DISTINCT b.iata, b.path_key FROM batch b CROSS JOIN LATERAL unnest(b.hash_prefix) AS hp JOIN amb a ON a.iata = b.iata AND a.len = length(hp) AND a.p = hp ), deleted AS ( DELETE FROM known_routes kr USING dead d WHERE kr.iata = d.iata AND kr.path_key = d.path_key ) UPDATE known_routes kr SET last_reconfirmed_at = NOW() FROM batch b WHERE kr.iata = b.iata AND kr.path_key = b.path_key AND NOT EXISTS ( SELECT 1 FROM dead d WHERE d.iata = b.iata AND d.path_key = b.path_key ) ` // Checks the $1 least-recently-reconfirmed routes: deletes those with a departed // hop node or a hop prefix now matching >1 node in that IATA (length-aware: // 1/2/3/4-byte hop prefixes check prefix_1/2/3/4), and stamps the survivors. func (q *Queries) ReconfirmRoutes(ctx context.Context, limit int32) error { _, err := q.db.Exec(ctx, reconfirmRoutes, limit) return err } const refreshHourlyStats = `-- name: RefreshHourlyStats :exec REFRESH MATERIALIZED VIEW CONCURRENTLY mv_hourly_iata_stats ` func (q *Queries) RefreshHourlyStats(ctx context.Context) error { _, err := q.db.Exec(ctx, refreshHourlyStats) 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 ` func (q *Queries) RefreshRadioPresets(ctx context.Context) error { _, err := q.db.Exec(ctx, refreshRadioPresets) return err } const refreshTopAdvertisers = `-- name: RefreshTopAdvertisers :exec REFRESH MATERIALIZED VIEW CONCURRENTLY mv_top_advertisers_by_iata ` func (q *Queries) RefreshTopAdvertisers(ctx context.Context) error { _, err := q.db.Exec(ctx, refreshTopAdvertisers) return err } const refreshTopNodes = `-- name: RefreshTopNodes :exec REFRESH MATERIALIZED VIEW CONCURRENTLY mv_top_nodes_by_iata ` func (q *Queries) RefreshTopNodes(ctx context.Context) error { _, err := q.db.Exec(ctx, refreshTopNodes) 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 refreshTopTalkers = `-- name: RefreshTopTalkers :exec REFRESH MATERIALIZED VIEW CONCURRENTLY mv_top_talkers_by_iata ` func (q *Queries) RefreshTopTalkers(ctx context.Context) error { _, err := q.db.Exec(ctx, refreshTopTalkers) return err } const resolvePathHashesP1 = `-- name: ResolvePathHashesP1 :many SELECT ns.prefix_4 AS hash, n.id AS node_id, n.name, n.latitude, n.longitude, n.public_key FROM node_short_ids ns JOIN nodes n ON n.id = ns.node_id WHERE ns.iata = $1 AND n.node_type IN (2, 3) AND ns.prefix_1 = ANY($2::bytea[]) ` type ResolvePathHashesP1Params struct { Iata string `json:"iata"` Column2 [][]byte `json:"column_2"` } type ResolvePathHashesP1Row struct { Hash []byte `json:"hash"` NodeID uuid.UUID `json:"node_id"` Name *string `json:"name"` Latitude *float64 `json:"latitude"` Longitude *float64 `json:"longitude"` PublicKey []byte `json:"public_key"` } // ============================================================ // HELPERS // ============================================================ // Path hash resolution is split per prefix width so each query gets a // cacheable generic plan on its (iata, prefix_N) index; a single CASE // predicate forced a fresh custom plan on every call. func (q *Queries) ResolvePathHashesP1(ctx context.Context, arg ResolvePathHashesP1Params) ([]ResolvePathHashesP1Row, error) { rows, err := q.db.Query(ctx, resolvePathHashesP1, arg.Iata, arg.Column2) if err != nil { return nil, err } defer rows.Close() items := []ResolvePathHashesP1Row{} for rows.Next() { var i ResolvePathHashesP1Row if err := rows.Scan( &i.Hash, &i.NodeID, &i.Name, &i.Latitude, &i.Longitude, &i.PublicKey, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const resolvePathHashesP2 = `-- name: ResolvePathHashesP2 :many SELECT ns.prefix_4 AS hash, n.id AS node_id, n.name, n.latitude, n.longitude, n.public_key FROM node_short_ids ns JOIN nodes n ON n.id = ns.node_id WHERE ns.iata = $1 AND n.node_type IN (2, 3) AND ns.prefix_2 = ANY($2::bytea[]) ` type ResolvePathHashesP2Params struct { Iata string `json:"iata"` Column2 [][]byte `json:"column_2"` } type ResolvePathHashesP2Row struct { Hash []byte `json:"hash"` NodeID uuid.UUID `json:"node_id"` Name *string `json:"name"` Latitude *float64 `json:"latitude"` Longitude *float64 `json:"longitude"` PublicKey []byte `json:"public_key"` } func (q *Queries) ResolvePathHashesP2(ctx context.Context, arg ResolvePathHashesP2Params) ([]ResolvePathHashesP2Row, error) { rows, err := q.db.Query(ctx, resolvePathHashesP2, arg.Iata, arg.Column2) if err != nil { return nil, err } defer rows.Close() items := []ResolvePathHashesP2Row{} for rows.Next() { var i ResolvePathHashesP2Row if err := rows.Scan( &i.Hash, &i.NodeID, &i.Name, &i.Latitude, &i.Longitude, &i.PublicKey, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const resolvePathHashesP3 = `-- name: ResolvePathHashesP3 :many SELECT ns.prefix_4 AS hash, n.id AS node_id, n.name, n.latitude, n.longitude, n.public_key FROM node_short_ids ns JOIN nodes n ON n.id = ns.node_id WHERE ns.iata = $1 AND n.node_type IN (2, 3) AND ns.prefix_3 = ANY($2::bytea[]) ` type ResolvePathHashesP3Params struct { Iata string `json:"iata"` Column2 [][]byte `json:"column_2"` } type ResolvePathHashesP3Row struct { Hash []byte `json:"hash"` NodeID uuid.UUID `json:"node_id"` Name *string `json:"name"` Latitude *float64 `json:"latitude"` Longitude *float64 `json:"longitude"` PublicKey []byte `json:"public_key"` } func (q *Queries) ResolvePathHashesP3(ctx context.Context, arg ResolvePathHashesP3Params) ([]ResolvePathHashesP3Row, error) { rows, err := q.db.Query(ctx, resolvePathHashesP3, arg.Iata, arg.Column2) if err != nil { return nil, err } defer rows.Close() items := []ResolvePathHashesP3Row{} for rows.Next() { var i ResolvePathHashesP3Row if err := rows.Scan( &i.Hash, &i.NodeID, &i.Name, &i.Latitude, &i.Longitude, &i.PublicKey, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const resolvePathHashesP4 = `-- name: ResolvePathHashesP4 :many SELECT ns.prefix_4 AS hash, n.id AS node_id, n.name, n.latitude, n.longitude, n.public_key FROM node_short_ids ns JOIN nodes n ON n.id = ns.node_id WHERE ns.iata = $1 AND n.node_type IN (2, 3) AND ns.prefix_4 = ANY($2::bytea[]) ` type ResolvePathHashesP4Params struct { Iata string `json:"iata"` Column2 [][]byte `json:"column_2"` } type ResolvePathHashesP4Row struct { Hash []byte `json:"hash"` NodeID uuid.UUID `json:"node_id"` Name *string `json:"name"` Latitude *float64 `json:"latitude"` Longitude *float64 `json:"longitude"` PublicKey []byte `json:"public_key"` } func (q *Queries) ResolvePathHashesP4(ctx context.Context, arg ResolvePathHashesP4Params) ([]ResolvePathHashesP4Row, error) { rows, err := q.db.Query(ctx, resolvePathHashesP4, arg.Iata, arg.Column2) if err != nil { return nil, err } defer rows.Close() items := []ResolvePathHashesP4Row{} for rows.Next() { var i ResolvePathHashesP4Row if err := rows.Scan( &i.Hash, &i.NodeID, &i.Name, &i.Latitude, &i.Longitude, &i.PublicKey, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const searchKnownRoutes = `-- name: SearchKnownRoutes :many SELECT id, node_ids, hash_prefix, iata, hop_count, first_seen, last_seen, observation_count FROM known_routes WHERE iata = $1 AND array_position(hash_prefix, $2::bytea) IS NOT NULL AND array_position(hash_prefix, $3::bytea) IS NOT NULL AND array_position(hash_prefix, $2::bytea) < array_position(hash_prefix, $3::bytea) ORDER BY hop_count ASC, last_seen DESC ` type SearchKnownRoutesParams struct { Iata string `json:"iata"` Column2 []byte `json:"column_2"` Column3 []byte `json:"column_3"` } type SearchKnownRoutesRow struct { ID int64 `json:"id"` NodeIds []uuid.UUID `json:"node_ids"` HashPrefix [][]byte `json:"hash_prefix"` Iata string `json:"iata"` HopCount int32 `json:"hop_count"` FirstSeen pgtype.Timestamptz `json:"first_seen"` LastSeen pgtype.Timestamptz `json:"last_seen"` ObservationCount int64 `json:"observation_count"` } // Returns known routes containing a subsequence from source to destination hash prefix. // Verifies source appears before destination in the route. func (q *Queries) SearchKnownRoutes(ctx context.Context, arg SearchKnownRoutesParams) ([]SearchKnownRoutesRow, error) { rows, err := q.db.Query(ctx, searchKnownRoutes, arg.Iata, arg.Column2, arg.Column3) if err != nil { return nil, err } defer rows.Close() items := []SearchKnownRoutesRow{} for rows.Next() { var i SearchKnownRoutesRow if err := rows.Scan( &i.ID, &i.NodeIds, &i.HashPrefix, &i.Iata, &i.HopCount, &i.FirstSeen, &i.LastSeen, &i.ObservationCount, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const setNodeDefaultScope = `-- name: SetNodeDefaultScope :exec UPDATE nodes SET default_scope_id = $2 WHERE id = $1 ` type SetNodeDefaultScopeParams struct { ID uuid.UUID `json:"id"` DefaultScopeID *int32 `json:"default_scope_id"` } func (q *Queries) SetNodeDefaultScope(ctx context.Context, arg SetNodeDefaultScopeParams) error { _, err := q.db.Exec(ctx, setNodeDefaultScope, arg.ID, arg.DefaultScopeID) return err } const setNodeMultibytePaths = `-- name: SetNodeMultibytePaths :exec UPDATE nodes SET supports_multibyte_paths = TRUE WHERE id = $1 AND supports_multibyte_paths = FALSE ` func (q *Queries) SetNodeMultibytePaths(ctx context.Context, id uuid.UUID) error { _, err := q.db.Exec(ctx, setNodeMultibytePaths, id) return err } const setNodeMultibyteTraces = `-- name: SetNodeMultibyteTraces :exec UPDATE nodes SET supports_multibyte_traces = TRUE WHERE id = $1 AND supports_multibyte_traces = FALSE ` func (q *Queries) SetNodeMultibyteTraces(ctx context.Context, id uuid.UUID) error { _, err := q.db.Exec(ctx, setNodeMultibyteTraces, id) return err } const setPacketDecrypted = `-- name: SetPacketDecrypted :exec UPDATE packets SET decrypted = true WHERE packet_hash = $1 ` func (q *Queries) SetPacketDecrypted(ctx context.Context, packetHash []byte) error { _, err := q.db.Exec(ctx, setPacketDecrypted, packetHash) return err } const touchObserverBrokers = `-- name: TouchObserverBrokers :exec UPDATE observer_brokers ob SET last_seen = GREATEST(ob.last_seen, v.seen), last_packet_at = GREATEST(ob.last_packet_at, v.seen) FROM ( SELECT unnest($1::uuid[]) AS observer_id, unnest($2::text[]) AS broker_name, unnest($3::timestamptz[]) AS seen ) v WHERE ob.observer_id = v.observer_id AND ob.broker_name = v.broker_name ` type TouchObserverBrokersParams struct { Column1 []uuid.UUID `json:"column_1"` Column2 []string `json:"column_2"` Column3 []pgtype.Timestamptz `json:"column_3"` } func (q *Queries) TouchObserverBrokers(ctx context.Context, arg TouchObserverBrokersParams) error { _, err := q.db.Exec(ctx, touchObserverBrokers, arg.Column1, arg.Column2, arg.Column3) return err } const touchObservers = `-- name: TouchObservers :exec UPDATE observers o SET last_seen = GREATEST(o.last_seen, v.seen), observation_count = COALESCE(o.observation_count, 0) + v.delta FROM ( SELECT unnest($1::uuid[]) AS id, unnest($2::timestamptz[]) AS seen, unnest($3::int[]) AS delta ) v WHERE o.id = v.id ` type TouchObserversParams struct { Column1 []uuid.UUID `json:"column_1"` Column2 []pgtype.Timestamptz `json:"column_2"` Column3 []int32 `json:"column_3"` } // Batched flush of coalesced presence bumps. GREATEST keeps a late flush // from regressing a newer write-through (e.g. a status update). func (q *Queries) TouchObservers(ctx context.Context, arg TouchObserversParams) error { _, err := q.db.Exec(ctx, touchObservers, arg.Column1, arg.Column2, arg.Column3) return err } const touchPackets = `-- name: TouchPackets :exec UPDATE packets p SET last_heard_at = GREATEST(p.last_heard_at, v.heard) FROM ( SELECT unnest($1::bytea[]) AS packet_hash, unnest($2::timestamptz[]) AS heard ) v WHERE p.packet_hash = v.packet_hash ` type TouchPacketsParams struct { Column1 [][]byte `json:"column_1"` Column2 []pgtype.Timestamptz `json:"column_2"` } func (q *Queries) TouchPackets(ctx context.Context, arg TouchPacketsParams) error { _, err := q.db.Exec(ctx, touchPackets, arg.Column1, arg.Column2) return err } const updateObserverRegionScope = `-- name: UpdateObserverRegionScope :exec UPDATE observers SET region_scope = $2 WHERE id = $1 ` type UpdateObserverRegionScopeParams struct { ID uuid.UUID `json:"id"` RegionScope *string `json:"region_scope"` } // Records the observer's own OTA-reported region scope, from the "self" // field of a /neighbors report. Always known (not queried OTA), so this // unconditionally overwrites, unlike the neighbor-side region_scope. func (q *Queries) UpdateObserverRegionScope(ctx context.Context, arg UpdateObserverRegionScopeParams) error { _, err := q.db.Exec(ctx, updateObserverRegionScope, arg.ID, arg.RegionScope) return err } const updateObserverStatus = `-- name: UpdateObserverStatus :one UPDATE observers SET display_name = COALESCE(NULLIF($2, ''), display_name), observer_type = COALESCE(NULLIF($3, ''), observer_type), software_version = COALESCE($4, software_version), hardware_model = COALESCE($5, hardware_model), firmware_version = COALESCE($6, firmware_version), firmware_build = COALESCE($7, firmware_build), radio_freq_mhz = COALESCE($8, radio_freq_mhz), radio_sf = COALESCE($9, radio_sf), radio_bw_khz = COALESCE($10, radio_bw_khz), radio_cr = COALESCE($11, radio_cr), battery_level = COALESCE($12, battery_level), uptime_seconds = COALESCE($13, uptime_seconds), status_metadata = $14, last_status_at = NOW(), last_seen = NOW() WHERE public_key = $1 RETURNING id ` type UpdateObserverStatusParams struct { PublicKey []byte `json:"public_key"` Column2 interface{} `json:"column_2"` Column3 interface{} `json:"column_3"` SoftwareVersion *string `json:"software_version"` HardwareModel *string `json:"hardware_model"` FirmwareVersion *string `json:"firmware_version"` FirmwareBuild *string `json:"firmware_build"` RadioFreqMhz *float32 `json:"radio_freq_mhz"` RadioSf *int16 `json:"radio_sf"` RadioBwKhz *float32 `json:"radio_bw_khz"` RadioCr *int16 `json:"radio_cr"` BatteryLevel *float32 `json:"battery_level"` UptimeSeconds *int64 `json:"uptime_seconds"` StatusMetadata []byte `json:"status_metadata"` } func (q *Queries) UpdateObserverStatus(ctx context.Context, arg UpdateObserverStatusParams) (uuid.UUID, error) { row := q.db.QueryRow(ctx, updateObserverStatus, arg.PublicKey, arg.Column2, arg.Column3, arg.SoftwareVersion, arg.HardwareModel, arg.FirmwareVersion, arg.FirmwareBuild, arg.RadioFreqMhz, arg.RadioSf, arg.RadioBwKhz, arg.RadioCr, arg.BatteryLevel, arg.UptimeSeconds, arg.StatusMetadata, ) var id uuid.UUID err := row.Scan(&id) return id, err } const upsertChannel = `-- name: UpsertChannel :one INSERT INTO channels (channel_hash, key_fingerprint, name, hashtag, is_hashtag, key_known, last_seen) VALUES ($1, $2::bytea, $3, $4, $5, ($2 IS NOT NULL), NOW()) ON CONFLICT (channel_hash, key_fingerprint) DO UPDATE SET last_seen = NOW(), name = COALESCE(EXCLUDED.name, channels.name), message_count = CASE WHEN $6 THEN channels.message_count + 1 ELSE channels.message_count END RETURNING id, channel_hash, key_fingerprint, name, hashtag, is_hashtag, is_public, key_known, first_seen, last_seen, message_count ` type UpsertChannelParams struct { ChannelHash []byte `json:"channel_hash"` Column2 []byte `json:"column_2"` Name *string `json:"name"` Hashtag *string `json:"hashtag"` IsHashtag *bool `json:"is_hashtag"` MessageCount *int64 `json:"message_count"` } // ============================================================ // CHANNELS // ============================================================ // Upsert a channel by (hash, key_fingerprint). Pass NULL fingerprint for // hash-only records (key unknown). Returns the channel row. func (q *Queries) UpsertChannel(ctx context.Context, arg UpsertChannelParams) (Channel, error) { row := q.db.QueryRow(ctx, upsertChannel, arg.ChannelHash, arg.Column2, arg.Name, arg.Hashtag, arg.IsHashtag, arg.MessageCount, ) var i Channel err := row.Scan( &i.ID, &i.ChannelHash, &i.KeyFingerprint, &i.Name, &i.Hashtag, &i.IsHashtag, &i.IsPublic, &i.KeyKnown, &i.FirstSeen, &i.LastSeen, &i.MessageCount, ) return i, err } const upsertChannelHashOnly = `-- name: UpsertChannelHashOnly :one INSERT INTO channels (channel_hash, last_seen) VALUES ($1, NOW()) ON CONFLICT (channel_hash) WHERE key_fingerprint IS NULL DO UPDATE SET last_seen = NOW() RETURNING id ` func (q *Queries) UpsertChannelHashOnly(ctx context.Context, channelHash []byte) (int32, error) { row := q.db.QueryRow(ctx, upsertChannelHashOnly, channelHash) var id int32 err := row.Scan(&id) 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 = EXCLUDED.last_heard WHERE EXCLUDED.last_heard > channel_iatas.last_heard + INTERVAL '1 hour' ` type UpsertChannelIATAParams struct { ChannelHash []byte `json:"channel_hash"` Iata string `json:"iata"` 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 } const upsertIATA = `-- name: UpsertIATA :exec INSERT INTO iata_codes (iata) VALUES ($1) ON CONFLICT (iata) DO NOTHING ` // Copyright 2026 Beacon Contributors // SPDX-License-Identifier: agpl // ============================================================ // IATA CODES // ============================================================ func (q *Queries) UpsertIATA(ctx context.Context, iata string) error { _, err := q.db.Exec(ctx, upsertIATA, iata) return err } const upsertIATABorder = `-- name: UpsertIATABorder :exec INSERT INTO iata_codes (iata, border) VALUES ($1, $2) ON CONFLICT (iata) DO UPDATE SET border = EXCLUDED.border ` type UpsertIATABorderParams struct { Iata string `json:"iata"` Border []byte `json:"border"` } // Written by the config-file-driven seeder (internal/config/seed.go), not a // runtime HTTP path. border is a full, pre-validated GeoJSON Feature with // bbox already computed -- see internal/config/border.go. func (q *Queries) UpsertIATABorder(ctx context.Context, arg UpsertIATABorderParams) error { _, err := q.db.Exec(ctx, upsertIATABorder, arg.Iata, arg.Border) return err } const upsertIATADetails = `-- name: UpsertIATADetails :exec INSERT INTO iata_codes (iata, display_name, approx_lat, approx_lng) VALUES ($1, $2, $3, $4) ON CONFLICT (iata) DO UPDATE SET display_name = EXCLUDED.display_name, approx_lat = EXCLUDED.approx_lat, approx_lng = EXCLUDED.approx_lng ` type UpsertIATADetailsParams struct { Iata string `json:"iata"` DisplayName *string `json:"display_name"` ApproxLat *float64 `json:"approx_lat"` ApproxLng *float64 `json:"approx_lng"` } func (q *Queries) UpsertIATADetails(ctx context.Context, arg UpsertIATADetailsParams) error { _, err := q.db.Exec(ctx, upsertIATADetails, arg.Iata, arg.DisplayName, arg.ApproxLat, arg.ApproxLng, ) return err } const upsertKnownRoute = `-- name: UpsertKnownRoute :exec INSERT INTO known_routes (path_key, node_ids, hash_prefix, iata, hop_count) VALUES ($1, $2, $3, $4, $5) ON CONFLICT (iata, path_key) DO UPDATE SET last_seen = NOW(), observation_count = known_routes.observation_count + 1 ` type UpsertKnownRouteParams struct { PathKey []byte `json:"path_key"` NodeIds []uuid.UUID `json:"node_ids"` HashPrefix [][]byte `json:"hash_prefix"` Iata string `json:"iata"` HopCount int32 `json:"hop_count"` } // ============================================================ // ROUTES // ============================================================ // Route identity is path_key, an md5 of node_ids computed by the caller. // On conflict, observation_count and last_seen are bumped. func (q *Queries) UpsertKnownRoute(ctx context.Context, arg UpsertKnownRouteParams) error { _, err := q.db.Exec(ctx, upsertKnownRoute, arg.PathKey, arg.NodeIds, arg.HashPrefix, arg.Iata, arg.HopCount, ) return err } const upsertNode = `-- name: UpsertNode :one INSERT INTO nodes (public_key, node_type, name, latitude, longitude, location_source, last_advert_at, last_seen, radio_freq_mhz, radio_sf, radio_bw_khz, device_clock_drift_seconds) VALUES ($1, $2, $3, $4, $5, 'advert', NOW(), NOW(), $6, $7, $8, $9) ON CONFLICT (public_key) DO UPDATE SET node_type = EXCLUDED.node_type, name = COALESCE(EXCLUDED.name, nodes.name), latitude = COALESCE(EXCLUDED.latitude, nodes.latitude), longitude = COALESCE(EXCLUDED.longitude, nodes.longitude), location_source = CASE WHEN EXCLUDED.latitude IS NOT NULL THEN 'advert' ELSE nodes.location_source END, last_advert_at = NOW(), last_seen = NOW(), radio_freq_mhz = EXCLUDED.radio_freq_mhz, radio_sf = EXCLUDED.radio_sf, radio_bw_khz = EXCLUDED.radio_bw_khz, device_clock_drift_seconds = EXCLUDED.device_clock_drift_seconds RETURNING id, public_key, node_type, name, latitude, longitude, location_source, last_advert_at, supports_multibyte_paths, supports_multibyte_traces, default_scope_id, min_firmware_version, first_seen, last_seen, radio_freq_mhz, radio_sf, radio_bw_khz, metadata, device_clock_drift_seconds ` type UpsertNodeParams struct { PublicKey []byte `json:"public_key"` NodeType int16 `json:"node_type"` Name *string `json:"name"` Latitude *float64 `json:"latitude"` Longitude *float64 `json:"longitude"` RadioFreqMhz *float32 `json:"radio_freq_mhz"` RadioSf *int16 `json:"radio_sf"` RadioBwKhz *float32 `json:"radio_bw_khz"` DeviceClockDriftSeconds *int32 `json:"device_clock_drift_seconds"` } // ============================================================ // NODES // ============================================================ func (q *Queries) UpsertNode(ctx context.Context, arg UpsertNodeParams) (Node, error) { row := q.db.QueryRow(ctx, upsertNode, arg.PublicKey, arg.NodeType, arg.Name, arg.Latitude, arg.Longitude, arg.RadioFreqMhz, arg.RadioSf, arg.RadioBwKhz, arg.DeviceClockDriftSeconds, ) var i Node err := row.Scan( &i.ID, &i.PublicKey, &i.NodeType, &i.Name, &i.Latitude, &i.Longitude, &i.LocationSource, &i.LastAdvertAt, &i.SupportsMultibytePaths, &i.SupportsMultibyteTraces, &i.DefaultScopeID, &i.MinFirmwareVersion, &i.FirstSeen, &i.LastSeen, &i.RadioFreqMhz, &i.RadioSf, &i.RadioBwKhz, &i.Metadata, &i.DeviceClockDriftSeconds, ) return i, err } const upsertNodeIATA = `-- name: UpsertNodeIATA :exec INSERT INTO node_iatas (node_id, iata, last_heard, observation_count) VALUES ($1, $2, NOW(), 1) ON CONFLICT (node_id, iata) DO UPDATE SET last_heard = NOW(), observation_count = node_iatas.observation_count + 1 ` type UpsertNodeIATAParams struct { NodeID uuid.UUID `json:"node_id"` Iata string `json:"iata"` } // ============================================================ // NODE IATAS // ============================================================ func (q *Queries) UpsertNodeIATA(ctx context.Context, arg UpsertNodeIATAParams) error { _, err := q.db.Exec(ctx, upsertNodeIATA, arg.NodeID, arg.Iata) return err } const upsertNodeNeighbor = `-- name: UpsertNodeNeighbor :exec INSERT INTO node_neighbors (node_id, neighbor_id, iata, observation_count, snr, region_scope) VALUES ($1, $2, $3, 1, $4, $5) ON CONFLICT (node_id, neighbor_id, iata) DO UPDATE SET last_seen = NOW(), observation_count = node_neighbors.observation_count + 1, snr = COALESCE(EXCLUDED.snr, node_neighbors.snr), region_scope = COALESCE(EXCLUDED.region_scope, node_neighbors.region_scope) ` type UpsertNodeNeighborParams struct { NodeID uuid.UUID `json:"node_id"` NeighborID uuid.UUID `json:"neighbor_id"` Iata string `json:"iata"` Snr *float32 `json:"snr"` RegionScope *string `json:"region_scope"` } // ============================================================ // NEIGHBORS // ============================================================ // Records or updates a neighbor relationship between two nodes observed in the same IATA. // node_id is the advertising node, neighbor_id is the first-hop forwarder. // snr is optional; pass NULL when no signal reading is available (the // common case). regionScope is optional too; pass NULL whenever the OTA // scope query for this neighbor didn't succeed (status != "responded"), // so a failed/timed-out query doesn't erase a previously known scope. // On conflict, snr and region_scope are only overwritten when a new // non-null value is supplied. func (q *Queries) UpsertNodeNeighbor(ctx context.Context, arg UpsertNodeNeighborParams) error { _, err := q.db.Exec(ctx, upsertNodeNeighbor, arg.NodeID, arg.NeighborID, arg.Iata, arg.Snr, arg.RegionScope, ) return err } const upsertNodeShortID = `-- name: UpsertNodeShortID :exec INSERT INTO node_short_ids (node_id, iata, prefix_4) VALUES ($1, $2, $3) ON CONFLICT (node_id, iata) DO NOTHING ` type UpsertNodeShortIDParams struct { NodeID uuid.UUID `json:"node_id"` Iata string `json:"iata"` Prefix4 []byte `json:"prefix_4"` } func (q *Queries) UpsertNodeShortID(ctx context.Context, arg UpsertNodeShortIDParams) error { _, err := q.db.Exec(ctx, upsertNodeShortID, arg.NodeID, arg.Iata, arg.Prefix4) return err } const upsertObserver = `-- name: UpsertObserver :one INSERT INTO observers (public_key, observer_type, last_seen) VALUES ($1, 'unknown', NOW()) ON CONFLICT (public_key) DO UPDATE SET last_seen = NOW(), observation_count = observers.observation_count + 1 RETURNING id, public_key, display_name, observer_type, software_version, hardware_model, firmware_version, firmware_build, radio_freq_mhz, radio_sf, radio_bw_khz, radio_cr, battery_level, uptime_seconds, status_metadata, last_status_at, first_seen, last_seen, observation_count, metadata, region_scope ` // ============================================================ // OBSERVERS // ============================================================ func (q *Queries) UpsertObserver(ctx context.Context, publicKey []byte) (Observer, error) { row := q.db.QueryRow(ctx, upsertObserver, publicKey) var i Observer err := row.Scan( &i.ID, &i.PublicKey, &i.DisplayName, &i.ObserverType, &i.SoftwareVersion, &i.HardwareModel, &i.FirmwareVersion, &i.FirmwareBuild, &i.RadioFreqMhz, &i.RadioSf, &i.RadioBwKhz, &i.RadioCr, &i.BatteryLevel, &i.UptimeSeconds, &i.StatusMetadata, &i.LastStatusAt, &i.FirstSeen, &i.LastSeen, &i.ObservationCount, &i.Metadata, &i.RegionScope, ) return i, err } const upsertObserverBroker = `-- name: UpsertObserverBroker :exec INSERT INTO observer_brokers (observer_id, broker_name, last_seen, last_packet_at) VALUES ($1, $2, NOW(), NOW()) ON CONFLICT (observer_id, broker_name) DO UPDATE SET last_seen = NOW(), last_packet_at = NOW() ` type UpsertObserverBrokerParams struct { ObserverID uuid.UUID `json:"observer_id"` BrokerName string `json:"broker_name"` } // ============================================================ // OBSERVER BROKERS // ============================================================ func (q *Queries) UpsertObserverBroker(ctx context.Context, arg UpsertObserverBrokerParams) error { _, err := q.db.Exec(ctx, upsertObserverBroker, arg.ObserverID, arg.BrokerName) return err } const upsertObserverScope = `-- name: UpsertObserverScope :exec INSERT INTO observer_scopes (observer_id, scope_id, last_seen) VALUES ($1, $2, NOW()) ON CONFLICT (observer_id, scope_id) DO UPDATE SET last_seen = NOW() ` type UpsertObserverScopeParams struct { ObserverID uuid.UUID `json:"observer_id"` ScopeID int32 `json:"scope_id"` } func (q *Queries) UpsertObserverScope(ctx context.Context, arg UpsertObserverScopeParams) error { _, err := q.db.Exec(ctx, upsertObserverScope, arg.ObserverID, arg.ScopeID) return err } const upsertPacket = `-- name: UpsertPacket :one INSERT INTO packets ( packet_hash, payload_type, payload_version, route_type, transport_codes_present, region_code, sub_region_code, origin_pubkey, raw_payload, raw_header, parsed_payload, channel_hash, scope_id, trace_tag, first_heard_at, last_heard_at ) VALUES ( $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, NOW(), NOW() ) ON CONFLICT (packet_hash) DO UPDATE SET last_heard_at = NOW() RETURNING packet_hash, payload_type, payload_version, route_type, transport_codes_present, region_code, sub_region_code, origin_pubkey, raw_payload, raw_header, parsed_payload, decrypted, channel_hash, first_heard_at, last_heard_at, (xmax = 0) AS inserted ` type UpsertPacketParams struct { PacketHash []byte `json:"packet_hash"` PayloadType int16 `json:"payload_type"` PayloadVersion int16 `json:"payload_version"` RouteType int16 `json:"route_type"` TransportCodesPresent *bool `json:"transport_codes_present"` RegionCode *int32 `json:"region_code"` SubRegionCode *int32 `json:"sub_region_code"` OriginPubkey []byte `json:"origin_pubkey"` RawPayload []byte `json:"raw_payload"` RawHeader []byte `json:"raw_header"` ParsedPayload []byte `json:"parsed_payload"` ChannelHash []byte `json:"channel_hash"` ScopeID *int32 `json:"scope_id"` TraceTag []byte `json:"trace_tag"` } type UpsertPacketRow struct { PacketHash []byte `json:"packet_hash"` PayloadType int16 `json:"payload_type"` PayloadVersion int16 `json:"payload_version"` RouteType int16 `json:"route_type"` TransportCodesPresent *bool `json:"transport_codes_present"` RegionCode *int32 `json:"region_code"` SubRegionCode *int32 `json:"sub_region_code"` OriginPubkey []byte `json:"origin_pubkey"` RawPayload []byte `json:"raw_payload"` RawHeader []byte `json:"raw_header"` ParsedPayload []byte `json:"parsed_payload"` Decrypted *bool `json:"decrypted"` ChannelHash []byte `json:"channel_hash"` FirstHeardAt pgtype.Timestamptz `json:"first_heard_at"` LastHeardAt pgtype.Timestamptz `json:"last_heard_at"` Inserted bool `json:"inserted"` } // ============================================================ // PACKETS // ============================================================ func (q *Queries) UpsertPacket(ctx context.Context, arg UpsertPacketParams) (UpsertPacketRow, error) { row := q.db.QueryRow(ctx, upsertPacket, arg.PacketHash, arg.PayloadType, arg.PayloadVersion, arg.RouteType, arg.TransportCodesPresent, arg.RegionCode, arg.SubRegionCode, arg.OriginPubkey, arg.RawPayload, arg.RawHeader, arg.ParsedPayload, arg.ChannelHash, arg.ScopeID, arg.TraceTag, ) var i UpsertPacketRow err := row.Scan( &i.PacketHash, &i.PayloadType, &i.PayloadVersion, &i.RouteType, &i.TransportCodesPresent, &i.RegionCode, &i.SubRegionCode, &i.OriginPubkey, &i.RawPayload, &i.RawHeader, &i.ParsedPayload, &i.Decrypted, &i.ChannelHash, &i.FirstHeardAt, &i.LastHeardAt, &i.Inserted, ) return i, err } const upsertRegion = `-- name: UpsertRegion :one INSERT INTO regions (slug, name, description, display_order, center_lat, center_lng, zoom_level, updated_at) VALUES ($1, $2, $3, $4, $5, $6, $7, NOW()) ON CONFLICT (slug) DO UPDATE SET name = EXCLUDED.name, description = EXCLUDED.description, display_order = EXCLUDED.display_order, center_lat = EXCLUDED.center_lat, center_lng = EXCLUDED.center_lng, zoom_level = EXCLUDED.zoom_level, updated_at = NOW() RETURNING id ` type UpsertRegionParams struct { Slug string `json:"slug"` Name string `json:"name"` Description *string `json:"description"` DisplayOrder *int32 `json:"display_order"` CenterLat *float64 `json:"center_lat"` CenterLng *float64 `json:"center_lng"` ZoomLevel *int32 `json:"zoom_level"` } func (q *Queries) UpsertRegion(ctx context.Context, arg UpsertRegionParams) (int32, error) { row := q.db.QueryRow(ctx, upsertRegion, arg.Slug, arg.Name, arg.Description, arg.DisplayOrder, arg.CenterLat, arg.CenterLng, arg.ZoomLevel, ) var id int32 err := row.Scan(&id) return id, err } const upsertRegionIATA = `-- name: UpsertRegionIATA :exec INSERT INTO region_iatas (region_id, iata) VALUES ($1, $2) ON CONFLICT (region_id, iata) DO NOTHING ` type UpsertRegionIATAParams struct { RegionID int32 `json:"region_id"` Iata string `json:"iata"` } func (q *Queries) UpsertRegionIATA(ctx context.Context, arg UpsertRegionIATAParams) error { _, err := q.db.Exec(ctx, upsertRegionIATA, arg.RegionID, arg.Iata) 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 = EXCLUDED.last_heard WHERE EXCLUDED.last_heard > trace_iatas.last_heard + INTERVAL '1 hour' ` type UpsertTraceIATAParams struct { TraceTag []byte `json:"trace_tag"` Iata string `json:"iata"` 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 } const upsertTransportScope = `-- name: UpsertTransportScope :exec INSERT INTO transport_scopes (name, display_name, transport_key, key_fingerprint) VALUES ($1, $2, $3, $4) ON CONFLICT (name) DO UPDATE SET display_name = EXCLUDED.display_name, transport_key = EXCLUDED.transport_key, key_fingerprint = EXCLUDED.key_fingerprint ` type UpsertTransportScopeParams struct { Name string `json:"name"` DisplayName *string `json:"display_name"` TransportKey []byte `json:"transport_key"` KeyFingerprint []byte `json:"key_fingerprint"` } // ============================================================ // TRANSPORT CODES // ============================================================ func (q *Queries) UpsertTransportScope(ctx context.Context, arg UpsertTransportScopeParams) error { _, err := q.db.Exec(ctx, upsertTransportScope, arg.Name, arg.DisplayName, arg.TransportKey, arg.KeyFingerprint, ) return err }