Files
beacon-server/db/sqlc/queries.sql.go
T
gadgethd 87239bd4b8 fix(seed): make UpsertIATADetails a true upsert
UpsertIATADetails was a plain UPDATE that silently affected zero rows
when the iata_codes row did not yet exist.  Adding a new IATA to the
config for the first time resulted in display_name, approx_lat, and
approx_lng remaining NULL.

Change the query to INSERT ... ON CONFLICT DO UPDATE so the row is
atomically created (if missing) or updated (if present), matching the
function's name and the pattern used by UpsertTransportScope.
2026-06-12 08:25:14 -07:00

3486 lines
103 KiB
Go

// 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 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 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 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
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"`
}
// 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,
); 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 ($1::text = '' OR iata = ANY(string_to_array($1::text, ',')))
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 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,
)
return i, 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"`
}
func (q *Queries) GetKnownRoutesByNode(ctx context.Context, arg GetKnownRoutesByNodeParams) ([]KnownRoute, error) {
rows, err := q.db.Query(ctx, getKnownRoutesByNode, arg.Iata, arg.Column2)
if err != nil {
return nil, err
}
defer rows.Close()
items := []KnownRoute{}
for rows.Next() {
var i KnownRoute
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, 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(*) 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"`
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.DefaultScopeName,
&i.IsObserver,
&i.ObserverID,
&i.Iatas,
&i.KnownNeighborCount,
)
return i, 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
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"`
}
// 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,
); 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 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,
)
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 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,
)
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 ($2::text = '' OR iata = ANY(string_to_array($2::text, ',')))
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,
COUNT(DISTINCT p.packet_hash) AS packet_count,
COUNT(DISTINCT os.observer_id) AS observer_count,
COUNT(DISTINCT n.id) AS node_count
FROM transport_scopes ts
LEFT JOIN packets p ON p.scope_id = ts.id
LEFT JOIN observer_scopes os ON os.scope_id = ts.id
LEFT JOIN nodes n ON n.default_scope_id = ts.id
GROUP BY ts.name
ORDER BY ts.name
`
type GetScopeStatsRow struct {
Name string `json:"name"`
PacketCount int64 `json:"packet_count"`
ObserverCount int64 `json:"observer_count"`
NodeCount int64 `json:"node_count"`
}
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 ($1::text = '' OR po.iata = ANY(string_to_array($1::text, ',')))
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 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 ($1::text = '' OR ni.iata = ANY(string_to_array($1::text, ',')))
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 ($1::text = '' OR po.iata = ANY(string_to_array($1::text, ',')))
`
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
p.payload_type,
COUNT(*) AS count
FROM packet_observations po
JOIN packets p ON p.packet_hash = po.packet_hash
WHERE po.heard_at > $1
AND ($2::text = '' OR po.iata = ANY(string_to_array($2::text, ',')))
GROUP BY p.payload_type
ORDER BY count DESC
`
type GetStatsPayloadBreakdownParams struct {
HeardAt pgtype.Timestamptz `json:"heard_at"`
Column2 string `json:"column_2"`
}
type GetStatsPayloadBreakdownRow struct {
PayloadType int16 `json:"payload_type"`
Count int64 `json:"count"`
}
// Returns observation counts grouped by payload type for the given window and IATA.
func (q *Queries) GetStatsPayloadBreakdown(ctx context.Context, arg GetStatsPayloadBreakdownParams) ([]GetStatsPayloadBreakdownRow, error) {
rows, err := q.db.Query(ctx, getStatsPayloadBreakdown, arg.HeardAt, arg.Column2)
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 getStatsTopObservers = `-- name: GetStatsTopObservers :many
SELECT
o.id,
o.display_name,
o.observer_type,
COUNT(*) AS observation_count,
COALESCE((
SELECT po2.iata FROM packet_observations po2
WHERE po2.observer_id = o.id
ORDER BY po2.heard_at DESC LIMIT 1
), '') AS iata
FROM packet_observations po
JOIN observers o ON o.id = po.observer_id
WHERE po.heard_at > $1
AND ($2::text = '' OR po.iata = ANY(string_to_array($2::text, ',')))
GROUP BY o.id
ORDER BY observation_count DESC
LIMIT $3
`
type GetStatsTopObserversParams struct {
HeardAt pgtype.Timestamptz `json:"heard_at"`
Column2 string `json:"column_2"`
Limit int32 `json:"limit"`
}
type GetStatsTopObserversRow struct {
ID uuid.UUID `json:"id"`
DisplayName *string `json:"display_name"`
ObserverType *string `json:"observer_type"`
ObservationCount int64 `json:"observation_count"`
Iata interface{} `json:"iata"`
}
// Returns the top N observers by observation count for the given window and IATA.
func (q *Queries) GetStatsTopObservers(ctx context.Context, arg GetStatsTopObserversParams) ([]GetStatsTopObserversRow, error) {
rows, err := q.db.Query(ctx, getStatsTopObservers, arg.HeardAt, 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 getTopNodes = `-- name: GetTopNodes :many
SELECT iata, node_id, name, node_type, observation_count, last_heard FROM mv_top_nodes_by_iata
WHERE ($1::text = '' OR iata = ANY(string_to_array($1::text, ',')))
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
) VALUES (
$1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16
)
ON CONFLICT (packet_hash, observer_id, heard_at) 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
`
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"`
}
// ============================================================
// 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,
)
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,
)
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 ($2::text = '' OR po.iata = ANY(string_to_array($2::text, ',')))
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 ($3::text = '' OR po.iata = ANY(string_to_array($3::text, ',')))
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 ($3::text = '' OR po.iata = ANY(string_to_array($3::text, ',')))
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 DISTINCT c.id, c.channel_hash, c.key_fingerprint, c.name, c.hashtag, c.is_hashtag, c.is_public, c.key_known, c.first_seen, c.last_seen, c.message_count FROM channels c
WHERE ($1::bytea IS NULL OR c.channel_hash = $1)
AND ($2 = '' OR EXISTS (
SELECT 1 FROM packets p
JOIN packet_observations po ON po.packet_hash = p.packet_hash
WHERE p.channel_hash = c.channel_hash
AND po.iata ILIKE $2
))
AND ($3::timestamptz IS NULL OR c.last_seen < $3)
ORDER BY c.last_seen DESC
LIMIT $4
`
type ListChannelsParams struct {
Column1 []byte `json:"column_1"`
Column2 interface{} `json:"column_2"`
Column3 pgtype.Timestamptz `json:"column_3"`
Limit int32 `json:"limit"`
}
// Returns channels ordered by last seen, optionally filtered by hash and/or IATA.
// Pass NULL for hash to skip hash filtering. Pass empty string for iata to skip IATA filtering.
// IATA filter returns channels that have active packets in that IATA (case-insensitive).
// Pass cursor=0 to start from the beginning (cursor is last_seen epoch ms).
func (q *Queries) ListChannels(ctx context.Context, arg ListChannelsParams) ([]Channel, error) {
rows, err := q.db.Query(ctx, listChannels,
arg.Column1,
arg.Column2,
arg.Column3,
arg.Limit,
)
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 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,
); 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"`
}
func (q *Queries) ListKnownRoutes(ctx context.Context, arg ListKnownRoutesParams) ([]KnownRoute, 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 := []KnownRoute{}
for rows.Next() {
var i KnownRoute
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 ($2::text = '' OR po.iata = ANY(string_to_array($2::text, ',')))
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(*) FROM node_neighbors nn WHERE nn.node_id = n.id)::bigint AS known_neighbor_count
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 ($2::text = '' OR n.id IN (SELECT node_id FROM node_iatas WHERE iata = ANY(string_to_array($2::text, ','))))
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)
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"`
}
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"`
}
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,
)
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,
); 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, 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"`
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.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 o.last_status_at > 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
($1::text = '' OR (
SELECT po.iata FROM packet_observations po
WHERE po.observer_id = o.id
ORDER BY po.heard_at DESC LIMIT 1
) = ANY(string_to_array($1::text, ',')))
AND ($2 = '' OR o.observer_type = $2)
AND ($3 = '' OR ob.broker_name = $3)
AND ($4 = '' OR CASE
WHEN o.last_status_at > 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
FROM packets p
LEFT JOIN LATERAL (
SELECT observer_id, iata
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
($1::smallint = -1 OR p.payload_type = $1::smallint)
AND ($2::smallint = -1 OR p.route_type = $2::smallint)
AND ($3::text = '' OR EXISTS (
SELECT 1 FROM packet_observations po3
WHERE po3.packet_hash = p.packet_hash
AND po3.iata = ANY(string_to_array($3::text, ','))
))
AND ($4::timestamptz IS NULL OR p.first_heard_at >= $4)
AND ($5::timestamptz IS NULL OR p.first_heard_at <= $5)
AND ($6::timestamptz IS NULL OR p.last_heard_at < $6)
AND ($8::text = '' OR ts.name = $8::text)
ORDER BY p.last_heard_at DESC
LIMIT $7
`
type ListPacketsParams struct {
Column1 int16 `json:"column_1"`
Column2 int16 `json:"column_2"`
Column3 string `json:"column_3"`
Column4 pgtype.Timestamptz `json:"column_4"`
Column5 pgtype.Timestamptz `json:"column_5"`
Column6 pgtype.Timestamptz `json:"column_6"`
Limit int32 `json:"limit"`
Column8 string `json:"column_8"`
}
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"`
}
// Returns packets with the latest observation rolled in for display.
// Pass cursor=0 to start from the beginning.
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.Column6,
arg.Limit,
arg.Column8,
)
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,
); 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,
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 ($4::text = '' OR po.iata = ANY(string_to_array($4::text, ',')))
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"`
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.ScopeName,
); 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
SELECT
encode(p.trace_tag, 'hex') AS trace_tag,
MIN(p.first_heard_at)::timestamptz AS first_heard_at,
MAX(p.last_heard_at)::timestamptz AS last_heard_at,
COUNT(DISTINCT p.packet_hash) AS packet_count,
COUNT(DISTINCT po.iata) AS iata_count,
MAX(p.parsed_payload->>'type')::text AS trace_type,
best.parsed_payload AS best_payload
FROM packets p
LEFT JOIN packet_observations po ON po.packet_hash = p.packet_hash
LEFT JOIN LATERAL (
SELECT parsed_payload
FROM packets p2
WHERE p2.trace_tag = p.trace_tag
ORDER BY jsonb_array_length(p2.parsed_payload->'pathHashes') DESC
LIMIT 1
) best ON true
WHERE p.trace_tag IS NOT NULL
AND ($1::text = '' OR po.iata = ANY(string_to_array($1::text, ',')))
AND ($2::text = '' OR p.scope_id = (SELECT id FROM transport_scopes WHERE name = $2))
AND ($3::timestamptz IS NULL OR p.first_heard_at >= $3)
AND ($4::timestamptz IS NULL OR p.first_heard_at <= $4)
AND ($5::timestamptz IS NULL OR p.last_heard_at < $5)
AND ($7::text = '' OR p.parsed_payload->>'type' = $7)
GROUP BY p.trace_tag, best.parsed_payload
ORDER BY MAX(p.last_heard_at) DESC
LIMIT $6
`
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.
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 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 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 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 resolvePathHashes = `-- name: ResolvePathHashes :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 CASE
WHEN cardinality($2::bytea[]) > 0 AND length($2[1]) = 1 THEN ns.prefix_1 = ANY($2)
WHEN cardinality($2::bytea[]) > 0 AND length($2[1]) = 2 THEN ns.prefix_2 = ANY($2)
WHEN cardinality($2::bytea[]) > 0 AND length($2[1]) = 3 THEN ns.prefix_3 = ANY($2)
WHEN cardinality($2::bytea[]) > 0 AND length($2[1]) = 4 THEN ns.prefix_4 = ANY($2)
ELSE FALSE
END
`
type ResolvePathHashesParams struct {
Iata string `json:"iata"`
Column2 [][]byte `json:"column_2"`
}
type ResolvePathHashesRow 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
// ============================================================
func (q *Queries) ResolvePathHashes(ctx context.Context, arg ResolvePathHashesParams) ([]ResolvePathHashesRow, error) {
rows, err := q.db.Query(ctx, resolvePathHashes, arg.Iata, arg.Column2)
if err != nil {
return nil, err
}
defer rows.Close()
items := []ResolvePathHashesRow{}
for rows.Next() {
var i ResolvePathHashesRow
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"`
}
// 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) ([]KnownRoute, error) {
rows, err := q.db.Query(ctx, searchKnownRoutes, arg.Iata, arg.Column2, arg.Column3)
if err != nil {
return nil, err
}
defer rows.Close()
items := []KnownRoute{}
for rows.Next() {
var i KnownRoute
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 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 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 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 (node_ids, hash_prefix, iata, hop_count)
VALUES ($1, $2, $3, $4)
ON CONFLICT (node_ids, iata) DO UPDATE SET
last_seen = NOW(),
observation_count = known_routes.observation_count + 1
`
type UpsertKnownRouteParams struct {
NodeIds []uuid.UUID `json:"node_ids"`
HashPrefix [][]byte `json:"hash_prefix"`
Iata string `json:"iata"`
HopCount int32 `json:"hop_count"`
}
// ============================================================
// ROUTES
// ============================================================
// Inserts or updates a known route (all hops resolved to high confidence).
// node_ids and hash_prefix are ordered arrays of the resolved node UUIDs and
// their hash bytes. last_seen is bumped on conflict.
func (q *Queries) UpsertKnownRoute(ctx context.Context, arg UpsertKnownRouteParams) error {
_, err := q.db.Exec(ctx, upsertKnownRoute,
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)
VALUES ($1, $2, $3, $4, $5, 'advert', NOW(), NOW(), $6, $7, $8)
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
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
`
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"`
}
// ============================================================
// 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,
)
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,
)
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)
VALUES ($1, $2, $3, 1)
ON CONFLICT (node_id, neighbor_id, iata) DO UPDATE SET
last_seen = NOW(),
observation_count = node_neighbors.observation_count + 1
`
type UpsertNodeNeighborParams struct {
NodeID uuid.UUID `json:"node_id"`
NeighborID uuid.UUID `json:"neighbor_id"`
Iata string `json:"iata"`
}
// ============================================================
// 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.
func (q *Queries) UpsertNodeNeighbor(ctx context.Context, arg UpsertNodeNeighborParams) error {
_, err := q.db.Exec(ctx, upsertNodeNeighbor, arg.NodeID, arg.NeighborID, arg.Iata)
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
`
// ============================================================
// 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,
)
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 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
}