Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions db/migrations/030_observation_airtime.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
-- Time-on-air per observation so per-observer activity never joins back to packets.
ALTER TABLE packet_observations ADD COLUMN airtime_ms REAL;
30 changes: 30 additions & 0 deletions db/migrations/031_observation_airtime_backfill.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
-- Separate from 030 so the ALTER's ACCESS EXCLUSIVE lock is released before this backfill runs.

-- Backfill only: mirrors internal/lora.TimeOnAirMs (RadioLib getTimeOnAir, CRC on, explicit header).
CREATE FUNCTION lora_airtime_ms_backfill(frame_len int, sf int, bw real, cr int) RETURNS real
LANGUAGE sql IMMUTABLE AS $$
SELECT (
(CASE WHEN sf <= 8 THEN 32 ELSE 16 END) + 4.25
+ 8 + GREATEST(ceil((8*frame_len - 4*sf + 44)::double precision
/ (4*(sf - 2*(CASE WHEN power(2, sf)/bw >= 16 THEN 1 ELSE 0 END)))) * cr, 0)
) * (power(2, sf)/bw)
$$;

-- Last 7 days like 015; parallelism off so the join spills to disk, not /dev/shm.
SET max_parallel_workers_per_gather = 0;

UPDATE packet_observations po
SET airtime_ms = lora_airtime_ms_backfill(
1 + CASE WHEN p.transport_codes_present THEN 4 ELSE 0 END + 1
+ COALESCE(octet_length(po.path_bytes), 0) + octet_length(p.raw_payload),
po.spread_factor, po.bandwidth_khz, po.coding_rate)
FROM packets p
WHERE p.packet_hash = po.packet_hash
AND po.heard_at > NOW() - INTERVAL '7 days'
AND po.airtime_ms IS NULL
AND po.spread_factor BETWEEN 7 AND 12
AND po.bandwidth_khz > 0
AND po.coding_rate BETWEEN 5 AND 8;

RESET max_parallel_workers_per_gather;
DROP FUNCTION lora_airtime_ms_backfill(int, int, real, int);
24 changes: 24 additions & 0 deletions db/migrations/032_mv_observer_activity.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
-- Hourly per-observer rollup so 7d/30d activity sums a few thousand rows instead of scanning
-- a month of observations. payload_type is a key so the breakdown comes from the same rows.
-- sqlc types the cast aggregates as non-null even though airtime_ms/snr_*/rssi_sum are NULL for
-- empty groups, so readers must gate on the matching *_n counts.
CREATE MATERIALIZED VIEW mv_observer_activity_hourly AS
SELECT
observer_id,
payload_type,
date_trunc('hour', heard_at)::timestamptz AS bucket,
COUNT(*)::bigint AS observations,
SUM(airtime_ms)::real AS airtime_ms,
COUNT(airtime_ms)::bigint AS airtime_n,
SUM(snr) FILTER (WHERE NOT (COALESCE(rssi, 0) = 0 AND COALESCE(snr, 0) = 0))::real AS snr_sum,
COUNT(snr) FILTER (WHERE NOT (COALESCE(rssi, 0) = 0 AND COALESCE(snr, 0) = 0))::bigint AS snr_n,
MIN(snr) FILTER (WHERE NOT (COALESCE(rssi, 0) = 0 AND COALESCE(snr, 0) = 0))::real AS snr_min,
SUM(rssi) FILTER (WHERE NOT (COALESCE(rssi, 0) = 0 AND COALESCE(snr, 0) = 0))::bigint AS rssi_sum,
COUNT(rssi) FILTER (WHERE NOT (COALESCE(rssi, 0) = 0 AND COALESCE(snr, 0) = 0))::bigint AS rssi_n
FROM packet_observations
WHERE heard_at > NOW() - INTERVAL '30 days'
AND payload_type IS NOT NULL
GROUP BY observer_id, payload_type, date_trunc('hour', heard_at);

CREATE UNIQUE INDEX idx_mv_observer_activity_hourly
ON mv_observer_activity_hourly(observer_id, payload_type, bucket);
124 changes: 124 additions & 0 deletions db/observers.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ import (
sqlc "github.com/MeshCore-Beacon/beacon-server/db/sqlc"
"github.com/MeshCore-Beacon/beacon-server/internal/api"
"github.com/MeshCore-Beacon/beacon-server/internal/ingest"
"github.com/MeshCore-Beacon/beacon-server/internal/lora"
"github.com/google/uuid"
"github.com/jackc/pgx/v5/pgtype"
)
Expand Down Expand Up @@ -226,6 +227,129 @@ func (s *Store) GetObserverTelemetryBucketed(ctx context.Context, observerID uui
return points, nil
}

