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
64 changes: 52 additions & 12 deletions internal/modules/mqtt_ha.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,16 +33,20 @@ type mqttState struct {
IPv6Enabled bool `json:"ipv6_enabled"`
SQMEnabled bool `json:"sqm_enabled"`
Mode string `json:"mode"`
// Aggregated client telemetry (#401).
ClientsTotal int `json:"clients_total"`
ClientsWifi int `json:"clients_wifi"`
Clients24 int `json:"clients_24"`
Clients5 int `json:"clients_5"`
ClientsCable int `json:"clients_cable"`
ClientsWeak int `json:"clients_weak"`
RxBytes int64 `json:"rx_bytes"`
TxBytes int64 `json:"tx_bytes"`
Ts int64 `json:"ts"`
// Aggregated client telemetry (#401) and Wi-Fi/LAN traffic telemetry (#403).
ClientsTotal int `json:"clients_total"`
ClientsWifi int `json:"clients_wifi"`
Clients24 int `json:"clients_24"`
Clients5 int `json:"clients_5"`
ClientsCable int `json:"clients_cable"`
ClientsWeak int `json:"clients_weak"`
WifiMinSignal int `json:"wifi_min_signal"`
WifiAvgSignal int `json:"wifi_avg_signal"`
RxBytes int64 `json:"rx_bytes"`
TxBytes int64 `json:"tx_bytes"`
RxMbps float64 `json:"rx_mbps"`
TxMbps float64 `json:"tx_mbps"`
Ts int64 `json:"ts"`
}

