feat(mqtt): add device health telemetry and migrate to v2 topic spec

Two efforts that landed together because the v2 topic work extends
tables the health-telemetry effort added days earlier in the same
files/functions, making them impractical to separate cleanly:

Health/diagnostics telemetry (schema, Jul 13-17):
- New Postgres tables: device_alert_events, device_boot_events,
  device_ping_samples, device_diagnostics_reports, plus a `source`
  column on device_logs to distinguish log origins
- Query/service layer in pg_mqtt.py and database/__init__.py for
  inserting and listing this history, plus a "latest metrics" endpoint
  combining most-recent diagnostics + ping RTT per device
- mqtt/router.py gains list endpoints for alert/boot/ping/diagnostics
  history, consumed by the upcoming Health tab

MQTT v2 topic migration (Sep 21):
- Heartbeat payload flattened per vesper_mqtt_topic_spec_v2.md, adding
  rssi/free_heap/state/ok fields
- Command replies move to control/ack, device-initiated events to
  control/reports; mqtt/client.py subscribes to the new topic set and
  runs a ping_loop (wired up in main.py) for RTT sampling
- mqtt/logger.py and pg_mqtt.py updated to parse and persist the new
  payload shape alongside the legacy fields for backwards compatibility

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
2026-09-21 18:24:36 +03:00
co-authored by Claude Sonnet 5
parent e5d556fee1
commit 7c533b9245
12 changed files with 1456 additions and 55 deletions
+50
View File
@@ -16,6 +16,31 @@ from database.pg_mqtt import (
upsert_alert,
delete_alert,
get_alerts,
insert_alert_event,
get_latest_alert_events,
get_latest_alert_event,
get_alert_events,
insert_boot_event,
get_boot_events,
get_latest_boot_event,
insert_ping_sample,
get_ping_samples,
get_latest_ping_samples,
insert_diagnostics_report,
get_diagnostics_reports,
get_latest_diagnostics_report,
get_latest_diagnostics_reports,
insert_report,
get_reports,
delete_device_logs,
delete_device_heartbeats,
delete_device_commands,
delete_device_alert_events,
delete_device_current_alerts,
delete_device_boot_events,
delete_device_ping_samples,
delete_device_diagnostics_reports,
delete_device_reports,
partition_manager_loop,
ensure_current_partitions,
)
@@ -42,6 +67,31 @@ __all__ = [
"upsert_alert",
"delete_alert",
"get_alerts",
"insert_alert_event",
"get_latest_alert_events",
"get_latest_alert_event",
"get_alert_events",
"insert_boot_event",
"get_boot_events",
"get_latest_boot_event",
"insert_ping_sample",
"get_ping_samples",
"get_latest_ping_samples",
"insert_diagnostics_report",
"get_diagnostics_reports",
"get_latest_diagnostics_report",
"get_latest_diagnostics_reports",
"insert_report",
"get_reports",
"delete_device_logs",
"delete_device_heartbeats",
"delete_device_commands",
"delete_device_alert_events",
"delete_device_current_alerts",
"delete_device_boot_events",
"delete_device_ping_samples",
"delete_device_diagnostics_reports",
"delete_device_reports",
"partition_manager_loop",
"ensure_current_partitions",
]
+662 -26
View File
@@ -26,15 +26,15 @@ logger = logging.getLogger("database.pg_mqtt")
# ---------------------------------------------------------------------------
async def insert_log(device_serial: str, level: str, message: str,
device_timestamp: int | None = None) -> int:
device_timestamp: int | None = None, source: str = "log") -> int:
async with AsyncSessionLocal() as session:
result = await session.execute(
text("""
INSERT INTO device_logs (device_serial, level, message, device_timestamp, received_at)
VALUES (:serial, :level, :message, :ts, now())
INSERT INTO device_logs (device_serial, level, message, device_timestamp, received_at, source)
VALUES (:serial, :level, :message, :ts, now(), :source)
RETURNING id
"""),
{"serial": device_serial, "level": level, "message": message, "ts": device_timestamp},
{"serial": device_serial, "level": level, "message": message, "ts": device_timestamp, "source": source},
)
row = result.fetchone()
await session.commit()
@@ -43,15 +43,17 @@ async def insert_log(device_serial: str, level: str, message: str,
async def insert_heartbeat(device_serial: str, device_id: str,
firmware_version: str, ip_address: str,
gateway: str, uptime_ms: int, uptime_display: str) -> int:
gateway: str, uptime_ms: int, uptime_display: str,
rssi: int | None = None, free_heap: int | None = None,
state: str | None = None, ok: bool | None = None) -> int:
async with AsyncSessionLocal() as session:
result = await session.execute(
text("""
INSERT INTO heartbeats
(device_serial, device_id, firmware_version, ip_address,
gateway, uptime_ms, uptime_display, received_at)
gateway, uptime_ms, uptime_display, rssi, free_heap, state, ok, received_at)
VALUES
(:serial, :device_id, :fw, :ip, :gw, :uptime_ms, :uptime_display, now())
(:serial, :device_id, :fw, :ip, :gw, :uptime_ms, :uptime_display, :rssi, :free_heap, :state, :ok, now())
RETURNING id
"""),
{
@@ -62,6 +64,10 @@ async def insert_heartbeat(device_serial: str, device_id: str,
"gw": gateway,
"uptime_ms": uptime_ms,
"uptime_display": uptime_display,
"rssi": rssi,
"free_heap": free_heap,
"state": state,
"ok": ok,
},
)
row = result.fetchone()
@@ -113,18 +119,48 @@ async def update_command_response(command_id: int, status: str,
# Query operations
# ---------------------------------------------------------------------------
LEVEL_ORDER = ["INFO", "WARN", "ERROR"]
async def get_logs(device_serial: str, level: str | None = None,
search: str | None = None,
limit: int = 100, offset: int = 0) -> tuple[list, int]:
search: str | None = None, source: str | list[str] | None = None,
min_level: bool = False,
limit: int = 100, offset: int = 0,
since: datetime | None = None, until: datetime | None = None) -> tuple[list, int]:
"""
level: exact level match (legacy behaviour).
min_level: if True, `level` is treated as a floor — also returns higher-severity
levels (e.g. min_level=True, level='WARN' returns WARN and ERROR).
source: single source string, or a list of sources (e.g. ['log', 'info']) to
include several channels in one query — used by the unified feed.
since/until: same per-device time-range filter as the other history
endpoints (heartbeats, boot events, etc.) — lets the console's
TimeRangeSelect scope the Logs sub-tab too.
"""
where = "device_serial = :serial"
params: dict = {"serial": device_serial, "limit": limit, "offset": offset}
if level:
where += " AND level = :level"
params["level"] = level
if min_level and level in LEVEL_ORDER:
allowed = LEVEL_ORDER[LEVEL_ORDER.index(level):]
where += " AND level = ANY(:levels)"
params["levels"] = allowed
else:
where += " AND level = :level"
params["level"] = level
if search:
where += " AND message ILIKE :search"
params["search"] = f"%{search}%"
if source:
sources = [source] if isinstance(source, str) else list(source)
where += " AND source = ANY(:sources)"
params["sources"] = sources
if since is not None:
where += " AND received_at >= :since"
params["since"] = since
if until is not None:
where += " AND received_at <= :until"
params["until"] = until
async with AsyncSessionLocal() as session:
count_result = await session.execute(
@@ -134,7 +170,7 @@ async def get_logs(device_serial: str, level: str | None = None,
rows_result = await session.execute(
text(f"""
SELECT id, device_serial, level, message, device_timestamp,
SELECT id, device_serial, level, message, device_timestamp, source,
received_at AT TIME ZONE 'UTC' AS received_at
FROM device_logs
WHERE {where}
@@ -148,26 +184,47 @@ async def get_logs(device_serial: str, level: str | None = None,
return [_row_to_dict(r) for r in rows], total
async def get_heartbeats(device_serial: str, limit: int = 100,
offset: int = 0) -> tuple[list, int]:
def _time_range_clause(column: str, since: datetime | None, until: datetime | None,
params: dict) -> str:
"""Appends `since`/`until` bind params to `params` in place and returns an
` AND {column} >= :since AND {column} <= :until`-style SQL fragment (only
the bounds actually provided). Shared by every history/telemetry query
that the console's per-device time-range filter needs to scope — see
TimeRangeSelect on the frontend."""
clause = ""
if since is not None:
clause += f" AND {column} >= :since"
params["since"] = since
if until is not None:
clause += f" AND {column} <= :until"
params["until"] = until
return clause
async def get_heartbeats(device_serial: str, limit: int = 100, offset: int = 0,
since: datetime | None = None,
until: datetime | None = None) -> tuple[list, int]:
async with AsyncSessionLocal() as session:
params: dict = {"serial": device_serial}
range_clause = _time_range_clause("received_at", since, until, params)
count_result = await session.execute(
text("SELECT COUNT(*) FROM heartbeats WHERE device_serial = :serial"),
{"serial": device_serial},
text(f"SELECT COUNT(*) FROM heartbeats WHERE device_serial = :serial{range_clause}"),
params,
)
total = count_result.scalar()
rows_result = await session.execute(
text("""
text(f"""
SELECT id, device_serial, device_id, firmware_version, ip_address,
gateway, uptime_ms, uptime_display,
gateway, uptime_ms, uptime_display, rssi, free_heap, state, ok,
received_at AT TIME ZONE 'UTC' AS received_at
FROM heartbeats
WHERE device_serial = :serial
WHERE device_serial = :serial{range_clause}
ORDER BY received_at DESC
LIMIT :limit OFFSET :offset
"""),
{"serial": device_serial, "limit": limit, "offset": offset},
{**params, "limit": limit, "offset": offset},
)
rows = rows_result.mappings().all()
@@ -207,7 +264,7 @@ async def get_latest_heartbeats() -> list:
text("""
SELECT DISTINCT ON (device_serial)
id, device_serial, device_id, firmware_version, ip_address,
gateway, uptime_ms, uptime_display,
gateway, uptime_ms, uptime_display, rssi, free_heap, state, ok,
received_at AT TIME ZONE 'UTC' AS received_at
FROM heartbeats
ORDER BY device_serial, received_at DESC
@@ -286,6 +343,548 @@ async def get_alerts(device_serial: str) -> list:
return [_row_to_dict(r) for r in rows]
# ---------------------------------------------------------------------------
# Device alert events (insert-only history — WARNING/CRITICAL/FAILED only,
# CLEARED is not logged here). Used to answer "when was the most recent
# issue on this device", surviving past the alert being resolved.
# ---------------------------------------------------------------------------
async def insert_alert_event(device_serial: str, subsystem: str, state: str,
message: str | None = None):
async with AsyncSessionLocal() as session:
await session.execute(
text("""
INSERT INTO device_alert_events (device_serial, subsystem, state, message, occurred_at)
VALUES (:serial, :subsystem, :state, :message, now())
"""),
{"serial": device_serial, "subsystem": subsystem, "state": state, "message": message},
)
await session.commit()
async def get_latest_alert_events() -> list:
"""Most recent alert event per device, across all devices — one row each."""
async with AsyncSessionLocal() as session:
result = await session.execute(
text("""
SELECT DISTINCT ON (device_serial)
id, device_serial, subsystem, state, message,
occurred_at AT TIME ZONE 'UTC' AS occurred_at
FROM device_alert_events
ORDER BY device_serial, occurred_at DESC
""")
)
rows = result.mappings().all()
return [_row_to_dict(r) for r in rows]
async def get_alert_events(device_serial: str, limit: int = 100,
offset: int = 0) -> tuple[list, int]:
"""Full alert transition history for one device (WARNING/CRITICAL/FAILED — no CLEARED)."""
async with AsyncSessionLocal() as session:
count_result = await session.execute(
text("SELECT COUNT(*) FROM device_alert_events WHERE device_serial = :serial"),
{"serial": device_serial},
)
total = count_result.scalar()
rows_result = await session.execute(
text("""
SELECT id, device_serial, subsystem, state, message,
occurred_at AT TIME ZONE 'UTC' AS occurred_at
FROM device_alert_events
WHERE device_serial = :serial
ORDER BY occurred_at DESC
LIMIT :limit OFFSET :offset
"""),
{"serial": device_serial, "limit": limit, "offset": offset},
)
rows = rows_result.mappings().all()
return [_row_to_dict(r) for r in rows], total
async def get_latest_alert_event(device_serial: str) -> dict | None:
async with AsyncSessionLocal() as session:
result = await session.execute(
text("""
SELECT id, device_serial, subsystem, state, message,
occurred_at AT TIME ZONE 'UTC' AS occurred_at
FROM device_alert_events
WHERE device_serial = :serial
ORDER BY occurred_at DESC
LIMIT 1
"""),
{"serial": device_serial},
)
row = result.mappings().first()
return _row_to_dict(row) if row else None
# ---------------------------------------------------------------------------
# Device boot events (insert-only history from firmware's boot_report — see
# mqtt/logger.py::_handle_info). One row per boot, fault or not. Powers the
# Health tab's boot timeline + restart chart.
# ---------------------------------------------------------------------------
async def insert_boot_event(device_serial: str, boot_count: int | None,
reset_reason: str | None, is_fault: bool,
free_heap: int | None,
crash_task: str | None = None,
crash_pc: int | None = None,
crash_exc_cause: int | None = None,
crash_exc_vaddr: int | None = None) -> int:
async with AsyncSessionLocal() as session:
result = await session.execute(
text("""
INSERT INTO device_boot_events
(device_serial, boot_count, reset_reason, is_fault, free_heap,
crash_task, crash_pc, crash_exc_cause, crash_exc_vaddr, occurred_at)
VALUES
(:serial, :boot_count, :reset_reason, :is_fault, :free_heap,
:crash_task, :crash_pc, :crash_exc_cause, :crash_exc_vaddr, now())
RETURNING id
"""),
{
"serial": device_serial,
"boot_count": boot_count,
"reset_reason": reset_reason,
"is_fault": is_fault,
"free_heap": free_heap,
"crash_task": crash_task,
"crash_pc": crash_pc,
"crash_exc_cause": crash_exc_cause,
"crash_exc_vaddr": crash_exc_vaddr,
},
)
row = result.fetchone()
await session.commit()
return row[0]
async def get_boot_events(device_serial: str, limit: int = 100, offset: int = 0,
since: datetime | None = None,
until: datetime | None = None) -> tuple[list, int]:
async with AsyncSessionLocal() as session:
params: dict = {"serial": device_serial}
range_clause = _time_range_clause("occurred_at", since, until, params)
count_result = await session.execute(
text(f"SELECT COUNT(*) FROM device_boot_events WHERE device_serial = :serial{range_clause}"),
params,
)
total = count_result.scalar()
rows_result = await session.execute(
text(f"""
SELECT id, device_serial, boot_count, reset_reason, is_fault, free_heap,
crash_task, crash_pc, crash_exc_cause, crash_exc_vaddr,
occurred_at AT TIME ZONE 'UTC' AS occurred_at
FROM device_boot_events
WHERE device_serial = :serial{range_clause}
ORDER BY occurred_at DESC
LIMIT :limit OFFSET :offset
"""),
{**params, "limit": limit, "offset": offset},
)
rows = rows_result.mappings().all()
return [_row_to_dict(r) for r in rows], total
async def get_latest_boot_event(device_serial: str) -> dict | None:
async with AsyncSessionLocal() as session:
result = await session.execute(
text("""
SELECT id, device_serial, boot_count, reset_reason, is_fault, free_heap,
crash_task, crash_pc, crash_exc_cause, crash_exc_vaddr,
occurred_at AT TIME ZONE 'UTC' AS occurred_at
FROM device_boot_events
WHERE device_serial = :serial
ORDER BY occurred_at DESC
LIMIT 1
"""),
{"serial": device_serial},
)
row = result.mappings().first()
return _row_to_dict(row) if row else None
# ---------------------------------------------------------------------------
# Device ping samples — backend-computed RTT, one row per ping/pong round
# trip. Sampler loop lives in mqtt/client.py; see MqttManager.ping_loop().
# ---------------------------------------------------------------------------
async def insert_ping_sample(device_serial: str, rtt_ms: int) -> int:
async with AsyncSessionLocal() as session:
result = await session.execute(
text("""
INSERT INTO device_ping_samples (device_serial, rtt_ms, sampled_at)
VALUES (:serial, :rtt_ms, now())
RETURNING id
"""),
{"serial": device_serial, "rtt_ms": rtt_ms},
)
row = result.fetchone()
await session.commit()
return row[0]
async def get_ping_samples(device_serial: str, limit: int = 200, offset: int = 0,
since: datetime | None = None,
until: datetime | None = None) -> tuple[list, int]:
async with AsyncSessionLocal() as session:
params: dict = {"serial": device_serial}
range_clause = _time_range_clause("sampled_at", since, until, params)
count_result = await session.execute(
text(f"SELECT COUNT(*) FROM device_ping_samples WHERE device_serial = :serial{range_clause}"),
params,
)
total = count_result.scalar()
rows_result = await session.execute(
text(f"""
SELECT id, device_serial, rtt_ms,
sampled_at AT TIME ZONE 'UTC' AS sampled_at
FROM device_ping_samples
WHERE device_serial = :serial{range_clause}
ORDER BY sampled_at DESC
LIMIT :limit OFFSET :offset
"""),
{**params, "limit": limit, "offset": offset},
)
rows = rows_result.mappings().all()
return [_row_to_dict(r) for r in rows], total
async def get_latest_ping_samples() -> list:
# Fleet-wide "last known" ping RTT for DeviceList — same DISTINCT ON
# pattern as get_latest_heartbeats(); devices with no RTC never appear.
async with AsyncSessionLocal() as session:
rows_result = await session.execute(
text("""
SELECT DISTINCT ON (device_serial)
id, device_serial, rtt_ms,
sampled_at AT TIME ZONE 'UTC' AS sampled_at
FROM device_ping_samples
ORDER BY device_serial, sampled_at DESC
""")
)
rows = rows_result.mappings().all()
return [_row_to_dict(r) for r in rows]
# ---------------------------------------------------------------------------
# Device diagnostics reports — firmware's periodic (5-min) diagnostics_report
# event: CPU temperature, WiFi reconnects, OTA state, stack high-water marks.
# See project-vesper's feature catalog F-056 and CommunicationRouter README
# for the full design rationale (why this is separate from the 30s heartbeat).
# ---------------------------------------------------------------------------
async def insert_diagnostics_report(
device_serial: str,
cpu_temp_avg: float | None = None,
cpu_temp_min: float | None = None,
cpu_temp_max: float | None = None,
cpu_temp_samples: int | None = None,
wifi_reconnect_count: int | None = None,
wifi_last_disconnect_reason: str | None = None,
wifi_last_disconnect_uptime_ms: int | None = None,
ota_current_version: str | None = None,
ota_update_available: bool | None = None,
ota_available_version: str | None = None,
ota_last_check_uptime_ms: int | None = None,
ota_last_error: str | None = None,
stack_high_water: dict | None = None,
bell_strikes: dict | None = None,
bell_loads: dict | None = None,
cooling_active: bool | None = None,
) -> int:
async with AsyncSessionLocal() as session:
result = await session.execute(
text("""
INSERT INTO device_diagnostics_reports
(device_serial,
cpu_temp_avg, cpu_temp_min, cpu_temp_max, cpu_temp_samples,
wifi_reconnect_count, wifi_last_disconnect_reason, wifi_last_disconnect_uptime_ms,
ota_current_version, ota_update_available, ota_available_version,
ota_last_check_uptime_ms, ota_last_error,
stack_high_water, bell_strikes, bell_loads, cooling_active, received_at)
VALUES
(:serial,
:cpu_temp_avg, :cpu_temp_min, :cpu_temp_max, :cpu_temp_samples,
:wifi_reconnect_count, :wifi_last_disconnect_reason, :wifi_last_disconnect_uptime_ms,
:ota_current_version, :ota_update_available, :ota_available_version,
:ota_last_check_uptime_ms, :ota_last_error,
:stack_high_water, :bell_strikes, :bell_loads, :cooling_active, now())
RETURNING id
"""),
{
"serial": device_serial,
"cpu_temp_avg": cpu_temp_avg,
"cpu_temp_min": cpu_temp_min,
"cpu_temp_max": cpu_temp_max,
"cpu_temp_samples": cpu_temp_samples,
"wifi_reconnect_count": wifi_reconnect_count,
"wifi_last_disconnect_reason": wifi_last_disconnect_reason,
"wifi_last_disconnect_uptime_ms": wifi_last_disconnect_uptime_ms,
"ota_current_version": ota_current_version,
"ota_update_available": ota_update_available,
"ota_available_version": ota_available_version,
"ota_last_check_uptime_ms": ota_last_check_uptime_ms,
"ota_last_error": ota_last_error,
"stack_high_water": json.dumps(stack_high_water) if stack_high_water else None,
"bell_strikes": json.dumps(bell_strikes) if bell_strikes else None,
"bell_loads": json.dumps(bell_loads) if bell_loads else None,
"cooling_active": cooling_active,
},
)
row = result.fetchone()
await session.commit()
return row[0]
async def get_diagnostics_reports(device_serial: str, limit: int = 200, offset: int = 0,
since: datetime | None = None,
until: datetime | None = None) -> tuple[list, int]:
async with AsyncSessionLocal() as session:
params: dict = {"serial": device_serial}
range_clause = _time_range_clause("received_at", since, until, params)
count_result = await session.execute(
text(f"SELECT COUNT(*) FROM device_diagnostics_reports WHERE device_serial = :serial{range_clause}"),
params,
)
total = count_result.scalar()
rows_result = await session.execute(
text(f"""
SELECT id, device_serial,
cpu_temp_avg, cpu_temp_min, cpu_temp_max, cpu_temp_samples,
wifi_reconnect_count, wifi_last_disconnect_reason, wifi_last_disconnect_uptime_ms,
ota_current_version, ota_update_available, ota_available_version,
ota_last_check_uptime_ms, ota_last_error,
stack_high_water, bell_strikes, bell_loads, cooling_active,
received_at AT TIME ZONE 'UTC' AS received_at
FROM device_diagnostics_reports
WHERE device_serial = :serial{range_clause}
ORDER BY received_at DESC
LIMIT :limit OFFSET :offset
"""),
{**params, "limit": limit, "offset": offset},
)
rows = rows_result.mappings().all()
return [_row_to_dict(r) for r in rows], total
async def get_latest_diagnostics_report(device_serial: str) -> dict | None:
async with AsyncSessionLocal() as session:
result = await session.execute(
text("""
SELECT id, device_serial,
cpu_temp_avg, cpu_temp_min, cpu_temp_max, cpu_temp_samples,
wifi_reconnect_count, wifi_last_disconnect_reason, wifi_last_disconnect_uptime_ms,
ota_current_version, ota_update_available, ota_available_version,
ota_last_check_uptime_ms, ota_last_error,
stack_high_water, bell_strikes, bell_loads, cooling_active,
received_at AT TIME ZONE 'UTC' AS received_at
FROM device_diagnostics_reports
WHERE device_serial = :serial
ORDER BY received_at DESC
LIMIT 1
"""),
{"serial": device_serial},
)
row = result.mappings().first()
return _row_to_dict(row) if row else None
async def get_latest_diagnostics_reports() -> list:
# Fleet-wide "last known" CPU temp for DeviceList — one DISTINCT ON query
# instead of N per-device round trips, same pattern as get_latest_heartbeats().
async with AsyncSessionLocal() as session:
rows_result = await session.execute(
text("""
SELECT DISTINCT ON (device_serial)
id, device_serial,
cpu_temp_avg, cpu_temp_min, cpu_temp_max, cpu_temp_samples,
received_at AT TIME ZONE 'UTC' AS received_at
FROM device_diagnostics_reports
ORDER BY device_serial, received_at DESC
""")
)
rows = rows_result.mappings().all()
return [_row_to_dict(r) for r in rows]
# ---------------------------------------------------------------------------
# Per-device stat reset — clears QA/bench test data for a single device
# before shipping to a customer. Deliberately narrow: only tables that are
# pure accumulated history/telemetry are here. device_alerts (live current
# state) is included but callers should treat it as opt-in/advanced, not a
# default — clearing it could hide a real active fault. Nothing here touches
# device identity, customer assignment, notes, warranty, or audit logs — see
# backend/devices/service.py for the separate device_stats (Firestore) reset.
# ---------------------------------------------------------------------------
async def delete_device_logs(device_serial: str) -> int:
async with AsyncSessionLocal() as session:
result = await session.execute(
text("DELETE FROM device_logs WHERE device_serial = :serial"),
{"serial": device_serial},
)
await session.commit()
return result.rowcount
async def delete_device_heartbeats(device_serial: str) -> int:
async with AsyncSessionLocal() as session:
result = await session.execute(
text("DELETE FROM heartbeats WHERE device_serial = :serial"),
{"serial": device_serial},
)
await session.commit()
return result.rowcount
async def delete_device_commands(device_serial: str) -> int:
async with AsyncSessionLocal() as session:
result = await session.execute(
text("DELETE FROM commands WHERE device_serial = :serial"),
{"serial": device_serial},
)
await session.commit()
return result.rowcount
async def delete_device_alert_events(device_serial: str) -> int:
async with AsyncSessionLocal() as session:
result = await session.execute(
text("DELETE FROM device_alert_events WHERE device_serial = :serial"),
{"serial": device_serial},
)
await session.commit()
return result.rowcount
async def delete_device_current_alerts(device_serial: str) -> int:
"""Clears LIVE alert state (device_alerts), not history. Opt-in only —
see module note above. If the device currently has a real active fault,
clearing this hides it from the console until the next state change."""
async with AsyncSessionLocal() as session:
result = await session.execute(
text("DELETE FROM device_alerts WHERE device_serial = :serial"),
{"serial": device_serial},
)
await session.commit()
return result.rowcount
async def delete_device_boot_events(device_serial: str) -> int:
async with AsyncSessionLocal() as session:
result = await session.execute(
text("DELETE FROM device_boot_events WHERE device_serial = :serial"),
{"serial": device_serial},
)
await session.commit()
return result.rowcount
async def delete_device_ping_samples(device_serial: str) -> int:
async with AsyncSessionLocal() as session:
result = await session.execute(
text("DELETE FROM device_ping_samples WHERE device_serial = :serial"),
{"serial": device_serial},
)
await session.commit()
return result.rowcount
async def delete_device_diagnostics_reports(device_serial: str) -> int:
async with AsyncSessionLocal() as session:
result = await session.execute(
text("DELETE FROM device_diagnostics_reports WHERE device_serial = :serial"),
{"serial": device_serial},
)
await session.commit()
return result.rowcount
# ---------------------------------------------------------------------------
# Device reports — critical, unsolicited, board-initiated events from the new
# control/reports MQTT topic (currently only bell_overload). Console-side
# storage is deliberately light: insert + list for historical/audit purposes.
# The physical tablets, not this console, are the real-time consumer of
# control/reports. See project-vesper's vesper_mqtt_topic_spec_v2.md.
# ---------------------------------------------------------------------------
async def insert_report(device_serial: str, report_type: str, payload: dict | None) -> int:
async with AsyncSessionLocal() as session:
result = await session.execute(
text("""
INSERT INTO device_reports (device_serial, report_type, payload, occurred_at)
VALUES (:serial, :report_type, :payload, now())
RETURNING id
"""),
{
"serial": device_serial,
"report_type": report_type,
"payload": json.dumps(payload) if payload else None,
},
)
row = result.fetchone()
await session.commit()
return row[0]
async def get_reports(device_serial: str, limit: int = 200, offset: int = 0,
since: datetime | None = None,
until: datetime | None = None) -> tuple[list, int]:
async with AsyncSessionLocal() as session:
params: dict = {"serial": device_serial}
range_clause = _time_range_clause("occurred_at", since, until, params)
count_result = await session.execute(
text(f"SELECT COUNT(*) FROM device_reports WHERE device_serial = :serial{range_clause}"),
params,
)
total = count_result.scalar()
rows_result = await session.execute(
text(f"""
SELECT id, device_serial, report_type, payload,
occurred_at AT TIME ZONE 'UTC' AS occurred_at
FROM device_reports
WHERE device_serial = :serial{range_clause}
ORDER BY occurred_at DESC
LIMIT :limit OFFSET :offset
"""),
{**params, "limit": limit, "offset": offset},
)
rows = rows_result.mappings().all()
return [_row_to_dict(r) for r in rows], total
async def delete_device_reports(device_serial: str) -> int:
async with AsyncSessionLocal() as session:
result = await session.execute(
text("DELETE FROM device_reports WHERE device_serial = :serial"),
{"serial": device_serial},
)
await session.commit()
return result.rowcount
# ---------------------------------------------------------------------------
# Partition management
# ---------------------------------------------------------------------------
@@ -314,8 +913,20 @@ async def ensure_current_partitions():
logger.info("Partition check complete")
async def drop_old_partitions(keep_months: int = 6):
"""Drop device_logs partitions older than keep_months."""
def _global_retention_days() -> int:
"""Read the admin-configured global log retention (days), falling back to
the env-configured default if Firestore is unreachable or unset."""
try:
from settings.log_retention_service import get_log_retention
return get_log_retention().days
except Exception:
return settings.mqtt_data_retention_days
async def drop_old_partitions(keep_months: int | None = None):
"""Drop device_logs partitions older than keep_months (default: global retention setting, in months)."""
if keep_months is None:
keep_months = max(1, round(_global_retention_days() / 30))
cutoff = _add_months(date.today().replace(day=1), -keep_months)
async with AsyncSessionLocal() as session:
result = await session.execute(text("""
@@ -361,7 +972,7 @@ async def partition_manager_loop():
# ---------------------------------------------------------------------------
async def purge_old_data(retention_days: int | None = None):
days = retention_days or settings.mqtt_data_retention_days
days = retention_days or _global_retention_days()
cutoff = datetime.now(timezone.utc) - timedelta(days=days)
async with AsyncSessionLocal() as session:
await session.execute(
@@ -372,8 +983,20 @@ async def purge_old_data(retention_days: int | None = None):
text("DELETE FROM commands WHERE sent_at < :cutoff"),
{"cutoff": cutoff},
)
await session.execute(
text("DELETE FROM device_ping_samples WHERE sampled_at < :cutoff"),
{"cutoff": cutoff},
)
await session.execute(
text("DELETE FROM device_diagnostics_reports WHERE received_at < :cutoff"),
{"cutoff": cutoff},
)
await session.execute(
text("DELETE FROM device_reports WHERE occurred_at < :cutoff"),
{"cutoff": cutoff},
)
await session.commit()
logger.info(f"Purged heartbeats and commands older than {days} days")
logger.info(f"Purged heartbeats, commands, ping samples, diagnostics reports, and device reports older than {days} days")
async def purge_loop():
@@ -403,9 +1026,22 @@ async def close_db():
# ---------------------------------------------------------------------------
def _row_to_dict(row) -> dict:
"""Convert a SQLAlchemy RowMapping to a plain dict with ISO string timestamps."""
"""Convert a SQLAlchemy RowMapping to a plain dict with ISO string timestamps.
Every timestamp column here is selected via `AT TIME ZONE 'UTC'`, which hands
back a naive (tzinfo-less) datetime that IS UTC but doesn't say so. Without
attaching tzinfo before .isoformat(), the resulting string has no offset
suffix (e.g. "2026-07-14T16:03:21") — browsers then parse that as LOCAL
time, not UTC, silently shifting every displayed time by the browser's
UTC offset (3h for Greek summer/EEST, 2h for winter/EET). Attaching
timezone.utc makes isoformat() emit the explicit "+00:00" suffix so
`new Date(...)` on the frontend converts to the viewer's local time
correctly year-round, including across DST transitions.
"""
d = dict(row)
for key, val in d.items():
if isinstance(val, datetime):
if val.tzinfo is None:
val = val.replace(tzinfo=timezone.utc)
d[key] = val.isoformat()
return d