feat(metrics): implement metrics data retrieval with RRDtool fallback

This commit is contained in:
Rightup
2026-07-20 13:44:14 +01:00
parent 7de15b135e
commit efb7e4a319
8 changed files with 738 additions and 36 deletions
+5
View File
@@ -140,6 +140,11 @@ repeater:
policy:
policy_file: "policy.yaml"
metrics:
# Store historical metrics in RRDtool when available.
# When disabled, chart APIs calculate metrics from repeater.db.
rrd_enabled: true
# Local GPS receiver. When enabled, the daemon reads NMEA sentences from the
# configured source and exposes parsed data at /api/gps.
gps:
+169
View File
@@ -1642,6 +1642,175 @@ class SQLiteHandler:
logger.error(f"Failed to get packet stats: {e}")
return {}
def get_metrics_data(
self,
start_time: Optional[int] = None,
end_time: Optional[int] = None,
resolution: str = "average",
) -> dict:
resolution_key = str(resolution or "average").lower()
gauge_aggregates = {
"average": "AVG",
"max": "MAX",
"min": "MIN",
}
gauge_aggregate = gauge_aggregates.get(resolution_key, "AVG")
if end_time is None:
end_ts = int(time.time())
else:
end_ts = int(end_time)
if start_time is None:
start_ts = end_ts - (24 * 3600)
else:
start_ts = int(start_time)
if end_ts < start_ts:
start_ts, end_ts = end_ts, start_ts
range_seconds = max(0, end_ts - start_ts)
if range_seconds <= 7 * 24 * 3600:
bucket_seconds = 60
elif range_seconds <= 30 * 24 * 3600:
bucket_seconds = 300
else:
bucket_seconds = 3600
aligned_start = int(start_ts / bucket_seconds) * bucket_seconds
aligned_end = int(end_ts / bucket_seconds) * bucket_seconds
timestamps = list(range(aligned_start, aligned_end + bucket_seconds, bucket_seconds))
metric_names = [
"rx_count",
"tx_count",
"drop_count",
"avg_rssi",
"avg_snr",
"avg_length",
"avg_score",
"neighbor_count",
]
packet_type_names = [f"type_{i}" for i in range(16)] + ["type_other"]
metrics = {
"rx_count": [],
"tx_count": [],
"drop_count": [],
"avg_rssi": [],
"avg_snr": [],
"avg_length": [],
"avg_score": [],
# Historical neighbor counts are not stored in packets, so the
# existing schema cannot reconstruct past values per time bucket.
"neighbor_count": [],
}
packet_types = {name: [] for name in packet_type_names}
bucket_metrics = {
ts: {
"rx_count": 0,
"tx_count": 0,
"drop_count": 0,
"avg_rssi": None,
"avg_snr": None,
"avg_length": None,
"avg_score": None,
"neighbor_count": None,
}
for ts in timestamps
}
bucket_packet_types = {ts: {name: 0 for name in packet_type_names} for ts in timestamps}
try:
with self._connect() as conn:
conn.row_factory = sqlite3.Row
aggregate_rows = conn.execute(
f"""
SELECT
CAST(timestamp / ? AS INTEGER) * ? AS bucket_ts,
COUNT(*) AS rx_count,
SUM(CASE WHEN transmitted = 1 THEN 1 ELSE 0 END) AS tx_count,
SUM(CASE WHEN transmitted = 0 THEN 1 ELSE 0 END) AS drop_count,
{gauge_aggregate}(rssi) AS avg_rssi,
{gauge_aggregate}(snr) AS avg_snr,
{gauge_aggregate}(length) AS avg_length,
{gauge_aggregate}(score) AS avg_score
FROM packets INDEXED BY idx_packets_timestamp
WHERE timestamp >= ? AND timestamp <= ?
GROUP BY bucket_ts
ORDER BY bucket_ts ASC
""",
(bucket_seconds, bucket_seconds, start_ts, end_ts),
).fetchall()
packet_type_rows = conn.execute(
"""
SELECT
CAST(timestamp / ? AS INTEGER) * ? AS bucket_ts,
CASE
WHEN type BETWEEN 0 AND 15 THEN CAST(type AS INTEGER)
ELSE 16
END AS type_bucket,
COUNT(*) AS count
FROM packets INDEXED BY idx_packets_timestamp
WHERE timestamp >= ? AND timestamp <= ?
GROUP BY bucket_ts, type_bucket
ORDER BY bucket_ts ASC, type_bucket ASC
""",
(bucket_seconds, bucket_seconds, start_ts, end_ts),
).fetchall()
for row in aggregate_rows:
bucket_ts = int(row["bucket_ts"])
if bucket_ts not in bucket_metrics:
continue
bucket_metrics[bucket_ts] = {
"rx_count": int(row["rx_count"] or 0),
"tx_count": int(row["tx_count"] or 0),
"drop_count": int(row["drop_count"] or 0),
"avg_rssi": row["avg_rssi"],
"avg_snr": row["avg_snr"],
"avg_length": row["avg_length"],
"avg_score": row["avg_score"],
"neighbor_count": None,
}
for row in packet_type_rows:
bucket_ts = int(row["bucket_ts"])
if bucket_ts not in bucket_packet_types:
continue
type_bucket = int(row["type_bucket"])
type_name = f"type_{type_bucket}" if 0 <= type_bucket <= 15 else "type_other"
bucket_packet_types[bucket_ts][type_name] = int(row["count"] or 0)
for timestamp in timestamps:
bucket = bucket_metrics[timestamp]
for name in metric_names:
metrics[name].append(bucket[name])
packet_bucket = bucket_packet_types[timestamp]
for name in packet_type_names:
packet_types[name].append(packet_bucket[name])
return {
"start_time": aligned_start,
"end_time": aligned_end,
"step": bucket_seconds,
"timestamps": timestamps,
"data_sources": metric_names + packet_type_names,
"packet_types": packet_types,
"metrics": metrics,
"data_source": "sqlite",
"counter_mode": "bucket_count",
}
except Exception as e:
logger.error(f"Failed to get SQLite metrics data: {e}", exc_info=True)
raise
def get_recent_packets(self, limit: int = 100) -> list:
try:
with self._connect() as conn:
+69 -8
View File
@@ -22,6 +22,11 @@ class StorageCollector:
self.glass_publish_callback = None
self._pending_tasks = set()
metrics_config = config.get("metrics")
if not isinstance(metrics_config, dict):
metrics_config = {}
self.rrd_enabled = bool(metrics_config.get("rrd_enabled", True))
# Dedicated single writer thread for all blocking storage work (the SQLite
# write, the cumulative-counts aggregate, RRD updates, and network
# publishing). This keeps that work off the asyncio event loop, which it
@@ -37,7 +42,16 @@ class StorageCollector:
self.storage_dir.mkdir(parents=True, exist_ok=True)
self.sqlite_handler = SQLiteHandler(self.storage_dir)
self.rrd_handler = RRDToolHandler(self.storage_dir)
self.rrd_handler = None
if self.rrd_enabled:
candidate_rrd_handler = RRDToolHandler(self.storage_dir)
if candidate_rrd_handler.available and candidate_rrd_handler.rrd_path.exists():
self.rrd_handler = candidate_rrd_handler
logger.info("RRDtool metrics enabled")
else:
logger.warning("RRDtool requested but unavailable; using SQLite metrics fallback")
else:
logger.info("RRDtool metrics disabled; SQLite metrics fallback will be used")
# Initialize MQTT handler only when at least one broker is configured
self.mqtt_handler = None
@@ -218,8 +232,11 @@ class StorageCollector:
packet_id = self.sqlite_handler.store_packet(packet_record)
if packet_id is not None:
packet_record["id"] = packet_id
cumulative_counts = self.sqlite_handler.get_cumulative_counts()
self.rrd_handler.update_packet_metrics(packet_record, cumulative_counts)
if self.rrd_handler is not None:
cumulative_counts = self.sqlite_handler.get_cumulative_counts()
self.rrd_handler.update_packet_metrics(packet_record, cumulative_counts)
self._publish_packet_sync(packet_record, skip_mqtt)
def _publish_packet_sync(self, packet_record: dict, skip_mqtt: bool):
@@ -462,14 +479,58 @@ class StorageCollector:
end_time: Optional[int] = None,
resolution: str = "average",
) -> Optional[dict]:
return self.rrd_handler.get_data(start_time, end_time, resolution)
return self.get_metrics_data(start_time, end_time, resolution)
def get_metrics_data(
self,
start_time: Optional[int] = None,
end_time: Optional[int] = None,
resolution: str = "average",
) -> dict:
if self.rrd_handler is not None:
try:
rrd_data = self.rrd_handler.get_data(start_time, end_time, resolution)
except Exception as e:
logger.warning(
f"RRDtool metrics read failed; using SQLite metrics fallback: {e}",
exc_info=True,
)
else:
if self._metrics_data_is_valid(rrd_data):
rrd_data.setdefault("data_source", "rrd")
return rrd_data
logger.warning(
"RRDtool metrics read returned no usable data; using SQLite metrics fallback"
)
sqlite_data = self.sqlite_handler.get_metrics_data(start_time, end_time, resolution)
sqlite_data.setdefault("data_source", "sqlite")
return sqlite_data
def _metrics_data_is_valid(self, metrics_data: Optional[dict]) -> bool:
if not isinstance(metrics_data, dict):
return False
if not isinstance(metrics_data.get("metrics"), dict):
return False
if not isinstance(metrics_data.get("timestamps"), list):
return False
return True
def get_packet_type_stats(self, hours: int = 24) -> dict:
rrd_stats = self.rrd_handler.get_packet_type_stats(hours)
if rrd_stats:
return rrd_stats
if self.rrd_handler is not None:
try:
rrd_stats = self.rrd_handler.get_packet_type_stats(hours)
except Exception as e:
logger.warning(
f"RRDtool packet type stats failed; using SQLite fallback: {e}",
exc_info=True,
)
else:
if rrd_stats:
return rrd_stats
logger.warning("Falling back to SQLite for packet type stats")
logger.warning("Falling back to SQLite for packet type stats")
return self.sqlite_handler.get_packet_type_stats(hours)
def get_route_stats(self, hours: int = 24) -> dict:
+44 -25
View File
@@ -503,9 +503,15 @@ class APIEndpoints:
previous = item
return output
rx_values = _counter_delta(metrics.get("rx_count", []))
tx_values = _counter_delta(metrics.get("tx_count", []))
drop_values = _counter_delta(metrics.get("drop_count", []))
def _counter_values(values: list) -> list[float]:
return [float(item or 0.0) for item in values]
counter_mode = rrd_data.get("counter_mode")
counter_processor = _counter_values if counter_mode == "bucket_count" else _counter_delta
rx_values = counter_processor(metrics.get("rx_count", []))
tx_values = counter_processor(metrics.get("tx_count", []))
drop_values = counter_processor(metrics.get("drop_count", []))
rssi_values = metrics.get("avg_rssi", []) or []
snr_values = metrics.get("avg_snr", []) or []
@@ -3322,8 +3328,8 @@ class APIEndpoints:
params = self._get_params(
{"start_time": None, "end_time": None, "resolution": "average"}
)
data = self._get_storage().get_rrd_data(**params)
return self._success(data) if data else self._error("No RRD data available")
data = self._get_storage().get_metrics_data(**params)
return self._success(data)
except ValueError as e:
return self._error(f"Invalid parameter format: {e}")
except Exception as e:
@@ -3392,12 +3398,16 @@ class APIEndpoints:
hours = int(hours)
start_time, end_time = self._get_time_range(hours)
rrd_data = self._get_storage().get_rrd_data(
metrics_data = self._get_storage().get_metrics_data(
start_time=start_time, end_time=end_time, resolution=resolution
)
if not rrd_data or "metrics" not in rrd_data:
return self._error("No RRD data available")
if (
not metrics_data
or "metrics" not in metrics_data
or "timestamps" not in metrics_data
):
return self._error("No metrics data available")
metric_names = {
"rx_count": "Received Packets",
@@ -3416,18 +3426,18 @@ class APIEndpoints:
if metrics != "all":
requested_metrics = [m.strip() for m in metrics.split(",")]
else:
requested_metrics = list(rrd_data["metrics"].keys())
requested_metrics = list(metrics_data["metrics"].keys())
requested_metrics.append("policy_events")
timestamps_ms = [ts * 1000 for ts in rrd_data["timestamps"]]
timestamps_ms = [ts * 1000 for ts in metrics_data["timestamps"]]
series = []
for metric_key in requested_metrics:
if metric_key == "policy_events":
bucket_seconds = max(1, int(rrd_data.get("step", 60)))
bucket_seconds = max(1, int(metrics_data.get("step", 60)))
policy_rows = self._get_storage().get_policy_event_counts(
start_timestamp=rrd_data["start_time"],
end_timestamp=rrd_data["end_time"],
start_timestamp=metrics_data["start_time"],
end_timestamp=metrics_data["end_time"],
bucket_seconds=bucket_seconds,
)
policy_by_bucket = {
@@ -3435,7 +3445,7 @@ class APIEndpoints:
for row in policy_rows
}
chart_data = []
for ts in rrd_data["timestamps"]:
for ts in metrics_data["timestamps"]:
bucket_ts = int(ts / bucket_seconds) * bucket_seconds
chart_data.append([ts * 1000, policy_by_bucket.get(bucket_ts, 0)])
@@ -3448,14 +3458,19 @@ class APIEndpoints:
)
continue
if metric_key in rrd_data["metrics"]:
if metric_key in metrics_data["metrics"]:
if metric_key in counter_metrics:
chart_data = self._process_counter_data(
rrd_data["metrics"][metric_key], timestamps_ms
)
if metrics_data.get("counter_mode") == "bucket_count":
chart_data = self._process_gauge_data(
metrics_data["metrics"][metric_key], timestamps_ms
)
else:
chart_data = self._process_counter_data(
metrics_data["metrics"][metric_key], timestamps_ms
)
else:
chart_data = self._process_gauge_data(
rrd_data["metrics"][metric_key], timestamps_ms
metrics_data["metrics"][metric_key], timestamps_ms
)
series.append(
@@ -3467,11 +3482,12 @@ class APIEndpoints:
)
graph_data = {
"start_time": rrd_data["start_time"],
"end_time": rrd_data["end_time"],
"step": rrd_data["step"],
"timestamps": rrd_data["timestamps"],
"start_time": metrics_data["start_time"],
"end_time": metrics_data["end_time"],
"step": metrics_data["step"],
"timestamps": metrics_data["timestamps"],
"series": series,
"data_source": metrics_data.get("data_source", "rrd"),
}
return self._success(graph_data)
@@ -7559,9 +7575,12 @@ class APIEndpoints:
storage = self._get_storage()
stats = storage.sqlite_handler.get_table_stats()
# Add RRD file size if it exists
rrd_path = storage.sqlite_handler.storage_dir / "metrics.rrd"
stats["rrd_size_bytes"] = rrd_path.stat().st_size if rrd_path.exists() else 0
rrd_exists = rrd_path.exists()
stats["rrd_enabled"] = bool(getattr(storage, "rrd_enabled", True))
stats["rrd_available"] = bool(getattr(storage, "rrd_handler", None)) and rrd_exists
stats["metrics_data_source"] = "rrd" if stats["rrd_available"] else "sqlite"
stats["rrd_size_bytes"] = rrd_path.stat().st_size if rrd_exists else 0
return {"success": True, "data": stats}
except Exception as e:
+1 -1
View File
@@ -176,7 +176,7 @@ class LogBuffer(logging.Handler):
# Global log buffer instance
_log_buffer = LogBuffer(max_lines=1000)
_log_buffer = LogBuffer(max_lines=300)
class DocEndpoint:
+125 -2
View File
@@ -882,19 +882,73 @@ def test_db_stats_options_success_and_error(cherrypy_ctx, tmp_path):
get_table_stats=MagicMock(return_value={"packets": {"rows": 10}}),
storage_dir=tmp_path,
)
_attach_storage(api, SimpleNamespace(sqlite_handler=sqlite_handler))
_attach_storage(
api,
SimpleNamespace(
sqlite_handler=sqlite_handler,
rrd_enabled=True,
rrd_handler=object(),
),
)
result = api.db_stats()
assert result["success"] is True
assert result["data"]["packets"]["rows"] == 10
assert result["data"]["rrd_enabled"] is True
assert result["data"]["rrd_available"] is True
assert result["data"]["metrics_data_source"] == "rrd"
assert result["data"]["rrd_size_bytes"] == 6
sqlite_handler.storage_dir = tmp_path / "missing"
sqlite_handler.storage_dir.mkdir()
_attach_storage(
api,
SimpleNamespace(
sqlite_handler=sqlite_handler,
rrd_enabled=False,
rrd_handler=None,
),
)
missing_rrd = api.db_stats()
assert missing_rrd["success"] is True
assert missing_rrd["data"]["rrd_enabled"] is False
assert missing_rrd["data"]["rrd_available"] is False
assert missing_rrd["data"]["metrics_data_source"] == "sqlite"
assert missing_rrd["data"]["rrd_size_bytes"] == 0
sqlite_handler.get_table_stats.side_effect = RuntimeError("stats failed")
err = api.db_stats()
assert err["success"] is False
assert "stats failed" in err["error"]
def test_rrd_data_endpoint_succeeds_with_sqlite_metrics(cherrypy_ctx):
del cherrypy_ctx
api = _make_api()
storage = SimpleNamespace(
get_metrics_data=MagicMock(
return_value={
"start_time": 0,
"end_time": 60,
"step": 60,
"timestamps": [0, 60],
"data_sources": ["rx_count"],
"packet_types": {"type_0": [0, 0]},
"metrics": {"rx_count": [1, 2]},
"data_source": "sqlite",
"counter_mode": "bucket_count",
}
)
)
_attach_storage(api, storage)
out = api.rrd_data()
assert out["success"] is True
assert out["data"]["data_source"] == "sqlite"
def test_db_purge_validation_and_results(cherrypy_ctx):
request, _ = cherrypy_ctx
request.method = "POST"
@@ -1791,7 +1845,7 @@ def test_metrics_graph_data_includes_policy_events(cherrypy_ctx):
del cherrypy_ctx
api = _make_api()
storage = SimpleNamespace(
get_rrd_data=MagicMock(
get_metrics_data=MagicMock(
return_value={
"start_time": 100,
"end_time": 220,
@@ -1821,6 +1875,75 @@ def test_metrics_graph_data_includes_policy_events(cherrypy_ctx):
assert "policy_events" in series_by_type
assert series_by_type["policy_events"]["name"] == "Policy Events"
assert series_by_type["policy_events"]["data"] == [[100000, 1], [160000, 3], [220000, 2]]
assert out["data"]["data_source"] == "rrd"
def test_metrics_graph_data_uses_sqlite_bucket_counts_without_counter_delta(cherrypy_ctx):
del cherrypy_ctx
api = _make_api()
storage = SimpleNamespace(
get_metrics_data=MagicMock(
return_value={
"start_time": 0,
"end_time": 60,
"step": 60,
"timestamps": [0, 60],
"metrics": {
"rx_count": [2, 4],
"avg_rssi": [None, -80.0],
},
"packet_types": {},
"data_source": "sqlite",
"counter_mode": "bucket_count",
}
),
get_policy_event_counts=MagicMock(return_value=[]),
)
_attach_storage(api, storage)
out = api.metrics_graph_data(hours="24", resolution="average", metrics="rx_count,avg_rssi")
assert out["success"] is True
series_by_type = {item["type"]: item for item in out["data"]["series"]}
assert series_by_type["rx_count"]["data"] == [[0, 2], [60000, 4]]
assert series_by_type["avg_rssi"]["data"] == [[0, 0], [60000, -80.0]]
assert out["data"]["data_source"] == "sqlite"
def test_metrics_graph_data_empty_sqlite_payload_is_successful(cherrypy_ctx):
del cherrypy_ctx
api = _make_api()
storage = SimpleNamespace(
get_metrics_data=MagicMock(
return_value={
"start_time": 0,
"end_time": 0,
"step": 60,
"timestamps": [0],
"metrics": {
"rx_count": [0],
"tx_count": [0],
"drop_count": [0],
"avg_rssi": [None],
"avg_snr": [None],
"avg_length": [None],
"avg_score": [None],
"neighbor_count": [None],
},
"packet_types": {},
"data_source": "sqlite",
"counter_mode": "bucket_count",
}
),
get_policy_event_counts=MagicMock(return_value=[]),
)
_attach_storage(api, storage)
out = api.metrics_graph_data(hours="24", resolution="average", metrics="rx_count")
assert out["success"] is True
assert out["data"]["series"][0]["data"] == [[0, 0]]
assert out["data"]["data_source"] == "sqlite"
def test_lbt_diagnostics_aligns_with_rrd_and_returns_correlations(cherrypy_ctx):
+138
View File
@@ -702,3 +702,141 @@ def test_cleanup_old_data_accepts_companion_events_days(tmp_path):
with h._connect() as conn:
count = conn.execute("SELECT COUNT(*) FROM packets").fetchone()[0]
assert count == 0
def test_get_metrics_data_returns_aligned_buckets_and_equal_length_arrays(tmp_path):
h = _make_handler(tmp_path)
packets = [
{
"timestamp": 61,
"type": 0,
"route": 1,
"length": 10,
"rssi": -80,
"snr": 5.0,
"score": 0.2,
"transmitted": False,
"packet_hash": "metrics-1",
},
{
"timestamp": 62,
"type": 19,
"route": 1,
"length": 20,
"rssi": -90,
"snr": 3.0,
"score": 0.4,
"transmitted": True,
"packet_hash": "metrics-2",
},
{
"timestamp": 121,
"type": 5,
"route": 1,
"length": 30,
"rssi": -70,
"snr": 10.0,
"score": 0.8,
"transmitted": True,
"packet_hash": "metrics-3",
},
]
for record in packets:
h.store_packet(record)
out = h.get_metrics_data(start_time=61, end_time=239, resolution="average")
assert out["data_source"] == "sqlite"
assert out["counter_mode"] == "bucket_count"
assert out["start_time"] == 60
assert out["end_time"] == 180
assert out["step"] == 60
assert out["timestamps"] == [60, 120, 180]
expected_length = len(out["timestamps"])
for values in out["metrics"].values():
assert len(values) == expected_length
for values in out["packet_types"].values():
assert len(values) == expected_length
assert out["metrics"]["rx_count"] == [2, 1, 0]
assert out["metrics"]["tx_count"] == [1, 1, 0]
assert out["metrics"]["drop_count"] == [1, 0, 0]
assert out["metrics"]["avg_length"] == [15.0, 30.0, None]
assert out["metrics"]["neighbor_count"] == [None, None, None]
assert out["packet_types"]["type_0"] == [1, 0, 0]
assert out["packet_types"]["type_5"] == [0, 1, 0]
assert out["packet_types"]["type_other"] == [1, 0, 0]
assert out["packet_types"]["type_15"] == [0, 0, 0]
def test_get_metrics_data_empty_buckets_use_zero_counters_and_null_gauges(tmp_path):
h = _make_handler(tmp_path)
out = h.get_metrics_data(start_time=0, end_time=120, resolution="average")
assert out["timestamps"] == [0, 60, 120]
assert out["metrics"]["rx_count"] == [0, 0, 0]
assert out["metrics"]["tx_count"] == [0, 0, 0]
assert out["metrics"]["drop_count"] == [0, 0, 0]
assert out["metrics"]["avg_rssi"] == [None, None, None]
assert out["metrics"]["avg_snr"] == [None, None, None]
assert out["metrics"]["avg_length"] == [None, None, None]
assert out["metrics"]["avg_score"] == [None, None, None]
assert out["packet_types"]["type_0"] == [0, 0, 0]
assert out["packet_types"]["type_other"] == [0, 0, 0]
@pytest.mark.parametrize(
("resolution", "expected_rssi", "expected_snr", "expected_length", "expected_score"),
[
("average", -75.0, 4.0, 15.0, 0.5),
("max", -70, 6.0, 20, 0.8),
("min", -80, 2.0, 10, 0.2),
],
)
def test_get_metrics_data_applies_requested_gauge_aggregation(
tmp_path,
resolution,
expected_rssi,
expected_snr,
expected_length,
expected_score,
):
h = _make_handler(tmp_path)
h.store_packet(
{
"timestamp": 61,
"type": 1,
"route": 1,
"length": 10,
"rssi": -80,
"snr": 2.0,
"score": 0.2,
"transmitted": False,
"packet_hash": f"agg-{resolution}-1",
}
)
h.store_packet(
{
"timestamp": 62,
"type": 1,
"route": 1,
"length": 20,
"rssi": -70,
"snr": 6.0,
"score": 0.8,
"transmitted": True,
"packet_hash": f"agg-{resolution}-2",
}
)
out = h.get_metrics_data(start_time=60, end_time=119, resolution=resolution)
assert out["metrics"]["avg_rssi"] == [expected_rssi]
assert out["metrics"]["avg_snr"] == [expected_snr]
assert out["metrics"]["avg_length"] == [expected_length]
assert out["metrics"]["avg_score"] == [expected_score]
@@ -1,6 +1,9 @@
import sys
import threading
import types
import builtins
import importlib.util
from pathlib import Path
from types import SimpleNamespace
from unittest.mock import MagicMock, patch
@@ -51,6 +54,12 @@ def _make_collector() -> StorageCollector:
return collector
def _stop_collector_threads(collector: StorageCollector) -> None:
collector._stats_stop_event.set()
if collector._stats_thread is not None:
collector._stats_thread.join(timeout=1)
def test_publish_packet_sync_broadcasts_packet_event_not_stats():
# The per-packet path must stay fast: it broadcasts the packet event but never
# runs the heavy aggregate or the stats broadcast (those moved to the loop).
@@ -98,3 +107,181 @@ def test_stats_loop_skips_when_no_clients():
collector._stats_broadcast_loop()
assert collector._broadcast_stats_once.call_count == 0
def test_rrd_defaults_to_enabled_when_metrics_config_missing():
rrd_candidate = MagicMock()
rrd_candidate.available = True
rrd_candidate.rrd_path.exists.return_value = True
with (
patch("repeater.data_acquisition.storage_collector.SQLiteHandler"),
patch(
"repeater.data_acquisition.storage_collector.RRDToolHandler",
return_value=rrd_candidate,
) as rrd_cls,
patch("repeater.data_acquisition.hardware_stats.HardwareStatsCollector"),
):
collector = StorageCollector(
config={"storage": {"storage_dir": "/tmp/openhop_repeater_test"}}
)
_stop_collector_threads(collector)
assert collector.rrd_enabled is True
rrd_cls.assert_called_once()
assert collector.rrd_handler is rrd_candidate
def test_rrd_disabled_does_not_instantiate_handler():
with (
patch("repeater.data_acquisition.storage_collector.SQLiteHandler"),
patch("repeater.data_acquisition.storage_collector.RRDToolHandler") as rrd_cls,
patch("repeater.data_acquisition.hardware_stats.HardwareStatsCollector"),
):
collector = StorageCollector(
config={
"storage": {"storage_dir": "/tmp/openhop_repeater_test"},
"metrics": {"rrd_enabled": False},
}
)
_stop_collector_threads(collector)
assert collector.rrd_enabled is False
assert collector.rrd_handler is None
rrd_cls.assert_not_called()
def test_record_packet_skips_rrd_work_when_disabled():
collector = _make_collector()
collector.rrd_enabled = False
collector.rrd_handler = None
collector.sqlite_handler.store_packet.return_value = 123
collector.sqlite_handler.get_cumulative_counts = MagicMock()
collector._publish_packet_sync = MagicMock()
collector._record_packet_blocking({"type": 1, "transmitted": False}, skip_mqtt=False)
collector.sqlite_handler.store_packet.assert_called_once()
collector.sqlite_handler.get_cumulative_counts.assert_not_called()
collector._publish_packet_sync.assert_called_once()
def test_requested_but_unavailable_rrd_falls_back_to_sqlite():
rrd_candidate = MagicMock()
rrd_candidate.available = False
rrd_candidate.rrd_path.exists.return_value = False
with (
patch("repeater.data_acquisition.storage_collector.SQLiteHandler"),
patch(
"repeater.data_acquisition.storage_collector.RRDToolHandler",
return_value=rrd_candidate,
),
patch("repeater.data_acquisition.hardware_stats.HardwareStatsCollector"),
):
collector = StorageCollector(
config={"storage": {"storage_dir": "/tmp/openhop_repeater_test"}}
)
_stop_collector_threads(collector)
collector.sqlite_handler.get_metrics_data.return_value = {
"start_time": 0,
"end_time": 0,
"step": 60,
"timestamps": [],
"metrics": {},
"packet_types": {},
"data_source": "sqlite",
"counter_mode": "bucket_count",
}
result = collector.get_metrics_data()
assert collector.rrd_handler is None
collector.sqlite_handler.get_metrics_data.assert_called_once_with(None, None, "average")
assert result["data_source"] == "sqlite"
def test_metrics_data_falls_back_when_rrd_read_raises():
collector = _make_collector()
collector.rrd_handler = MagicMock()
collector.rrd_handler.get_data.side_effect = RuntimeError("boom")
collector.sqlite_handler.get_metrics_data.return_value = {
"start_time": 0,
"end_time": 60,
"step": 60,
"timestamps": [0, 60],
"metrics": {"rx_count": [1, 2]},
"packet_types": {},
"data_source": "sqlite",
"counter_mode": "bucket_count",
}
result = collector.get_metrics_data(start_time=0, end_time=60, resolution="average")
collector.sqlite_handler.get_metrics_data.assert_called_once_with(0, 60, "average")
assert result["data_source"] == "sqlite"
def test_metrics_data_prefers_healthy_rrd_result():
collector = _make_collector()
collector.rrd_handler = MagicMock()
collector.rrd_handler.get_data.return_value = {
"start_time": 0,
"end_time": 60,
"step": 60,
"timestamps": [0, 60],
"metrics": {"rx_count": [1, 2]},
}
result = collector.get_metrics_data(start_time=0, end_time=60, resolution="average")
collector.sqlite_handler.get_metrics_data.assert_not_called()
assert result["data_source"] == "rrd"
def test_metrics_data_falls_back_when_rrd_result_is_malformed():
collector = _make_collector()
collector.rrd_handler = MagicMock()
collector.rrd_handler.get_data.return_value = {"start_time": 0, "end_time": 60}
collector.sqlite_handler.get_metrics_data.return_value = {
"start_time": 0,
"end_time": 60,
"step": 60,
"timestamps": [0, 60],
"metrics": {"rx_count": [0, 1]},
"packet_types": {},
"data_source": "sqlite",
"counter_mode": "bucket_count",
}
result = collector.get_metrics_data(start_time=0, end_time=60, resolution="average")
collector.sqlite_handler.get_metrics_data.assert_called_once_with(0, 60, "average")
assert result["data_source"] == "sqlite"
def test_rrdtool_handler_import_succeeds_without_dependency(monkeypatch):
module_path = (
Path(__file__).resolve().parents[1] / "repeater" / "data_acquisition" / "rrdtool_handler.py"
)
real_import = builtins.__import__
def fake_import(name, globals=None, locals=None, fromlist=(), level=0):
if name == "rrdtool":
raise ImportError("missing rrdtool")
return real_import(name, globals, locals, fromlist, level)
monkeypatch.setattr(builtins, "__import__", fake_import)
spec = importlib.util.spec_from_file_location("rrdtool_handler_missing_dep", module_path)
module = importlib.util.module_from_spec(spec)
assert spec is not None
assert spec.loader is not None
spec.loader.exec_module(module)
assert module.RRDTOOL_AVAILABLE is False