type mqttBoard struct {
Expand Down Expand Up @@ -99,8 +103,12 @@ func buildMQTTState(node, version string) mqttState {
st.Clients5 = tel.Wifi5
st.ClientsCable = tel.Cable
st.ClientsWeak = tel.Weak
st.WifiMinSignal = tel.WifiMinSignal
st.WifiAvgSignal = tel.WifiAvgSignal
st.RxBytes = tel.RxBytes
st.TxBytes = tel.TxBytes
st.RxMbps = tel.RxMbps
st.TxMbps = tel.TxMbps

return st
}
Expand Down Expand Up @@ -305,8 +313,40 @@ func mqttDiscoveryEntities(node, version, model string) []mqttEntity {
"icon": "mdi:wifi-alert",
"entity_category": "diagnostic",
}),
newEntity("sensor", "wifi_signal_min", map[string]any{
"name": "Weakest client signal",
"value_template": "{{ value_json.wifi_min_signal }}",
"unit_of_measurement": "dBm",
"device_class": "signal_strength",
"state_class": "measurement",
"icon": "mdi:wifi-strength-1",
"entity_category": "diagnostic",
}),
newEntity("sensor", "wifi_signal_avg", map[string]any{
"name": "Average client signal",
"value_template": "{{ value_json.wifi_avg_signal }}",
"unit_of_measurement": "dBm",
"device_class": "signal_strength",
"state_class": "measurement",
"icon": "mdi:wifi-strength-3",
"entity_category": "diagnostic",
}),
newEntity("sensor", "rx_mbps", map[string]any{
"name": "LAN traffic in",
"value_template": "{{ value_json.rx_mbps }}",
"unit_of_measurement": "Mbit/s",
"state_class": "measurement",
"icon": "mdi:download-network",
}),
newEntity("sensor", "tx_mbps", map[string]any{
"name": "LAN traffic out",
"value_template": "{{ value_json.tx_mbps }}",
"unit_of_measurement": "Mbit/s",
"state_class": "measurement",
"icon": "mdi:upload-network",
}),
newEntity("sensor", "rx_bytes", map[string]any{
"name": "Clients received",
"name": "LAN received total",
"value_template": "{{ value_json.rx_bytes }}",
"unit_of_measurement": "B",
"device_class": "data_size",
Expand All @@ -315,7 +355,7 @@ func mqttDiscoveryEntities(node, version, model string) []mqttEntity {
"entity_category": "diagnostic",
}),
newEntity("sensor", "tx_bytes", map[string]any{
"name": "Clients sent",
"name": "LAN sent total",
"value_template": "{{ value_json.tx_bytes }}",
"unit_of_measurement": "B",
"device_class": "data_size",
Expand Down
107 changes: 92 additions & 15 deletions internal/modules/telemetry.go
Original file line number Diff line number Diff line change
@@ -1,19 +1,31 @@
// telemetry.go: aggregated client and traffic telemetry (#401) used by the
// MQTT integration to expose router activity to Home Assistant without a
// telemetry.go: aggregated client and traffic telemetry (#401, #403) used by
// the MQTT integration to expose router activity to Home Assistant without a
// NetPulse server.
package modules

import "time"
import (
"math"
"sync"
"time"
)

// TelemetryProbe is the aggregated view of the connected clients.
// TelemetryProbe is the aggregated view of the connected clients plus the LAN
// traffic counters.
type TelemetryProbe struct {
Total int `json:"total"`
Wifi24 int `json:"wifi24"`
Wifi5 int `json:"wifi5"`
Cable int `json:"cable"`
Weak int `json:"weak"`
RxBytes int64 `json:"rx_bytes"`
TxBytes int64 `json:"tx_bytes"`
Total int `json:"total"`
Wifi24 int `json:"wifi24"`
Wifi5 int `json:"wifi5"`
Cable int `json:"cable"`
Weak int `json:"weak"`
WifiMinSignal int `json:"wifi_min_signal"`
WifiAvgSignal int `json:"wifi_avg_signal"`
// LAN bridge counters: cumulative bytes and the rate derived from the
// previous sample. "Rx" is traffic entering the bridge (from the LAN
// clients), "tx" is traffic leaving it towards them.
RxBytes int64 `json:"rx_bytes"`
TxBytes int64 `json:"tx_bytes"`
RxMbps float64 `json:"rx_mbps"`
TxMbps float64 `json:"tx_mbps"`
}

// weakSignalDbm: clients at or below this signal count as weak. Chosen from
Expand All @@ -25,22 +37,77 @@ const weakSignalDbm = -75
// ListClients forks several commands, so it is cached briefly.
const telemetryTTL = 15 * time.Second

// trafficSample keeps the last LAN bridge counters so the next sample can
// turn them into a rate. Package-level on purpose: the probe is cached, so
// samples are spaced by telemetryTTL at the very least.
type trafficSample struct {
at time.Time
rx, tx int64
}

var (
trafficMu sync.Mutex
lastTraffic trafficSample
)

// ProbeTelemetry returns the aggregated client telemetry, cached for
// telemetryTTL.
func ProbeTelemetry() *TelemetryProbe {
p, _ := CachedRead("telemetry", telemetryTTL, func() (*TelemetryProbe, error) {
return aggregateClients(ListClients("")), nil
return sampleTelemetry(), nil
})
if p == nil {
return &TelemetryProbe{}
}
return p
}

// sampleTelemetry folds the client list and the LAN bridge counters into one
// probe, deriving the traffic rate from the previous sample.
func sampleTelemetry() *TelemetryProbe {
t := aggregateClients(ListClients(""))
t.RxBytes, t.TxBytes = bridgeCounters()

now := time.Now()
trafficMu.Lock()
if !lastTraffic.at.IsZero() {
secs := now.Sub(lastTraffic.at).Seconds()
t.RxMbps = rateMbps(t.RxBytes-lastTraffic.rx, secs)
t.TxMbps = rateMbps(t.TxBytes-lastTraffic.tx, secs)
}
lastTraffic = trafficSample{at: now, rx: t.RxBytes, tx: t.TxBytes}
trafficMu.Unlock()

return t
}

// bridgeCounters returns the cumulative LAN bridge counters, the same source
// the panel traffic graph uses. Per-client counters are not used: most devices
// do not account per station and report zero, and APs often have no WAN at all.
func bridgeCounters() (int64, int64) {
bridge := LANBridge()
for _, c := range NetDevCounters() {
if c.Name == bridge {
return c.RxBytes, c.TxBytes
}
}
return 0, 0
}

// rateMbps converts a byte delta over an interval into Mbit/s, rounded to two
// decimals. A negative delta (counter reset) reads as zero.
func rateMbps(deltaBytes int64, seconds float64) float64 {
if seconds <= 0 || deltaBytes <= 0 {
return 0
}
return math.Round(float64(deltaBytes)*8/seconds/1e6*100) / 100
}

// aggregateClients folds a client list into the telemetry counters. Pure so it
// can be tested without touching the router.
func aggregateClients(clients []Client) *TelemetryProbe {
t := &TelemetryProbe{}
var signalSum, signalCount int
for _, c := range clients {
t.Total++
switch c.Type {
Expand All @@ -51,11 +118,21 @@ func aggregateClients(clients []Client) *TelemetryProbe {
default:
t.Cable++
}
if c.Type != "cable" && c.Signal != 0 && c.Signal < weakSignalDbm {
// Only Wi-Fi clients with a real measurement feed the signal figures.
if (c.Type != "wifi24" && c.Type != "wifi5") || c.Signal == 0 {
continue
}
signalSum += c.Signal
signalCount++
if t.WifiMinSignal == 0 || c.Signal < t.WifiMinSignal {
t.WifiMinSignal = c.Signal
}
if c.Signal < weakSignalDbm {
t.Weak++
}
t.RxBytes += c.RxBytes
t.TxBytes += c.TxBytes
}
if signalCount > 0 {
t.WifiAvgSignal = int(math.Round(float64(signalSum) / float64(signalCount)))
}
return t
}
49 changes: 43 additions & 6 deletions internal/modules/telemetry_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,10 +4,10 @@ import "testing"

func TestAggregateClients(t *testing.T) {
tel := aggregateClients([]Client{
{Type: "wifi24", Signal: -50, RxBytes: 100, TxBytes: 200},
{Type: "wifi5", Signal: -80, RxBytes: 300, TxBytes: 400}, // weak
{Type: "wifi5", Signal: 0, RxBytes: 1, TxBytes: 2}, // no measurement
{Type: "cable", RxBytes: 5, TxBytes: 6},
{Type: "wifi24", Signal: -50},
{Type: "wifi5", Signal: -80}, // weak
{Type: "wifi5", Signal: 0}, // no measurement
{Type: "cable"},
})
if tel.Total != 4 {
t.Fatalf("total = %d, want 4", tel.Total)
Expand All @@ -18,8 +18,24 @@ func TestAggregateClients(t *testing.T) {
if tel.Weak != 1 {
t.Fatalf("weak = %d, want 1 (signal 0 must not count)", tel.Weak)
}
if tel.RxBytes != 406 || tel.TxBytes != 608 {
t.Fatalf("bytes = %d/%d, want 406/608", tel.RxBytes, tel.TxBytes)
if tel.WifiMinSignal != -80 {
t.Fatalf("min signal = %d, want -80", tel.WifiMinSignal)
}
if tel.WifiAvgSignal != -65 {
t.Fatalf("avg signal = %d, want -65 (only measured clients)", tel.WifiAvgSignal)
}
}

func TestAggregateClientsNoMeasurement(t *testing.T) {
tel := aggregateClients([]Client{
{Type: "cable"},
{Type: "wifi5", Signal: 0},
})
if tel.WifiMinSignal != 0 || tel.WifiAvgSignal != 0 {
t.Fatalf("signal figures = %d/%d, want 0/0", tel.WifiMinSignal, tel.WifiAvgSignal)
}
if tel.Weak != 0 {
t.Fatalf("weak = %d, want 0", tel.Weak)
}
}

Expand All @@ -30,6 +46,27 @@ func TestAggregateClientsEmpty(t *testing.T) {
}
}

func TestRateMbps(t *testing.T) {
cases := []struct {
name string
delta int64
seconds float64
want float64
}{
{"no time", 1000, 0, 0},
{"negative delta (counter reset)", -1000, 60, 0},
{"no traffic", 0, 60, 0},
{"one MB over a minute", 7_500_000, 60, 1},
{"one MB in one second", 1_000_000, 1, 8},
{"rounded to two decimals", 1_234_567, 60, 0.16},
}
for _, tc := range cases {
if got := rateMbps(tc.delta, tc.seconds); got != tc.want {
t.Errorf("%s: rateMbps(%d, %v) = %v, want %v", tc.name, tc.delta, tc.seconds, got, tc.want)
}
}
}

func TestProbeTelemetryNeverNil(t *testing.T) {
if tel := ProbeTelemetry(); tel == nil {
t.Fatal("ProbeTelemetry returned nil")
Expand Down
Loading