diff --git a/db/queries/queries.sql b/db/queries/queries.sql index 9ba4148..d18cd25 100644 --- a/db/queries/queries.sql +++ b/db/queries/queries.sql @@ -59,7 +59,8 @@ SELECT * FROM observers WHERE public_key = $1; SELECT * FROM observers WHERE id = $1; -- name: GetObserverBrokers :many -SELECT broker_name FROM observer_brokers +SELECT broker_name, last_seen, last_packet_at +FROM observer_brokers WHERE observer_id = $1 ORDER BY last_seen DESC; diff --git a/db/sqlc/queries.sql.go b/db/sqlc/queries.sql.go index 7d4046b..9015616 100644 --- a/db/sqlc/queries.sql.go +++ b/db/sqlc/queries.sql.go @@ -210,24 +210,31 @@ func (q *Queries) GetNodeByPubkey(ctx context.Context, publicKey []byte) (Node, } const getObserverBrokers = `-- name: GetObserverBrokers :many -SELECT broker_name FROM observer_brokers +SELECT broker_name, last_seen, last_packet_at +FROM observer_brokers WHERE observer_id = $1 ORDER BY last_seen DESC ` -func (q *Queries) GetObserverBrokers(ctx context.Context, observerID uuid.UUID) ([]string, error) { +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 := []string{} + items := []GetObserverBrokersRow{} for rows.Next() { - var broker_name string - if err := rows.Scan(&broker_name); err != nil { + var i GetObserverBrokersRow + if err := rows.Scan(&i.BrokerName, &i.LastSeen, &i.LastPacketAt); err != nil { return nil, err } - items = append(items, broker_name) + items = append(items, i) } if err := rows.Err(); err != nil { return nil, err diff --git a/db/store.go b/db/store.go index bdb0d0e..93fc18c 100644 --- a/db/store.go +++ b/db/store.go @@ -536,7 +536,7 @@ func (s *Store) GetObserver(ctx context.Context, observerID uuid.UUID) (*api.Obs if err != nil { return nil, err } - brokers, err := s.q.GetObserverBrokers(ctx, observerID) + brokerRows, err := s.q.GetObserverBrokers(ctx, observerID) if err != nil { return nil, err } @@ -562,8 +562,20 @@ func (s *Store) GetObserver(ctx context.Context, observerID uuid.UUID) (*api.Obs FirstSeen: obs.FirstSeen.Time.UnixMilli(), LastSeen: obs.LastSeen.Time.UnixMilli(), ObservationCount: *obs.ObservationCount, - Brokers: brokers, } + brokers := make([]api.ObserverBroker, 0, len(brokerRows)) + for _, v := range brokerRows { + var lastPacketAt int64 + if v.LastPacketAt.Valid { + lastPacketAt = v.LastPacketAt.Time.UnixMilli() + } + brokers = append(brokers, api.ObserverBroker{ + Name: v.BrokerName, + LastPacketAt: lastPacketAt, + LastSeenAt: v.LastSeen.Time.UnixMilli(), + }) + } + observer.Brokers = brokers if obs.LastStatusAt.Valid && time.Since(obs.LastStatusAt.Time) < 5*time.Minute { observer.Status = "online" } diff --git a/internal/api/reader.go b/internal/api/reader.go index 783ef03..424fdc3 100644 --- a/internal/api/reader.go +++ b/internal/api/reader.go @@ -18,27 +18,36 @@ type ObserverSummary struct { Status string `json:"status"` // "online" or "offline" derived from last_status_at } +// ObserverBroker represents a single MQTT broker an observer has been seen on, +// including timestamps for diagnosing partial outages — e.g. distinguishing +// "observer is down" from "one broker stopped delivering for this observer". +type ObserverBroker struct { + Name string `json:"name"` // broker name e.g. "mqtt1" + LastSeenAt int64 `json:"lastSeenAt"` // epoch ms, last time observer was seen on this broker + LastPacketAt int64 `json:"lastPacketAt"` // epoch ms, last packet received via this broker; 0 if none +} + // Observer is the full observer representation including radio config, // telemetry, broker memberships and raw status metadata. type Observer struct { ObserverSummary - PublicKey string `json:"publicKey"` // hex-encoded public key - SoftwareVersion *string `json:"softwareVersion,omitempty"` - HardwareModel *string `json:"hardwareModel,omitempty"` - FirmwareVersion *string `json:"firmwareVersion,omitempty"` - FirmwareBuild *string `json:"firmwareBuild,omitempty"` - RadioFreqMHz *float32 `json:"radioFreqMhz,omitempty"` // MHz e.g. 910.525 - RadioSF *int16 `json:"radioSf,omitempty"` // LoRa spreading factor - RadioBWKHz *float32 `json:"radioBwKhz,omitempty"` // bandwidth in kHz - RadioCR *int16 `json:"radioCr,omitempty"` // coding rate denominator - BatteryLevel *float32 `json:"batteryLevel,omitempty"` // volts, nil if mains powered - UptimeSeconds *int64 `json:"uptimeSeconds,omitempty"` - StatusMetadata json.RawMessage `json:"statusMetadata,omitempty"` // raw /status JSON payload - LastStatusAt *int64 `json:"lastStatusAt,omitempty"` // epoch ms - FirstSeen int64 `json:"firstSeen"` // epoch ms - LastSeen int64 `json:"lastSeen"` // epoch ms - ObservationCount int64 `json:"observationCount"` - Brokers []string `json:"brokers"` // broker names this observer has been seen on + PublicKey string `json:"publicKey"` // hex-encoded public key + SoftwareVersion *string `json:"softwareVersion,omitempty"` + HardwareModel *string `json:"hardwareModel,omitempty"` + FirmwareVersion *string `json:"firmwareVersion,omitempty"` + FirmwareBuild *string `json:"firmwareBuild,omitempty"` + RadioFreqMHz *float32 `json:"radioFreqMhz,omitempty"` // MHz e.g. 910.525 + RadioSF *int16 `json:"radioSf,omitempty"` // LoRa spreading factor + RadioBWKHz *float32 `json:"radioBwKhz,omitempty"` // bandwidth in kHz + RadioCR *int16 `json:"radioCr,omitempty"` // coding rate denominator + BatteryLevel *float32 `json:"batteryLevel,omitempty"` // volts, nil if mains powered + UptimeSeconds *int64 `json:"uptimeSeconds,omitempty"` + StatusMetadata json.RawMessage `json:"statusMetadata,omitempty"` // raw /status JSON payload + LastStatusAt *int64 `json:"lastStatusAt,omitempty"` // epoch ms + FirstSeen int64 `json:"firstSeen"` // epoch ms + LastSeen int64 `json:"lastSeen"` // epoch ms + ObservationCount int64 `json:"observationCount"` + Brokers []ObserverBroker `json:"brokers"` // broker names this observer has been seen on } // ChannelMessage represents a single decrypted channel message.