From 4147a93e8b1b21ff5bd77318fde6ea06f38b64e1 Mon Sep 17 00:00:00 2001 From: "Enot (ded) Skelly" Date: Mon, 8 Jun 2026 12:31:56 -0700 Subject: [PATCH] feat(observers): add telemetry bucketing to rest api closes #37 --- db/observers.go | 35 ++++++++++++++- db/queries/queries.sql | 18 ++++++++ db/sqlc/queries.sql.go | 71 ++++++++++++++++++++++++++++++ docs/docs.go | 2 +- docs/swagger.json | 2 +- docs/swagger.yaml | 3 +- internal/api/handlers/observers.go | 30 +++++++++++-- internal/api/reader.go | 4 +- 8 files changed, 155 insertions(+), 10 deletions(-) diff --git a/db/observers.go b/db/observers.go index efeee07..f37dba2 100644 --- a/db/observers.go +++ b/db/observers.go @@ -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, diff --git a/db/queries/queries.sql b/db/queries/queries.sql index 7971830..0dcc4d8 100644 --- a/db/queries/queries.sql +++ b/db/queries/queries.sql @@ -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. diff --git a/db/sqlc/queries.sql.go b/db/sqlc/queries.sql.go index 3360aa0..95c3046 100644 --- a/db/sqlc/queries.sql.go +++ b/db/sqlc/queries.sql.go @@ -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, diff --git a/docs/docs.go b/docs/docs.go index 827aed8..402b400 100644 --- a/docs/docs.go +++ b/docs/docs.go @@ -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" } diff --git a/docs/swagger.json b/docs/swagger.json index 24d2311..2aa6056 100644 --- a/docs/swagger.json +++ b/docs/swagger.json @@ -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" } diff --git a/docs/swagger.yaml b/docs/swagger.yaml index 0106fe8..0c1bc8f 100644 --- a/docs/swagger.yaml +++ b/docs/swagger.yaml @@ -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 diff --git a/internal/api/handlers/observers.go b/internal/api/handlers/observers.go index 5d376f9..40fac5f 100644 --- a/internal/api/handlers/observers.go +++ b/internal/api/handlers/observers.go @@ -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) } } diff --git a/internal/api/reader.go b/internal/api/reader.go index 3a7b080..bb7c7c7 100644 --- a/internal/api/reader.go +++ b/internal/api/reader.go @@ -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.