// GetObserverActivity returns bucketed heard-activity for an observer over the trailing window.
// Buckets of an hour or coarser come from the hourly rollup; anything finer reads observations directly.
// Range and Interval are left empty for the handler to fill.
func (s *Store) GetObserverActivity(ctx context.Context, observerID uuid.UUID, window, interval time.Duration) (*api.ObserverActivity, error) {
obs, err := s.q.GetObserverByID(ctx, observerID)
if err != nil {
return nil, err
}
activity := &api.ObserverActivity{}
// radio is non-nil only when airtime is actually costable, so radio != null implies costed buckets
if obs.RadioSf != nil && obs.RadioBwKhz != nil && obs.RadioCr != nil &&
*obs.RadioSf >= 7 && *obs.RadioSf <= 12 && *obs.RadioBwKhz > 0 && *obs.RadioCr > 0 {
activity.Radio = &api.ObserverActivityRadio{
FreqMHz: obs.RadioFreqMhz,
SF: *obs.RadioSf,
BWKHz: *obs.RadioBwKhz,
CR: *obs.RadioCr,
PreambleSymbols: lora.PreambleSymbols(int(*obs.RadioSf)),
}
}
// round the window start up to a bucket boundary so the first bucket is never a partial one
start := time.Now().Add(-window).UTC()
since := start.Truncate(interval)
if since.Before(start) {
since = since.Add(interval)
}
sinceTS := pgtype.Timestamptz{Time: since, Valid: true}
binWidth := pgtype.Interval{Microseconds: interval.Microseconds(), Valid: true}

if interval >= time.Hour {
rows, err := s.q.GetObserverActivityHourly(ctx, sqlc.GetObserverActivityHourlyParams{
ObserverID: observerID,
Column2: sinceTS,
Column3: binWidth,
})
if err != nil {
return nil, err
}
activity.Points = make([]api.ObserverActivityPoint, 0, len(rows))
for _, r := range rows {
p := api.ObserverActivityPoint{T: r.Bucket.Time.UnixMilli(), Observations: r.Observations}
if r.AirtimeN > 0 {
airtime := r.AirtimeMs
p.AirtimeMs = &airtime
}
if r.SnrN > 0 {
avg := r.SnrSum / float32(r.SnrN)
min := r.SnrMin
p.SNRAvg, p.SNRMin = &avg, &min
}
if r.RssiN > 0 {
avg := float32(r.RssiSum) / float32(r.RssiN)
p.RSSIAvg = &avg
}
activity.Points = append(activity.Points, p)
}
typeRows, err := s.q.GetObserverActivityHourlyPayloadTypes(ctx, sqlc.GetObserverActivityHourlyPayloadTypesParams{
ObserverID: observerID,
Column2: sinceTS,
})
if err != nil {
return nil, err
}
activity.PayloadTypes = make([]api.PayloadBreakdownItem, 0, len(typeRows))
for _, v := range typeRows {
if v.PayloadType == nil {
continue
}
activity.PayloadTypes = append(activity.PayloadTypes, api.PayloadBreakdownItem{
PayloadType: *v.PayloadType,
PayloadTypeName: api.PayloadTypeName(*v.PayloadType),
Count: v.Count,
})
}
return activity, nil
}

rows, err := s.q.GetObserverActivityRaw(ctx, sqlc.GetObserverActivityRawParams{
ObserverID: observerID,
Column2: sinceTS,
Column3: binWidth,
})
if err != nil {
return nil, err
}
activity.Points = make([]api.ObserverActivityPoint, 0, len(rows))
for _, r := range rows {
p := api.ObserverActivityPoint{T: r.Bucket.Time.UnixMilli(), Observations: r.Observations}
if r.AirtimeN > 0 {
airtime := r.AirtimeMs
p.AirtimeMs = &airtime
}
if r.SnrN > 0 {
avg, min := r.SnrAvg, r.SnrMin
p.SNRAvg, p.SNRMin = &avg, &min
}
if r.RssiN > 0 {
avg := r.RssiAvg
p.RSSIAvg = &avg
}
activity.Points = append(activity.Points, p)
}
typeRows, err := s.q.GetObserverActivityRawPayloadTypes(ctx, sqlc.GetObserverActivityRawPayloadTypesParams{
ObserverID: observerID,
Column2: sinceTS,
})
if err != nil {
return nil, err
}
activity.PayloadTypes = make([]api.PayloadBreakdownItem, 0, len(typeRows))
for _, v := range typeRows {
if v.PayloadType == nil {
continue
}
activity.PayloadTypes = append(activity.PayloadTypes, api.PayloadBreakdownItem{
PayloadType: *v.PayloadType,
PayloadTypeName: api.PayloadTypeName(*v.PayloadType),
Count: v.Count,
})
}
return activity, 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,
Expand Down
Loading
Loading