Files
bellsystems-cp/backend/database/pg_mqtt.py
T
bonaminandClaude Sonnet 5 7c533b9245 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>
2026-09-21 18:24:36 +03:00

1048 lines
42 KiB
Python

"""
Phase 5 — MQTT live data functions backed by Postgres.
device_logs is a partitioned table; heartbeats and commands are plain tables.
All three are accessed via raw SQL (not ORM) because device_logs partitioning
does not play well with SQLAlchemy's declarative ORM.
device_alerts is an ORM model (devices/orm.py) and is handled here via raw SQL
to keep a single consistent interface for callers that used to import from database.core.
"""
import asyncio
import json
import logging
from datetime import date, datetime, timedelta, timezone
from sqlalchemy import text
from config import settings
from database.postgres import AsyncSessionLocal
logger = logging.getLogger("database.pg_mqtt")
# ---------------------------------------------------------------------------
# Insert operations
# ---------------------------------------------------------------------------
async def insert_log(device_serial: str, level: str, message: str,
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, source)
VALUES (:serial, :level, :message, :ts, now(), :source)
RETURNING id
"""),
{"serial": device_serial, "level": level, "message": message, "ts": device_timestamp, "source": source},
)
row = result.fetchone()
await session.commit()
return row[0]
async def insert_heartbeat(device_serial: str, device_id: str,
firmware_version: str, ip_address: str,
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, rssi, free_heap, state, ok, received_at)
VALUES
(:serial, :device_id, :fw, :ip, :gw, :uptime_ms, :uptime_display, :rssi, :free_heap, :state, :ok, now())
RETURNING id
"""),
{
"serial": device_serial,
"device_id": device_id,
"fw": firmware_version,
"ip": ip_address,
"gw": gateway,
"uptime_ms": uptime_ms,
"uptime_display": uptime_display,
"rssi": rssi,
"free_heap": free_heap,
"state": state,
"ok": ok,
},
)
row = result.fetchone()
await session.commit()
return row[0]
async def insert_command(device_serial: str, command_name: str,
command_payload: dict) -> int:
async with AsyncSessionLocal() as session:
result = await session.execute(
text("""
INSERT INTO commands (device_serial, command_name, command_payload, sent_at)
VALUES (:serial, :name, :payload, now())
RETURNING id
"""),
{
"serial": device_serial,
"name": command_name,
"payload": json.dumps(command_payload),
},
)
row = result.fetchone()
await session.commit()
return row[0]
async def update_command_response(command_id: int, status: str,
response_payload: dict | None = None):
async with AsyncSessionLocal() as session:
await session.execute(
text("""
UPDATE commands
SET status = :status,
response_payload = :payload,
responded_at = now()
WHERE id = :id
"""),
{
"id": command_id,
"status": status,
"payload": json.dumps(response_payload) if response_payload else None,
},
)
await session.commit()
# ---------------------------------------------------------------------------
# Query operations
# ---------------------------------------------------------------------------
LEVEL_ORDER = ["INFO", "WARN", "ERROR"]
async def get_logs(device_serial: str, level: str | None = None,
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:
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(
text(f"SELECT COUNT(*) FROM device_logs WHERE {where}"), params
)
total = count_result.scalar()
rows_result = await session.execute(
text(f"""
SELECT id, device_serial, level, message, device_timestamp, source,
received_at AT TIME ZONE 'UTC' AS received_at
FROM device_logs
WHERE {where}
ORDER BY received_at DESC
LIMIT :limit OFFSET :offset
"""),
params,
)
rows = rows_result.mappings().all()
return [_row_to_dict(r) for r in rows], total
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(f"SELECT COUNT(*) FROM heartbeats WHERE device_serial = :serial{range_clause}"),
params,
)
total = count_result.scalar()
rows_result = await session.execute(
text(f"""
SELECT id, device_serial, device_id, firmware_version, ip_address,
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{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_commands(device_serial: str, limit: int = 100,
offset: int = 0) -> tuple[list, int]:
async with AsyncSessionLocal() as session:
count_result = await session.execute(
text("SELECT COUNT(*) FROM commands WHERE device_serial = :serial"),
{"serial": device_serial},
)
total = count_result.scalar()
rows_result = await session.execute(
text("""
SELECT id, device_serial, command_name, command_payload, status,
response_payload,
sent_at AT TIME ZONE 'UTC' AS sent_at,
responded_at AT TIME ZONE 'UTC' AS responded_at
FROM commands
WHERE device_serial = :serial
ORDER BY sent_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_heartbeats() -> list:
async with AsyncSessionLocal() as session:
rows_result = await session.execute(
text("""
SELECT DISTINCT ON (device_serial)
id, device_serial, device_id, firmware_version, ip_address,
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
""")
)
rows = rows_result.mappings().all()
return [_row_to_dict(r) for r in rows]
async def get_pending_command(device_serial: str) -> dict | None:
async with AsyncSessionLocal() as session:
result = await session.execute(
text("""
SELECT id, device_serial, command_name, command_payload, status,
response_payload,
sent_at AT TIME ZONE 'UTC' AS sent_at,
responded_at AT TIME ZONE 'UTC' AS responded_at
FROM commands
WHERE device_serial = :serial AND status = 'pending'
ORDER BY sent_at DESC
LIMIT 1
"""),
{"serial": device_serial},
)
row = result.mappings().fetchone()
return _row_to_dict(row) if row else None
# ---------------------------------------------------------------------------
# Device alerts
# ---------------------------------------------------------------------------
async def upsert_alert(device_serial: str, subsystem: str, state: str,
message: str | None = None):
async with AsyncSessionLocal() as session:
await session.execute(
text("""
INSERT INTO device_alerts (device_serial, subsystem, state, message, updated_at)
VALUES (:serial, :subsystem, :state, :message, now())
ON CONFLICT (device_serial, subsystem)
DO UPDATE SET
state = EXCLUDED.state,
message = EXCLUDED.message,
updated_at = EXCLUDED.updated_at
"""),
{"serial": device_serial, "subsystem": subsystem, "state": state, "message": message},
)
await session.commit()
async def delete_alert(device_serial: str, subsystem: str):
async with AsyncSessionLocal() as session:
await session.execute(
text("DELETE FROM device_alerts WHERE device_serial = :serial AND subsystem = :subsystem"),
{"serial": device_serial, "subsystem": subsystem},
)
await session.commit()
async def get_alerts(device_serial: str) -> list:
async with AsyncSessionLocal() as session:
result = await session.execute(
text("""
SELECT id, device_serial, subsystem, state, message,
updated_at AT TIME ZONE 'UTC' AS updated_at
FROM device_alerts
WHERE device_serial = :serial
ORDER BY updated_at DESC
"""),
{"serial": device_serial},
)
rows = result.mappings().all()
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
# ---------------------------------------------------------------------------
def _add_months(d: date, months: int) -> date:
month = d.month - 1 + months
year = d.year + month // 12
month = month % 12 + 1
return d.replace(year=year, month=month, day=1)
async def ensure_current_partitions():
"""Create device_logs partitions for the current and next month if missing."""
async with AsyncSessionLocal() as session:
for month_offset in (0, 1):
d = _add_months(date.today().replace(day=1), month_offset)
partition_name = f"device_logs_{d.strftime('%Y_%m')}"
start = d.isoformat()
end = _add_months(d, 1).isoformat()
await session.execute(text(f"""
CREATE TABLE IF NOT EXISTS {partition_name}
PARTITION OF device_logs
FOR VALUES FROM ('{start}') TO ('{end}')
"""))
await session.commit()
logger.info("Partition check complete")
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("""
SELECT tablename FROM pg_tables
WHERE schemaname = 'public'
AND tablename LIKE 'device_logs_%'
"""))
partitions = [r[0] for r in result.fetchall()]
for name in partitions:
# name format: device_logs_YYYY_MM
parts = name.split("_")
if len(parts) != 4:
continue
try:
partition_date = date(int(parts[2]), int(parts[3]), 1)
except ValueError:
continue
if partition_date < cutoff:
async with AsyncSessionLocal() as session:
await session.execute(text(f"DROP TABLE IF EXISTS {name}"))
await session.commit()
logger.info(f"Dropped old partition: {name}")
async def partition_manager_loop():
"""Runs once on startup, then monthly thereafter."""
await ensure_current_partitions()
while True:
# Sleep ~30 days, wake up and ensure next month's partition exists
await asyncio.sleep(30 * 24 * 3600)
try:
await ensure_current_partitions()
await drop_old_partitions()
except Exception as e:
logger.error(f"Partition manager error: {e}")
# ---------------------------------------------------------------------------
# Cleanup (replaces SQLite purge_loop — now a no-op since Postgres uses
# partition drops instead of row-by-row deletes for device_logs; heartbeats
# and commands are still purged by row deletion)
# ---------------------------------------------------------------------------
async def purge_old_data(retention_days: int | None = None):
days = retention_days or _global_retention_days()
cutoff = datetime.now(timezone.utc) - timedelta(days=days)
async with AsyncSessionLocal() as session:
await session.execute(
text("DELETE FROM heartbeats WHERE received_at < :cutoff"),
{"cutoff": cutoff},
)
await session.execute(
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, commands, ping samples, diagnostics reports, and device reports older than {days} days")
async def purge_loop():
while True:
await asyncio.sleep(86400)
try:
await purge_old_data()
except Exception as e:
logger.error(f"Purge failed: {e}")
# ---------------------------------------------------------------------------
# Stub — no longer needed but kept so nothing that imports init_db/close_db breaks
# ---------------------------------------------------------------------------
async def init_db():
"""No-op: Postgres schema is managed by Alembic, not runtime init."""
logger.info("Postgres MQTT backend active — no SQLite init needed")
async def close_db():
"""No-op: SQLAlchemy engine lifecycle is managed by the process."""
# ---------------------------------------------------------------------------
# Helpers
# ---------------------------------------------------------------------------
def _row_to_dict(row) -> dict:
"""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