Files
bellsystems-cp/backend/database/pg_mqtt.py
T
bonaminandClaude Opus 5.5 98dd16b597 fix(mqtt): stop logging retained boot_report/alerts as new events
The firmware publishes boot_report on system/info and alerts on
system/alerts with retain=true. Every time the backend (re)connects - which
under uvicorn --reload means every backend file save - the broker redelivers
the last retained message and we inserted it again with occurred_at=now().
Result: the Health tab showed fresh PANIC boots and "Device reset due to
fault" alerts for a device that had been up for 4 days.

- insert_boot_event skips the insert when a row with the same
  (device_serial, boot_count) already exists. boot_count is the firmware's
  lifetime counter, so it uniquely identifies a boot.
- upsert_alert only writes when state/message actually changed and returns
  whether it did; the alert-event history row is only added on a change.
  A redelivered identical alert no longer bumps updated_at either.

Existing duplicate rows are not touched by this commit.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-30 15:48:48 +03:00

1073 lines
44 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) -> bool:
"""Insert or update the current alert. Returns False (and leaves the row,
including updated_at, untouched) when it already holds this exact state and
message — which is what a retained system/alerts message looks like when the
broker redelivers it on every backend (re)subscribe."""
async with AsyncSessionLocal() as session:
result = 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
WHERE device_alerts.state IS DISTINCT FROM EXCLUDED.state
OR device_alerts.message IS DISTINCT FROM EXCLUDED.message
RETURNING id
"""),
{"serial": device_serial, "subsystem": subsystem, "state": state, "message": message},
)
changed = result.fetchone() is not None
await session.commit()
return changed
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 | None:
"""Insert a boot event, unless one with the same boot_count already exists
for this device (returns None then). boot_report is published retained on
system/info, so the broker redelivers the last one every time the backend
(re)subscribes — without this check each backend restart logged the
device's last boot (crash detail included) again as a brand-new event."""
async with AsyncSessionLocal() as session:
if boot_count is not None:
dup = await session.execute(
text("""
SELECT 1 FROM device_boot_events
WHERE device_serial = :serial AND boot_count = :boot_count
LIMIT 1
"""),
{"serial": device_serial, "boot_count": boot_count},
)
if dup.first() is not None:
return None
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] if row else None
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