// 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 getChannelByHashAndFingerprint = `-- name: GetChannelByHashAndFingerprint :one SELECT id, channel_hash, key_fingerprint, name, hashtag, is_hashtag, is_public, key_known, first_seen, last_seen, message_count FROM channels WHERE channel_hash = $1 AND key_fingerprint = $2 ` type GetChannelByHashAndFingerprintParams struct { ChannelHash []byte `json:"channel_hash"` KeyFingerprint []byte `json:"key_fingerprint"` } func (q *Queries) GetChannelByHashAndFingerprint(ctx context.Context, arg GetChannelByHashAndFingerprintParams) (Channel, error) { row := q.db.QueryRow(ctx, getChannelByHashAndFingerprint, arg.ChannelHash, arg.KeyFingerprint) 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 getChannelByHashtag = `-- name: GetChannelByHashtag :one SELECT id, channel_hash, key_fingerprint, name, hashtag, is_hashtag, is_public, key_known, first_seen, last_seen, message_count FROM channels WHERE hashtag = $1 ` func (q *Queries) GetChannelByHashtag(ctx context.Context, hashtag *string) (Channel, error) { row := q.db.QueryRow(ctx, getChannelByHashtag, hashtag) 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 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 getChannelsByHash = `-- name: GetChannelsByHash :many SELECT id, channel_hash, key_fingerprint, name, hashtag, is_hashtag, is_public, key_known, first_seen, last_seen, message_count FROM channels WHERE channel_hash = $1 ORDER BY last_seen DESC LIMIT $2 ` type GetChannelsByHashParams struct { ChannelHash []byte `json:"channel_hash"` Limit int32 `json:"limit"` } // Returns all channels for a given hash (may be multiple on hash collision). func (q *Queries) GetChannelsByHash(ctx context.Context, arg GetChannelsByHashParams) ([]Channel, error) { rows, err := q.db.Query(ctx, getChannelsByHash, arg.ChannelHash, 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 getHourlyStats = `-- name: GetHourlyStats :many SELECT iata, hour, observation_count, unique_packets, active_observers FROM mv_hourly_iata_stats WHERE ($1::char(3) IS NULL OR iata = $1) 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 getNodeByPubkey = `-- name: GetNodeByPubkey :one SELECT id, public_key, node_type, name, latitude, longitude, location_source, last_advert_at, supports_multibyte_paths, supports_multibyte_traces, min_firmware_version, first_seen, last_seen, metadata FROM nodes WHERE public_key = $1 ` func (q *Queries) GetNodeByPubkey(ctx context.Context, publicKey []byte) (Node, error) { row := q.db.QueryRow(ctx, getNodeByPubkey, publicKey) 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.MinFirmwareVersion, &i.FirstSeen, &i.LastSeen, &i.Metadata, ) return i, err } 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 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 getPacket = `-- name: GetPacket :one SELECT packet_hash, payload_type, payload_version, route_type, transport_codes_present, region_code, sub_region_code, origin_pubkey, raw_payload, parsed_payload, decrypted, channel_hash, first_heard_at, last_heard_at, observation_count FROM packets WHERE packet_hash = $1 ` func (q *Queries) GetPacket(ctx context.Context, packetHash []byte) (Packet, error) { row := q.db.QueryRow(ctx, getPacket, packetHash) var i Packet err := row.Scan( &i.PacketHash, &i.PayloadType, &i.PayloadVersion, &i.RouteType, &i.TransportCodesPresent, &i.RegionCode, &i.SubRegionCode, &i.OriginPubkey, &i.RawPayload, &i.ParsedPayload, &i.Decrypted, &i.ChannelHash, &i.FirstHeardAt, &i.LastHeardAt, &i.ObservationCount, ) return i, err } 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 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 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::char(3) IS NULL OR po.iata = $1) ` 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 getTopNodes = `-- name: GetTopNodes :many SELECT iata, node_id, name, node_type, observation_count, last_heard FROM mv_top_nodes_by_iata WHERE ($1::char(3) IS NULL OR iata = $1) 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 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 FROM channel_messages cm JOIN channels c ON c.id = cm.channel_id JOIN packet_observations po ON po.packet_hash = cm.packet_hash WHERE ($1::timestamptz IS NULL OR cm.sent_at >= $1) AND ($2 = '' OR po.iata = $2) ORDER BY cm.id, cm.sent_at DESC LIMIT $3 ` type ListAllChannelMessagesParams struct { Column1 pgtype.Timestamptz `json:"column_1"` Column2 interface{} `json:"column_2"` 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"` } // Returns all messages across all channels with optional time and IATA filters. // Pass empty string for iata to skip IATA filtering. func (q *Queries) ListAllChannelMessages(ctx context.Context, arg ListAllChannelMessagesParams) ([]ListAllChannelMessagesRow, error) { rows, err := q.db.Query(ctx, listAllChannelMessages, arg.Column1, arg.Column2, 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, ); 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 FROM channel_messages cm JOIN channels c ON c.id = cm.channel_id JOIN packet_observations po ON po.packet_hash = cm.packet_hash WHERE cm.channel_id = $1 AND ($2::timestamptz IS NULL OR cm.sent_at >= $2) AND ($3 = '' OR po.iata = $3) ORDER BY cm.id, cm.sent_at DESC LIMIT $4 ` type ListChannelMessagesParams struct { ChannelID int32 `json:"channel_id"` Column2 pgtype.Timestamptz `json:"column_2"` Column3 interface{} `json:"column_3"` 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"` } // 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. 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.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, ); 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 FROM channel_messages cm JOIN channels c ON c.id = cm.channel_id JOIN packet_observations po ON po.packet_hash = cm.packet_hash WHERE c.channel_hash = $1 AND ($2::timestamptz IS NULL OR cm.sent_at >= $2) AND ($3 = '' OR po.iata = $3) ORDER BY cm.id, cm.sent_at DESC LIMIT $4 ` type ListChannelMessagesByHashParams struct { ChannelHash []byte `json:"channel_hash"` Column2 pgtype.Timestamptz `json:"column_2"` Column3 interface{} `json:"column_3"` 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"` } // 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 to skip IATA filtering. 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.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, ); 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 = $2 )) ORDER BY c.last_seen DESC LIMIT $3 ` type ListChannelsParams struct { Column1 []byte `json:"column_1"` Column2 interface{} `json:"column_2"` 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 been active (have messages heard) in that IATA. func (q *Queries) ListChannels(ctx context.Context, arg ListChannelsParams) ([]Channel, error) { rows, err := q.db.Query(ctx, listChannels, arg.Column1, arg.Column2, 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 listNodes = `-- name: ListNodes :many SELECT id, public_key, node_type, name, latitude, longitude, location_source, last_advert_at, supports_multibyte_paths, supports_multibyte_traces, min_firmware_version, first_seen, last_seen, metadata FROM nodes WHERE ($1::smallint IS NULL OR node_type = $1) ORDER BY last_seen DESC LIMIT $2 ` type ListNodesParams struct { Column1 int16 `json:"column_1"` Limit int32 `json:"limit"` } func (q *Queries) ListNodes(ctx context.Context, arg ListNodesParams) ([]Node, error) { rows, err := q.db.Query(ctx, listNodes, arg.Column1, arg.Limit) if err != nil { return nil, err } defer rows.Close() items := []Node{} for rows.Next() { var i Node if err := rows.Scan( &i.ID, &i.PublicKey, &i.NodeType, &i.Name, &i.Latitude, &i.Longitude, &i.LocationSource, &i.LastAdvertAt, &i.SupportsMultibytePaths, &i.SupportsMultibyteTraces, &i.MinFirmwareVersion, &i.FirstSeen, &i.LastSeen, &i.Metadata, ); err != nil { return nil, err } items = append(items, i) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const listObservationsForObserver = `-- name: ListObservationsForObserver :many SELECT 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 FROM packet_observations WHERE observer_id = $1 AND ($2::timestamptz IS NULL OR heard_at >= $2) ORDER BY heard_at DESC LIMIT $3 ` type ListObservationsForObserverParams struct { ObserverID uuid.UUID `json:"observer_id"` Column2 pgtype.Timestamptz `json:"column_2"` Limit int32 `json:"limit"` } func (q *Queries) ListObservationsForObserver(ctx context.Context, arg ListObservationsForObserverParams) ([]PacketObservation, error) { rows, err := q.db.Query(ctx, listObservationsForObserver, arg.ObserverID, arg.Column2, arg.Limit) if err != nil { return nil, err } defer rows.Close() items := []PacketObservation{} for rows.Next() { var i PacketObservation 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, ); 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 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 FROM packet_observations WHERE packet_hash = $1 ORDER BY heard_at ASC ` func (q *Queries) ListObservationsForPacket(ctx context.Context, packetHash []byte) ([]PacketObservation, error) { rows, err := q.db.Query(ctx, listObservationsForPacket, packetHash) if err != nil { return nil, err } defer rows.Close() items := []PacketObservation{} for rows.Next() { var i PacketObservation 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, ); 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, 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 WHERE ($1 = '' OR ( SELECT po.iata FROM packet_observations po WHERE po.observer_id = o.id ORDER BY po.heard_at DESC LIMIT 1 ) = $1) 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) GROUP BY o.id ORDER BY o.last_seen DESC ` type ListObserversParams struct { Column1 interface{} `json:"column_1"` Column2 interface{} `json:"column_2"` Column3 interface{} `json:"column_3"` Column4 interface{} `json:"column_4"` } 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"` Status string `json:"status"` Iata string `json:"iata"` } 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, ) 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.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.payload_version, p.route_type, p.transport_codes_present, p.region_code, p.sub_region_code, p.origin_pubkey, p.raw_payload, p.parsed_payload, p.decrypted, p.channel_hash, p.first_heard_at, p.last_heard_at, p.observation_count FROM packets p WHERE ($1::smallint IS NULL OR p.payload_type = $1) AND ($2::smallint IS NULL OR p.route_type = $2) AND ($3::timestamptz IS NULL OR p.first_heard_at >= $3) AND ($4::timestamptz IS NULL OR p.first_heard_at <= $4) ORDER BY p.last_heard_at DESC LIMIT $5 ` type ListPacketsParams struct { Column1 int16 `json:"column_1"` Column2 int16 `json:"column_2"` Column3 pgtype.Timestamptz `json:"column_3"` Column4 pgtype.Timestamptz `json:"column_4"` Limit int32 `json:"limit"` } func (q *Queries) ListPackets(ctx context.Context, arg ListPacketsParams) ([]Packet, error) { rows, err := q.db.Query(ctx, listPackets, arg.Column1, arg.Column2, arg.Column3, arg.Column4, arg.Limit, ) if err != nil { return nil, err } defer rows.Close() items := []Packet{} for rows.Next() { var i Packet if err := rows.Scan( &i.PacketHash, &i.PayloadType, &i.PayloadVersion, &i.RouteType, &i.TransportCodesPresent, &i.RegionCode, &i.SubRegionCode, &i.OriginPubkey, &i.RawPayload, &i.ParsedPayload, &i.Decrypted, &i.ChannelHash, &i.FirstHeardAt, &i.LastHeardAt, &i.ObservationCount, ); 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.payload_version, p.route_type, p.transport_codes_present, p.region_code, p.sub_region_code, p.origin_pubkey, p.raw_payload, p.parsed_payload, p.decrypted, p.channel_hash, p.first_heard_at, p.last_heard_at, p.observation_count FROM packets p JOIN packet_observations po ON po.packet_hash = p.packet_hash WHERE po.id > $1 ORDER BY po.id ASC LIMIT $2 ` type ListPacketsAfterIDParams struct { ID int64 `json:"id"` Limit int32 `json:"limit"` } func (q *Queries) ListPacketsAfterID(ctx context.Context, arg ListPacketsAfterIDParams) ([]Packet, error) { rows, err := q.db.Query(ctx, listPacketsAfterID, arg.ID, arg.Limit) if err != nil { return nil, err } defer rows.Close() items := []Packet{} for rows.Next() { var i Packet if err := rows.Scan( &i.PacketHash, &i.PayloadType, &i.PayloadVersion, &i.RouteType, &i.TransportCodesPresent, &i.RegionCode, &i.SubRegionCode, &i.OriginPubkey, &i.RawPayload, &i.ParsedPayload, &i.Decrypted, &i.ChannelHash, &i.FirstHeardAt, &i.LastHeardAt, &i.ObservationCount, ); 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 resolvePathHashes = `-- name: ResolvePathHashes :many SELECT DISTINCT n.id 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"` } // ============================================================ // HELPERS // ============================================================ func (q *Queries) ResolvePathHashes(ctx context.Context, arg ResolvePathHashesParams) ([]uuid.UUID, error) { rows, err := q.db.Query(ctx, resolvePathHashes, arg.Iata, arg.Column2) if err != nil { return nil, err } defer rows.Close() items := []uuid.UUID{} for rows.Next() { var id uuid.UUID if err := rows.Scan(&id); err != nil { return nil, err } items = append(items, id) } if err := rows.Err(); err != nil { return nil, err } return items, nil } const setChannelKeyKnown = `-- name: SetChannelKeyKnown :exec UPDATE channels SET key_known = TRUE WHERE channel_hash = $1 AND key_fingerprint = $2 ` type SetChannelKeyKnownParams struct { ChannelHash []byte `json:"channel_hash"` KeyFingerprint []byte `json:"key_fingerprint"` } func (q *Queries) SetChannelKeyKnown(ctx context.Context, arg SetChannelKeyKnownParams) error { _, err := q.db.Exec(ctx, setChannelKeyKnown, arg.ChannelHash, arg.KeyFingerprint) 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 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 ` // ============================================================ // 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 UPDATE iata_codes SET display_name = $2, approx_lat = $3, approx_lng = $4 WHERE iata = $1 ` 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 upsertNode = `-- name: UpsertNode :one INSERT INTO nodes (public_key, node_type, name, latitude, longitude, location_source, last_advert_at, last_seen) VALUES ($1, $2, $3, $4, $5, 'advert', NOW(), NOW()) 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() RETURNING id, public_key, node_type, name, latitude, longitude, location_source, last_advert_at, supports_multibyte_paths, supports_multibyte_traces, min_firmware_version, first_seen, last_seen, 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"` } // ============================================================ // 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, ) 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.MinFirmwareVersion, &i.FirstSeen, &i.LastSeen, &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 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 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, parsed_payload, channel_hash, first_heard_at, last_heard_at, observation_count ) VALUES ( $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, NOW(), NOW(), 1 ) ON CONFLICT (packet_hash) DO UPDATE SET last_heard_at = NOW(), observation_count = packets.observation_count + 1 RETURNING packet_hash, payload_type, payload_version, route_type, transport_codes_present, region_code, sub_region_code, origin_pubkey, raw_payload, parsed_payload, decrypted, channel_hash, first_heard_at, last_heard_at, observation_count, (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"` ParsedPayload []byte `json:"parsed_payload"` ChannelHash []byte `json:"channel_hash"` } 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"` 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"` ObservationCount *int32 `json:"observation_count"` 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.ParsedPayload, arg.ChannelHash, ) 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.ParsedPayload, &i.Decrypted, &i.ChannelHash, &i.FirstHeardAt, &i.LastHeardAt, &i.ObservationCount, &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 }