Files
bellsystems-cp/backend/alembic/versions/b8c9d0e1f2a3_mqtt_v2_topic_migration.py
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

79 lines
3.7 KiB
Python

"""mqtt v2 topic migration
Firmware moved to a new MQTT topic spec (project-vesper's
docs/reference/vesper_mqtt_topic_spec_v2.md, feature-catalog F-062). Three
schema changes needed to keep ingestion correct against the new payload shapes:
1. heartbeats.state / heartbeats.ok — the v2 heartbeat payload is flat (no more
{"status":"INFO","type":"heartbeat","payload":{...}} wrapper) and adds two
new fields: state ("idle"/"playing"/"paused"/"error"/"booting") and ok
(overall health), both independently derived on the firmware side — not
mirrored from status/playback or system/alerts. Nullable: older rows
(pre-migration) and any device still reporting have no value for these.
2. device_diagnostics_reports gains bell_strikes / bell_loads / cooling_active.
The 5-min diagnostics report moved from status/info to system/metrics and
picked up per-bell strike/heat data along the way (the topic spec calls out
"bell heat ratings" explicitly for this topic). Stored as JSON-encoded TEXT,
matching this table's existing stack_high_water column — same reasoning:
16 fixed bell slots, but keeping the encoding consistent with the sibling
column is simpler than mixing raw JSONB and TEXT in one row shape.
3. device_reports — new table for the new control/reports topic. Board-
initiated, unsolicited, critical events (currently only bell_overload).
Distinct from device_alert_events (subsystem health transitions) and
device_logs (routine log lines) — this is a narrow, insert-only table for
a different kind of signal: urgent, bell-mechanism-specific events the
physical tablets act on in real time. The console side is deliberately
light for now — store + list for historical/audit purposes; the tablets
are the real-time consumer of control/reports, not this console.
Revision ID: b8c9d0e1f2a3
Revises: a7b8c9d0e1f2
Create Date: 2026-09-21 00:00:00.000000
"""
from typing import Sequence, Union
import sqlalchemy as sa
from alembic import op
revision: str = "b8c9d0e1f2a3"
down_revision: Union[str, None] = "a7b8c9d0e1f2"
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None
def upgrade() -> None:
op.add_column("heartbeats", sa.Column("state", sa.String(16), nullable=True))
op.add_column("heartbeats", sa.Column("ok", sa.Boolean(), nullable=True))
op.add_column("device_diagnostics_reports", sa.Column("bell_strikes", sa.Text(), nullable=True))
op.add_column("device_diagnostics_reports", sa.Column("bell_loads", sa.Text(), nullable=True))
op.add_column("device_diagnostics_reports", sa.Column("cooling_active", sa.Boolean(), nullable=True))
op.create_table(
"device_reports",
sa.Column("id", sa.BigInteger(), primary_key=True, autoincrement=True),
sa.Column("device_serial", sa.String(128), nullable=False),
sa.Column("report_type", sa.String(64), nullable=False),
sa.Column("payload", sa.Text(), nullable=True),
sa.Column("occurred_at", sa.DateTime(timezone=True), nullable=False,
server_default=sa.func.now()),
)
op.create_index(
"idx_device_reports_serial_occurred",
"device_reports",
["device_serial", sa.text("occurred_at DESC")],
)
def downgrade() -> None:
op.drop_index("idx_device_reports_serial_occurred", table_name="device_reports")
op.drop_table("device_reports")
op.drop_column("device_diagnostics_reports", "cooling_active")
op.drop_column("device_diagnostics_reports", "bell_loads")
op.drop_column("device_diagnostics_reports", "bell_strikes")
op.drop_column("heartbeats", "ok")
op.drop_column("heartbeats", "state")