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
7 changes: 7 additions & 0 deletions docs/operations.md
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,13 @@ Useful web routes:
`mqtt_ingest_health` reports callback counts, latency, and contention. Slow
callback warnings identify only the MQTT topic and never include payloads.

Dashboard queries use indexed latest-packet and recent-trend lookups so older
history does not need to be scanned for each sensor refresh. Historical sensor
ID spelling and case-sensitive trend grouping are preserved. The two-second
dashboard JSON cache is checked before database work; its lifetime, inventory
cache lifetimes, and statistics cache lifetimes start after computation finishes.
These optimizations use existing SQLite indexes and require no schema migration.

Dashboard values normally refresh every 15 seconds. A failed response or a
request exceeding 12 seconds (including the JSON body) shows **Live updates
delayed** above the dashboard. Existing values remain visible while automatic
Expand Down
2 changes: 1 addition & 1 deletion sensorius/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,4 +4,4 @@
release notes, and supporting tooling can report a consistent build identity.
"""

__version__ = "v0.26.262.1"
__version__ = "v0.26.262.2"
50 changes: 30 additions & 20 deletions sensorius/saiDataLogger.py
Original file line number Diff line number Diff line change
Expand Up @@ -1994,21 +1994,24 @@ def get_latest_values(self, sensor_id):
return {}

def get_available_sensors(self):
"""Return distinct historical sensor IDs in case-sensitive sorted order."""
now_mono = time.monotonic()
cached = self._available_sensors_cache
if cached and cached[0] > now_mono:
return list(cached[1])
query = "SELECT DISTINCT sensor_id FROM readings ORDER BY sensor_id"
# Sort only the small distinct result, not the entire historical index.
query = "SELECT DISTINCT sensor_id FROM readings"
try:
with self._open_conn() as conn:
result = [row[0] for row in conn.execute(query).fetchall()]
self._available_sensors_cache = (now_mono + 5.0, list(result))
result = sorted(row[0] for row in conn.execute(query).fetchall())
self._available_sensors_cache = (time.monotonic() + 5.0, list(result))
return result
except Exception as e:
printDM(f"Sensor ID query error: {e}", location=MODULE)
return []

def get_available_metrics(self, sensor_id):
"""Return cached historical and derived metric names for a sensor."""
sensor_key = str(sensor_id or "").strip().lower()
now_mono = time.monotonic()
cached = self._available_metrics_cache.get(sensor_key) if sensor_key else None
Expand All @@ -2020,7 +2023,7 @@ def get_available_metrics(self, sensor_id):
for sid, metrics in cached_by_sensor[1].items():
if str(sid or "").strip().lower() == sensor_key:
result = list(metrics)
self._available_metrics_cache[sensor_key] = (now_mono + 5.0, result)
self._available_metrics_cache[sensor_key] = (time.monotonic() + 5.0, result)
return result

try:
Expand All @@ -2036,13 +2039,14 @@ def get_available_metrics(self, sensor_id):
result = [row[0] for row in rows if row and row[0]]
result = self._with_available_derived_metrics(result)
if sensor_key:
self._available_metrics_cache[sensor_key] = (now_mono + 5.0, list(result))
self._available_metrics_cache[sensor_key] = (time.monotonic() + 5.0, list(result))
return result
except Exception as e:
printDM(f"Error fetching metrics for {sensor_id}: {e}", location=MODULE)
return []

def get_available_metrics_by_sensor(self):
"""Return a cached inventory of historical and derived sensor metrics."""
now_mono = time.monotonic()
cached = self._available_metrics_by_sensor_cache
if cached and cached[0] > now_mono:
Expand Down Expand Up @@ -2073,7 +2077,7 @@ def get_available_metrics_by_sensor(self):
for sid, metrics in list(result.items()):
result[sid] = self._with_available_derived_metrics(metrics)

expires = now_mono + 5.0
expires = time.monotonic() + 5.0
self._available_metrics_by_sensor_cache = (
expires,
{sid: list(metrics) for sid, metrics in result.items()},
Expand Down Expand Up @@ -2132,17 +2136,19 @@ def get_latest_timestamps(self, sensor_ids: list[str]) -> dict[str, str]:
return result

sid_map = {sid.lower(): sid for sid in missing}
placeholders = ",".join("?" for _ in sid_map)
placeholders = ",".join("(?)" for _ in sid_map)
try:
with self._open_conn() as conn:
rows = conn.execute(
f"""
WITH latest AS (
SELECT sensor_id COLLATE NOCASE AS sid_l,
MAX(ts_epoch) AS latest_ts_epoch
FROM readings
WHERE sensor_id COLLATE NOCASE IN ({placeholders})
GROUP BY sensor_id COLLATE NOCASE
WITH requested(sid_l) AS (VALUES {placeholders}),
latest AS (
SELECT sid_l, (
SELECT ts_epoch FROM readings
WHERE sensor_id = requested.sid_l COLLATE NOCASE
ORDER BY ts_epoch DESC LIMIT 1
) AS latest_ts_epoch
FROM requested
)
SELECT r.sensor_id, r.timestamp
FROM readings r
Expand Down Expand Up @@ -2186,21 +2192,25 @@ def get_latest_values_and_timestamps(self, sensor_ids: list[str]) -> tuple[dict[
# The database may be written through another logger instance or process.
# Always reconcile the dashboard snapshot with persisted data so a populated
# RAM cache cannot remain stale while graph queries continue to advance.
# Seek each newest epoch through the existing sensor/time index, then
# retrieve every metric tied at that epoch as in the grouped MAX query.
sid_map = {sid.lower(): sid for sid in clean_ids}
placeholders = ",".join("?" for _ in sid_map)
placeholders = ",".join("(?)" for _ in sid_map)
db_values: dict[str, dict] = {}
db_timestamps: dict[str, str] = {}
try:
with self._open_conn() as conn:
cur = conn.cursor()
cur.execute(
f"""
WITH latest AS (
SELECT sensor_id COLLATE NOCASE AS sid_l,
MAX(ts_epoch) AS latest_ts_epoch
FROM readings
WHERE sensor_id COLLATE NOCASE IN ({placeholders})
GROUP BY sensor_id COLLATE NOCASE
WITH requested(sid_l) AS (VALUES {placeholders}),
latest AS (
SELECT sid_l, (
SELECT ts_epoch FROM readings
WHERE sensor_id = requested.sid_l COLLATE NOCASE
ORDER BY ts_epoch DESC LIMIT 1
) AS latest_ts_epoch
FROM requested
)
SELECT r.sensor_id, r.timestamp, r.metric, r.value
FROM readings r
Expand Down
10 changes: 7 additions & 3 deletions sensorius/saiStats.py
Original file line number Diff line number Diff line change
Expand Up @@ -134,6 +134,8 @@ def _metric_trends(
params.append(str(sensor_id))
params.append(pressure_window_s)

# Match the index collation for bounded seeks, then preserve the original
# case-sensitive grouping when historical IDs differ only by case.
rows = conn.execute(
f"""
WITH latest AS (
Expand All @@ -148,7 +150,8 @@ def _metric_trends(
SELECT r.sensor_id, r.metric, r.value, r.ts_epoch, latest.end_ts
FROM readings AS r
JOIN latest
ON r.sensor_id = latest.sensor_id
ON r.sensor_id = latest.sensor_id COLLATE NOCASE
AND r.sensor_id = latest.sensor_id COLLATE BINARY
WHERE r.value IS NOT NULL
AND r.ts_epoch >= latest.end_ts - ?
AND r.ts_epoch <= latest.end_ts
Expand Down Expand Up @@ -280,6 +283,7 @@ def _get_stats_for_range_impl(self, sensor_id, start_epoch: float, end_epoch: fl
return results

def get_24hr_stats(self, sensor_id):
"""Return cached 24-hour extrema, averages, and trends for a sensor."""
sid = str(sensor_id or "").strip()
now_mono = time.monotonic()
cached = self._stats_cache.get(sid)
Expand All @@ -294,7 +298,7 @@ def get_24hr_stats(self, sensor_id):
for metric, trend in trends.items():
if metric in result:
result[metric]["trend"] = trend
self._stats_cache[sid] = (now_mono + self._stats_cache_ttl_sec, dict(result))
self._stats_cache[sid] = (time.monotonic() + self._stats_cache_ttl_sec, dict(result))
return result

def get_all_stats_fast(self):
Expand Down Expand Up @@ -395,7 +399,7 @@ def _get_all_stats_fast_impl(self):
if metric in sensor_stats:
sensor_stats[metric]["trend"] = trend

self._all_stats_cache = (now_mono + self._stats_cache_ttl_sec, dict(results))
self._all_stats_cache = (time.monotonic() + self._stats_cache_ttl_sec, dict(results))
return results

def create_stats_router(settings, gc_mgr, data_logger=None):
Expand Down
37 changes: 18 additions & 19 deletions sensorius/saiWebRoutes.py
Original file line number Diff line number Diff line change
Expand Up @@ -2421,6 +2421,22 @@ async def current_data_page(
_phase_started = time.monotonic()
global _cdp_debug_last_log
dashboard_cache_key = str(sensor_id or "All")
if json_only:
cache_key = (str(sensor_id or "All"), 1 if include_extras else 0)
now_mono = time.monotonic()
cached_json = _DASHBOARD_JSON_CACHE.get(cache_key)
if cached_json and cached_json[0] > now_mono:
cached_payload = cached_json[1]
_ui_profile_log(
"dashboard",
_route_started,
json_only=1,
include_extras=int(bool(include_extras)),
cache=1,
sensor_id=(sensor_id or "All"),
)
return JSONResponse(cached_payload)

if dashboard_return and not json_only:
cached_dashboard = _DASHBOARD_HTML_CACHE.get(dashboard_cache_key)
if cached_dashboard is not None:
Expand Down Expand Up @@ -2891,7 +2907,7 @@ def _switch_has_renderable_channels(
lambda: {sid: resolve_location_for_sid(sid) for sid in available}
)
_DASHBOARD_INVENTORY_CACHE = (
now_mono + _DASHBOARD_INVENTORY_CACHE_TTL_SEC,
time.monotonic() + _DASHBOARD_INVENTORY_CACHE_TTL_SEC,
{
"local_ids": list(local_ids),
"available": list(available),
Expand Down Expand Up @@ -3470,23 +3486,6 @@ def _resolve_meas_status_for_sid(sid: str) -> str:
return "unknown"

if json_only:
cache_key = (str(sensor_id or "All"), 1 if include_extras else 0)
now_mono = time.monotonic()
cached_json = _DASHBOARD_JSON_CACHE.get(cache_key)
if cached_json and cached_json[0] > now_mono:
cached_payload = cached_json[1]
_ui_profile_log(
"dashboard",
_route_started,
json_only=1,
include_extras=int(bool(include_extras)),
cache=1,
sensor_id=(sensor_id or "All"),
sensors=len(available),
switches=len(available_switches),
)
return JSONResponse(cached_payload)

timestamps = await asyncio.to_thread(
lambda: {
sid: (bulk_timestamps.get(sid) or data_logger.get_latest_timestamp(sid) or "")
Expand Down Expand Up @@ -3562,7 +3561,7 @@ def _resolve_meas_status_for_sid(sid: str) -> str:

payload = _dashboard_json_safe(payload)
_DASHBOARD_JSON_CACHE[cache_key] = (
now_mono + _DASHBOARD_JSON_CACHE_TTL_SEC,
time.monotonic() + _DASHBOARD_JSON_CACHE_TTL_SEC,
payload,
)
_ui_profile_log(
Expand Down
109 changes: 109 additions & 0 deletions testApparatus/test_dashboard_query_optimization.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,109 @@
"""Exercise dashboard query compatibility and bounded historical lookups.

Large synthetic history is inserted directly only for query-work regression
checks; functional packets use the normal logger write path.
"""

import importlib
import sqlite3
from datetime import datetime, timedelta

import pytest

from sensorius.saiDataLogger import saiDataLogger
from sensorius.saiStats import saiStats


@pytest.fixture
def logger(tmp_path):
saiDataLogger._schema_ready = False
instance = saiDataLogger(db_path=str(tmp_path / "dashboard.db"))
yield instance
instance.close()
saiDataLogger._schema_ready = False


def test_latest_packet_preserves_case_lookup_ties_and_missing_sensors(logger):
stamp = datetime.now().replace(microsecond=0)
logger.log_readings(stamp.isoformat(), "Sensor-A", {"Temperature": 10.0})
latest = (stamp + timedelta(minutes=1)).isoformat()
logger.log_readings(latest, "Sensor-A", {"Temperature": 20.0})
logger.log_readings(latest, "sensor-a", {"CO2": 700.0})
logger.sensor_values.clear()
logger.sensor_timestamps.clear()

values, timestamps = logger.get_latest_values_and_timestamps(["SENSOR-A", "missing"])

assert values["SENSOR-A"]["Temperature"] == 20.0
assert values["SENSOR-A"]["CO2"] == 700.0
assert "missing" not in values
assert "missing" not in timestamps
logger.sensor_timestamps.clear()
assert logger.get_latest_timestamps(["SENSOR-A", "missing"]) == timestamps
assert logger.get_available_sensors() == ["Sensor-A", "sensor-a"]


def test_latest_packet_seeks_past_large_history(logger, monkeypatch):
logger.log_readings(datetime.now().isoformat(), "sensor-a", {"Temperature": 20.0})
# A direct historical seed avoids thousands of irrelevant writer callbacks.
with logger._writer_conn:
logger._writer_conn.executemany(
"INSERT INTO readings(timestamp, ts_epoch, sensor_id, metric, value) VALUES (?, ?, ?, ?, ?)",
[("2000-01-01", float(i), "sensor-a", "Temperature", 1.0) for i in range(5000)],
)
logger.sensor_values.clear()
logger.sensor_timestamps.clear()
original_open = logger._open_conn
steps = []

def tracked_open():
conn = original_open()
conn.set_progress_handler(lambda: steps.append(1) or 0, 100)
return conn

monkeypatch.setattr(logger, "_open_conn", tracked_open)
values, _ = logger.get_latest_values_and_timestamps(["sensor-a"])
assert values["sensor-a"]["Temperature"] == 20.0
assert len(steps) < 20, "Latest packet retrieval scanned historical readings"


def test_trends_keep_case_distinct_histories_and_use_bounded_index(logger):
now = datetime.now().replace(microsecond=0)
for index in range(6):
stamp = (now - timedelta(minutes=5-index)).isoformat()
logger.log_readings(stamp, "Sensor-A", {"Temperature": 10.0 + index})
logger.log_readings(stamp, "sensor-a", {"Temperature": 30.0 - index})
# Old rows should not be visited by the recent-data join.
with logger._writer_conn:
logger._writer_conn.executemany(
"INSERT INTO readings(timestamp, ts_epoch, sensor_id, metric, value) VALUES (?, ?, ?, ?, ?)",
[("2000-01-01", float(i), "Sensor-A", "Temperature", 1.0) for i in range(5000)],
)
statter = saiStats(logger.db_path)
steps = []
with sqlite3.connect(logger.db_path) as conn:
conn.set_progress_handler(lambda: steps.append(1) or 0, 100)
trends = statter._metric_trends(conn)
assert trends["Sensor-A"]["Temperature"]["samples"] == 6
assert trends["sensor-a"]["Temperature"]["samples"] == 6
assert trends["Sensor-A"]["Temperature"]["rate_per_hour"] > 0
assert trends["sensor-a"]["Temperature"]["rate_per_hour"] < 0
assert len(steps) < 30, "Trend join scanned historical readings"


def test_stats_cache_lifetime_starts_after_computation(logger, monkeypatch):
module = importlib.import_module("sensorius.saiStats")
statter = saiStats(logger.db_path)
clock = [100.0]
monkeypatch.setattr(module.time, "monotonic", lambda: clock[0])
calls = []

def slow_trends(*args, **kwargs):
calls.append(1)
clock[0] += 10.0
return {}

monkeypatch.setattr(statter, "_metric_trends", slow_trends)
assert statter.get_all_stats_fast() == {}
assert statter.get_all_stats_fast() == {}
assert len(calls) == 1
31 changes: 31 additions & 0 deletions testApparatus/test_nodus_settings_schema_writes.py
Original file line number Diff line number Diff line change
Expand Up @@ -4106,6 +4106,37 @@ def _render_dashboard(*_args, **_kwargs):
assert render_calls == [1, 2]


@pytest.mark.asyncio
async def test_dashboard_json_cache_skips_queries_and_expires_after_completion(tmp_path, monkeypatch):
app, _ingest, _system_root, _sensor_root, _switch_root = await _build_app(tmp_path, monkeypatch)
monkeypatch.setattr(saiWebRoutes, "_DASHBOARD_INVENTORY_CACHE", None)
monkeypatch.setattr(saiWebRoutes, "_DASHBOARD_JSON_CACHE", {})
monkeypatch.setattr(saiWebRoutes.data_logger, "get_available_sensors", lambda: [])
monkeypatch.setattr(saiWebRoutes.data_logger, "get_switch_identities", lambda: [])
monkeypatch.setattr(saiWebRoutes.statter, "get_all_stats_fast", lambda: {})
original_clock = time.monotonic
offset = [0.0]
monkeypatch.setattr(time, "monotonic", lambda: original_clock() + offset[0])
calls = []

def slow_inventory():
calls.append(1)
offset[0] += 10.0
return []

monkeypatch.setattr(saiWebRoutes.data_logger, "get_available_sensors", slow_inventory)
async with AsyncClient(transport=ASGITransport(app=app), base_url="http://test") as client:
first = await client.get("/", params={"json_only": "true"})
second = await client.get("/", params={"json_only": "true", "sensor_id": "All"})
assert first.status_code == second.status_code == 200
assert second.json() == first.json()
assert len(calls) == 1
offset[0] += saiWebRoutes._DASHBOARD_JSON_CACHE_TTL_SEC + 1.0
expired = await client.get("/", params={"json_only": "true"})
assert expired.status_code == 200
assert len(calls) == 2


def test_dashboard_warming_payloads_are_not_stable_shell_content():
assert saiWebRoutes._dashboard_payload_is_warming({"reason": "warming"}) is True
assert saiWebRoutes._dashboard_payload_is_warming({"cache_status": "warming"}) is True
Expand Down
Loading