mirror of
https://github.com/MeshCore-Beacon/beacon-server.git
synced 2026-09-01 16:48:19 +00:00
adds insert telemetry for observers in ingestion adds new config for observer telemetry, moves packet retention into config for consistency with defaults on load adds cleanup for both packets and observer telemetry
1895 lines
54 KiB
Go
1895 lines
54 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 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 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
|
|
}
|