feat(observers): add telemetry bucketing to rest api

closes #37
This commit is contained in:
Enot (ded) Skelly
2026-06-08 12:32:22 -07:00
parent cace46b004
commit 4147a93e8b
8 changed files with 155 additions and 10 deletions
+33 -2
View File
@@ -166,8 +166,6 @@ func (s *Store) InsertObserverTelemetry(ctx context.Context, observerID uuid.UUI
}
func (s *Store) GetObserverTelemetry(ctx context.Context, observerID uuid.UUID, since, until time.Time, afterID int64) (*api.ObserverTelemetry, error) {
// TODO: implement server-side bucketing by interval when needed.
// Currently returns all points in the range at stored resolution.
rows, err := s.q.GetObserverTelemetry(ctx, sqlc.GetObserverTelemetryParams{
ObserverID: observerID,
Column2: pgtype.Timestamptz{Time: since, Valid: !since.IsZero()},
@@ -193,6 +191,39 @@ func (s *Store) GetObserverTelemetry(ctx context.Context, observerID uuid.UUID,
return &api.ObserverTelemetry{Points: points}, nil
}
func (s *Store) GetObserverTelemetryBucketed(ctx context.Context, observerID uuid.UUID, since, until time.Time, bucketHours int32) ([]api.ObserverTelemetryPoint, error) {
var sinceTS, untilTS pgtype.Timestamptz
if !since.IsZero() {
sinceTS = pgtype.Timestamptz{Time: since, Valid: true}
}
if !until.IsZero() {
untilTS = pgtype.Timestamptz{Time: until, Valid: true}
}
rows, err := s.q.GetObserverTelemetryBucketed(ctx, sqlc.GetObserverTelemetryBucketedParams{
ObserverID: observerID,
Column2: sinceTS,
Column3: untilTS,
Column4: bucketHours,
})
if err != nil {
return nil, err
}
points := make([]api.ObserverTelemetryPoint, 0, len(rows))
for _, r := range rows {
points = append(points, api.ObserverTelemetryPoint{
T: r.Bucket.Time.UnixMilli(),
BatteryMV: &r.BatteryVoltageMv,
AirtimeTxPct: &r.AirtimeTxPct,
AirtimeRxPct: &r.AirtimeRxPct,
NoiseFloorDB: &r.NoiseFloorDb,
UptimeSeconds: &r.UptimeSeconds,
QueueLength: &r.QueueLength,
ReceiveErrors: &r.ReceiveErrors,
})
}
return points, nil
}
func (s *Store) ListObserverAdverts(ctx context.Context, observerID uuid.UUID, cursor int64, limit int32) (api.Page[api.AdvertObservation], error) {
rows, err := s.q.ListObserverAdverts(ctx, sqlc.ListObserverAdvertsParams{
ObserverID: observerID,
+18
View File
@@ -210,6 +210,24 @@ WHERE observer_id = $1
AND ($4 = 0 OR id > $4)
ORDER BY reported_at ASC;
-- name: GetObserverTelemetryBucketed :many
SELECT
(date_trunc('hour', reported_at) +
(EXTRACT(HOUR FROM reported_at)::int / $4::int) * ($4::int * interval '1 hour'))::timestamptz AS bucket,
AVG(battery_voltage_mv)::int AS battery_voltage_mv,
AVG(airtime_tx_pct)::real AS airtime_tx_pct,
AVG(airtime_rx_pct)::real AS airtime_rx_pct,
AVG(noise_floor_db)::real AS noise_floor_db,
MAX(uptime_seconds)::bigint AS uptime_seconds,
AVG(queue_length)::int AS queue_length,
AVG(receive_errors)::int AS receive_errors
FROM observer_telemetry
WHERE observer_id = $1
AND ($2::timestamptz IS NULL OR reported_at >= $2)
AND ($3::timestamptz IS NULL OR reported_at <= $3)
GROUP BY bucket
ORDER BY bucket ASC;
-- name: ListObserverAdverts :many
-- Returns advert packets (payload_type=4) heard by a specific observer.
-- Pass cursor=0 to start from the beginning, or the last seen id for pagination.
+71
View File
@@ -783,6 +783,77 @@ func (q *Queries) GetObserverTelemetry(ctx context.Context, arg GetObserverTelem
return items, nil
}
const getObserverTelemetryBucketed = `-- name: GetObserverTelemetryBucketed :many
SELECT
(date_trunc('hour', reported_at) +
(EXTRACT(HOUR FROM reported_at)::int / $4::int) * ($4::int * interval '1 hour'))::timestamptz AS bucket,
AVG(battery_voltage_mv)::int AS battery_voltage_mv,
AVG(airtime_tx_pct)::real AS airtime_tx_pct,
AVG(airtime_rx_pct)::real AS airtime_rx_pct,
AVG(noise_floor_db)::real AS noise_floor_db,
MAX(uptime_seconds)::bigint AS uptime_seconds,
AVG(queue_length)::int AS queue_length,
AVG(receive_errors)::int AS receive_errors
FROM observer_telemetry
WHERE observer_id = $1
AND ($2::timestamptz IS NULL OR reported_at >= $2)
AND ($3::timestamptz IS NULL OR reported_at <= $3)
GROUP BY bucket
ORDER BY bucket ASC
`
type GetObserverTelemetryBucketedParams struct {
ObserverID uuid.UUID `json:"observer_id"`
Column2 pgtype.Timestamptz `json:"column_2"`
Column3 pgtype.Timestamptz `json:"column_3"`
Column4 int32 `json:"column_4"`
}
type GetObserverTelemetryBucketedRow struct {
Bucket pgtype.Timestamptz `json:"bucket"`
BatteryVoltageMv int32 `json:"battery_voltage_mv"`
AirtimeTxPct float32 `json:"airtime_tx_pct"`
AirtimeRxPct float32 `json:"airtime_rx_pct"`
NoiseFloorDb float32 `json:"noise_floor_db"`
UptimeSeconds int64 `json:"uptime_seconds"`
QueueLength int32 `json:"queue_length"`
ReceiveErrors int32 `json:"receive_errors"`
}
func (q *Queries) GetObserverTelemetryBucketed(ctx context.Context, arg GetObserverTelemetryBucketedParams) ([]GetObserverTelemetryBucketedRow, error) {
rows, err := q.db.Query(ctx, getObserverTelemetryBucketed,
arg.ObserverID,
arg.Column2,
arg.Column3,
arg.Column4,
)
if err != nil {
return nil, err
}
defer rows.Close()
items := []GetObserverTelemetryBucketedRow{}
for rows.Next() {
var i GetObserverTelemetryBucketedRow
if err := rows.Scan(
&i.Bucket,
&i.BatteryVoltageMv,
&i.AirtimeTxPct,
&i.AirtimeRxPct,
&i.NoiseFloorDb,
&i.UptimeSeconds,
&i.QueueLength,
&i.ReceiveErrors,
); err != nil {
return nil, err
}
items = append(items, i)
}
if err := rows.Err(); err != nil {
return nil, err
}
return items, nil
}
const getPacketByHash = `-- name: GetPacketByHash :one
SELECT p.packet_hash, p.payload_type, p.payload_version, p.route_type, p.transport_codes_present, p.region_code, p.sub_region_code, p.scope_id, p.origin_pubkey, p.raw_payload, p.raw_header, p.parsed_payload, p.decrypted, p.channel_hash, p.trace_tag, p.first_heard_at, p.last_heard_at, ts.name AS scope_name,
cm.sender_name AS cm_sender_name,
+1 -1
View File
@@ -913,7 +913,7 @@ const docTemplate = `{
},
{
"type": "string",
"description": "Bucketing interval, echoed back in the response; not yet applied server-side",
"description": "Bucketing interval: 1h (default), 6h, or 24h",
"name": "interval",
"in": "query"
}
+1 -1
View File
@@ -911,7 +911,7 @@
},
{
"type": "string",
"description": "Bucketing interval, echoed back in the response; not yet applied server-side",
"description": "Bucketing interval: 1h (default), 6h, or 24h",
"name": "interval",
"in": "query"
}
+1 -2
View File
@@ -1544,8 +1544,7 @@ paths:
in: query
name: afterId
type: integer
- description: Bucketing interval, echoed back in the response; not yet applied
server-side
- description: 'Bucketing interval: 1h (default), 6h, or 24h'
in: query
name: interval
type: string
+27 -3
View File
@@ -171,7 +171,7 @@ func listObserverAdverts(reader api.Reader) http.HandlerFunc {
// @Param observerId path string true "Observer UUID"
// @Param range query string false "Duration window e.g. 24h, 48h, 168h (default 24h)"
// @Param afterId query int false "Return points after this telemetry ID for WS reconnection backfill"
// @Param interval query string false "Bucketing interval, echoed back in the response; not yet applied server-side"
// @Param interval query string false "Bucketing interval: 1h (default), 6h, or 24h"
// @Success 200 {object} api.ObserverTelemetry
// @Failure 400 {object} handlers.APIError
// @Failure 500 {object} handlers.APIError
@@ -201,15 +201,39 @@ func getObserverTelemetry(reader api.Reader) http.HandlerFunc {
}
afterID = id
}
intervalParam := r.URL.Query().Get("interval")
if intervalParam == "" {
intervalParam = "1h"
}
var bucketHours int32
switch intervalParam {
case "1h":
bucketHours = 0 // use raw query
case "6h":
bucketHours = 6
case "24h":
bucketHours = 24
default:
respondError(w, http.StatusBadRequest, "invalid interval, use 1h, 6h or 24h")
return
}
since := time.Now().Add(-duration)
until := time.Time{} // no upper bound
telemetry, err := reader.GetObserverTelemetry(r.Context(), observerID, since, until, afterID)
var telemetry *api.ObserverTelemetry
if bucketHours == 0 {
telemetry, err = reader.GetObserverTelemetry(r.Context(), observerID, since, until, afterID)
} else {
points, err := reader.GetObserverTelemetryBucketed(r.Context(), observerID, since, until, bucketHours)
if err == nil {
telemetry = &api.ObserverTelemetry{Points: points}
}
}
if err != nil {
respondError(w, http.StatusInternalServerError, "internal server error")
return
}
telemetry.Range = rangeParam
telemetry.Interval = r.URL.Query().Get("interval") // echoed back, not used server-side yet
telemetry.Interval = intervalParam
respond(w, http.StatusOK, telemetry)
}
}
+3 -1
View File
@@ -81,7 +81,9 @@ type Reader interface {
// GetObserverTelemetry returns telemetry points for an observer within the given time range.
// since and until define the window; pass zero times to use defaults (last 24h).
GetObserverTelemetry(ctx context.Context, observerID uuid.UUID, since, until time.Time, afterID int64) (*ObserverTelemetry, error)
// TODO: add interval time.Duration param for server-side bucketing
// GetObserverTelemetryBucketed returns telemetry points for an observer bucketed into N-hour intervals.
GetObserverTelemetryBucketed(ctx context.Context, observerID uuid.UUID, since, until time.Time, bucketHours int32) ([]ObserverTelemetryPoint, error)
// GetObserverScopes returns the names of all transport scopes an observer has
// been seen forwarding packets for, ordered alphabetically.