Firmware F-070 adds a fuller `crash` object (backtrace, backtrace_corrupted, elf_sha256) and a new `pre_crash` snapshot (uptime, heap stats, optional abort_msg) to boot_report and to telemetry.get_boot_history entries. - Migration c9d0e1f2a3b4: device_boot_events gains crash / pre_crash JSONB, stored verbatim so later additive fields need no migration. The legacy crash_* columns are still filled; old rows and old firmware are unchanged. JSON is bound as CAST(CAST(:x AS TEXT) AS JSONB): a bare JSONB cast makes asyncpg JSON-encode the already-encoded string a second time. - boot_report ingestion passes both objects through. - telemetry.get_boot_history replies (control/ack) are merged into boot history whoever sent the command: an entry matches an existing row on boot_count + reset_reason + time within 30 min (boot_count alone is not unique - the counter gets reset), and only fills crash/pre_crash the row lacks; unmatched entries are boots we never saw live and are inserted at the device's timestamp; entries without ts are skipped. - GET /api/mqtt/crash-groups: fault boots grouped by abort_msg with hex addresses stripped, else task + exception cause, else reset reason. When abort_msg is present, pc/exc_cause describe abort() itself and are ignored for grouping. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
1242 lines
51 KiB
Python
1242 lines
51 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.
|
|
# ---------------------------------------------------------------------------
|
|
|
|
_BOOT_COLUMNS = """
|
|
id, device_serial, boot_count, reset_reason, is_fault, free_heap,
|
|
crash_task, crash_pc, crash_exc_cause, crash_exc_vaddr,
|
|
crash::text AS crash, pre_crash::text AS pre_crash,
|
|
occurred_at AT TIME ZONE 'UTC' AS occurred_at
|
|
"""
|
|
|
|
# How far apart a device-reported boot timestamp (telemetry.get_boot_history
|
|
# "ts", device clock) and our own occurred_at (broker receive time of the live
|
|
# boot_report) may be and still count as the same boot. The live report goes
|
|
# out on first MQTT connect, so it normally lands within a minute of boot.
|
|
_BOOT_MATCH_WINDOW = timedelta(minutes=30)
|
|
|
|
|
|
def _json_or_none(obj: dict | None) -> str | None:
|
|
"""Encode for a CAST(CAST(:x AS TEXT) AS JSONB) bind. Passed as text on
|
|
purpose: a bare CAST(:x AS JSONB) makes asyncpg type the parameter as jsonb
|
|
and JSON-encode our already-encoded string a second time."""
|
|
return json.dumps(obj) if obj else None
|
|
|
|
|
|
def _boot_row_to_dict(row) -> dict:
|
|
d = _row_to_dict(row)
|
|
for key in ("crash", "pre_crash"):
|
|
val = d.get(key)
|
|
if isinstance(val, str):
|
|
try:
|
|
d[key] = json.loads(val)
|
|
except ValueError:
|
|
d[key] = None
|
|
return d
|
|
|
|
|
|
def _crash_columns(crash: dict | None) -> dict:
|
|
"""The legacy flat crash_* columns, still filled so older readers and the
|
|
fleet grouping query keep working without unpacking JSON."""
|
|
crash = crash or {}
|
|
return {
|
|
"crash_task": crash.get("task"),
|
|
"crash_pc": crash.get("pc"),
|
|
"crash_exc_cause": crash.get("exc_cause"),
|
|
"crash_exc_vaddr": crash.get("exc_vaddr"),
|
|
}
|
|
|
|
|
|
async def insert_boot_event(device_serial: str, boot_count: int | None,
|
|
reset_reason: str | None, is_fault: bool,
|
|
free_heap: int | None,
|
|
crash: dict | None = None,
|
|
pre_crash: dict | None = None,
|
|
skip_if_latest: bool = False) -> int | None:
|
|
"""Insert a boot event. With skip_if_latest=True (a retained replay of the
|
|
device's last boot_report), skip it and return None when it matches the
|
|
device's most recent row: that boot is already recorded.
|
|
|
|
Only the LATEST row is compared, never "any row with this boot_count": the
|
|
firmware's lifetime counter gets reset (reflash / telemetry reset), so the
|
|
same boot_count legitimately appears again for a later, different boot.
|
|
|
|
`crash` / `pre_crash` are the firmware's objects stored verbatim (F-070);
|
|
both are omitted by older firmware and simply stay NULL."""
|
|
async with AsyncSessionLocal() as session:
|
|
if skip_if_latest and boot_count is not None:
|
|
latest = await session.execute(
|
|
text("""
|
|
SELECT boot_count, reset_reason FROM device_boot_events
|
|
WHERE device_serial = :serial
|
|
ORDER BY occurred_at DESC
|
|
LIMIT 1
|
|
"""),
|
|
{"serial": device_serial},
|
|
)
|
|
row = latest.first()
|
|
if row is not None and row.boot_count == boot_count and row.reset_reason == reset_reason:
|
|
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,
|
|
crash, pre_crash, occurred_at)
|
|
VALUES
|
|
(:serial, :boot_count, :reset_reason, :is_fault, :free_heap,
|
|
:crash_task, :crash_pc, :crash_exc_cause, :crash_exc_vaddr,
|
|
CAST(CAST(:crash AS TEXT) AS JSONB), CAST(CAST(:pre_crash AS TEXT) AS JSONB), now())
|
|
RETURNING id
|
|
"""),
|
|
{
|
|
"serial": device_serial,
|
|
"boot_count": boot_count,
|
|
"reset_reason": reset_reason,
|
|
"is_fault": is_fault,
|
|
"free_heap": free_heap,
|
|
**_crash_columns(crash),
|
|
"crash": _json_or_none(crash),
|
|
"pre_crash": _json_or_none(pre_crash),
|
|
},
|
|
)
|
|
row = result.fetchone()
|
|
await session.commit()
|
|
return row[0] if row else None
|
|
|
|
|
|
async def merge_device_boot_history(device_serial: str, boots: list[dict]) -> dict:
|
|
"""Merge the device's own SD boot log (telemetry.get_boot_history reply)
|
|
into device_boot_events.
|
|
|
|
Each entry is {"ts": epoch s, "boot": n, "reason": str, "crash"?, "pre_crash"?}.
|
|
An entry matches an existing row when boot_count and reset_reason are equal
|
|
and the times are within _BOOT_MATCH_WINDOW (boot_count alone is not unique,
|
|
see insert_boot_event). A match only fills in crash/pre_crash the row is
|
|
missing; nothing already stored is overwritten. An entry with no match is a
|
|
boot we never saw live and is inserted at the device's timestamp. Entries
|
|
without a usable ts can't be placed in time and are skipped.
|
|
"""
|
|
inserted = updated = skipped = 0
|
|
async with AsyncSessionLocal() as session:
|
|
for entry in boots:
|
|
if not isinstance(entry, dict):
|
|
skipped += 1
|
|
continue
|
|
ts = entry.get("ts")
|
|
boot_count = entry.get("boot")
|
|
reason = entry.get("reason")
|
|
if not isinstance(ts, (int, float)) or ts <= 0:
|
|
skipped += 1
|
|
continue
|
|
at = datetime.fromtimestamp(ts, tz=timezone.utc)
|
|
crash = entry.get("crash") if isinstance(entry.get("crash"), dict) else None
|
|
pre_crash = entry.get("pre_crash") if isinstance(entry.get("pre_crash"), dict) else None
|
|
|
|
match = (await session.execute(
|
|
text("""
|
|
SELECT id, crash IS NULL AS no_crash, pre_crash IS NULL AS no_pre_crash
|
|
FROM device_boot_events
|
|
WHERE device_serial = :serial
|
|
AND boot_count IS NOT DISTINCT FROM :boot_count
|
|
AND reset_reason IS NOT DISTINCT FROM :reason
|
|
AND occurred_at BETWEEN :lo AND :hi
|
|
ORDER BY abs(extract(epoch FROM occurred_at - :at))
|
|
LIMIT 1
|
|
"""),
|
|
{"serial": device_serial, "boot_count": boot_count, "reason": reason,
|
|
"lo": at - _BOOT_MATCH_WINDOW, "hi": at + _BOOT_MATCH_WINDOW, "at": at},
|
|
)).first()
|
|
|
|
if match is None:
|
|
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,
|
|
crash, pre_crash, occurred_at)
|
|
VALUES
|
|
(:serial, :boot_count, :reason, :is_fault, NULL,
|
|
:crash_task, :crash_pc, :crash_exc_cause, :crash_exc_vaddr,
|
|
CAST(CAST(:crash AS TEXT) AS JSONB), CAST(CAST(:pre_crash AS TEXT) AS JSONB), :at)
|
|
"""),
|
|
{
|
|
"serial": device_serial, "boot_count": boot_count, "reason": reason,
|
|
"is_fault": bool(entry.get("is_fault", reason in FAULT_RESET_REASONS)),
|
|
**_crash_columns(crash),
|
|
"crash": _json_or_none(crash), "pre_crash": _json_or_none(pre_crash),
|
|
"at": at,
|
|
},
|
|
)
|
|
inserted += 1
|
|
continue
|
|
|
|
fill_crash = crash is not None and match.no_crash
|
|
fill_pre = pre_crash is not None and match.no_pre_crash
|
|
if not (fill_crash or fill_pre):
|
|
continue
|
|
sets, params = [], {"id": match.id}
|
|
if fill_crash:
|
|
sets.append("crash = CAST(CAST(:crash AS TEXT) AS JSONB)")
|
|
sets.append("crash_task = COALESCE(crash_task, :crash_task)")
|
|
sets.append("crash_pc = COALESCE(crash_pc, :crash_pc)")
|
|
sets.append("crash_exc_cause = COALESCE(crash_exc_cause, :crash_exc_cause)")
|
|
sets.append("crash_exc_vaddr = COALESCE(crash_exc_vaddr, :crash_exc_vaddr)")
|
|
params.update(_crash_columns(crash))
|
|
params["crash"] = _json_or_none(crash)
|
|
if fill_pre:
|
|
sets.append("pre_crash = CAST(CAST(:pre_crash AS TEXT) AS JSONB)")
|
|
params["pre_crash"] = _json_or_none(pre_crash)
|
|
await session.execute(
|
|
text(f"UPDATE device_boot_events SET {', '.join(sets)} WHERE id = :id"),
|
|
params,
|
|
)
|
|
updated += 1
|
|
await session.commit()
|
|
return {"inserted": inserted, "updated": updated, "skipped": skipped}
|
|
|
|
|
|
# Mirrors the firmware's isFaultReset() — used only when a history entry
|
|
# doesn't say is_fault itself.
|
|
FAULT_RESET_REASONS = {"PANIC", "TASK_WATCHDOG", "INTERRUPT_WATCHDOG", "OTHER_WATCHDOG", "BROWNOUT"}
|
|
|
|
|
|
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 {_BOOT_COLUMNS}
|
|
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 [_boot_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(f"""
|
|
SELECT {_BOOT_COLUMNS}
|
|
FROM device_boot_events
|
|
WHERE device_serial = :serial
|
|
ORDER BY occurred_at DESC
|
|
LIMIT 1
|
|
"""),
|
|
{"serial": device_serial},
|
|
)
|
|
row = result.mappings().first()
|
|
|
|
return _boot_row_to_dict(row) if row else None
|
|
|
|
|
|
async def get_crash_events(since: datetime | None = None,
|
|
until: datetime | None = None,
|
|
limit: int = 5000) -> list:
|
|
"""Fleet-wide fault boots that carry any crash detail (coredump summary
|
|
and/or pre-crash snapshot), newest first. Grouped in mqtt/crash_groups.py."""
|
|
async with AsyncSessionLocal() as session:
|
|
params: dict = {"limit": limit}
|
|
range_clause = _time_range_clause("occurred_at", since, until, params)
|
|
rows_result = await session.execute(
|
|
text(f"""
|
|
SELECT {_BOOT_COLUMNS}
|
|
FROM device_boot_events
|
|
WHERE is_fault
|
|
AND (crash IS NOT NULL OR pre_crash IS NOT NULL OR crash_task IS NOT NULL)
|
|
{range_clause}
|
|
ORDER BY occurred_at DESC
|
|
LIMIT :limit
|
|
"""),
|
|
params,
|
|
)
|
|
rows = rows_result.mappings().all()
|
|
|
|
return [_boot_row_to_dict(r) for r in rows]
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# 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
|