From 038ab471906963eacd6c4177dae3ac9f0108e4ea Mon Sep 17 00:00:00 2001 From: MarekWo Date: Tue, 14 Jul 2026 08:15:55 +0200 Subject: [PATCH] feat(observer): observer_brokers table and DB CRUD New observer_brokers table (name/host/port/credentials/TLS flags) plus Database CRUD mirroring the analyzers pattern. Groundwork for the MQTT packet-capture observer; no behavior change yet. Co-Authored-By: Claude Fable 5 --- app/database.py | 67 +++++++++++++++++++++++++++++++++++++++++++++++++ app/schema.sql | 16 ++++++++++++ 2 files changed, 83 insertions(+) diff --git a/app/database.py b/app/database.py index b2d4433..6a01acf 100644 --- a/app/database.py +++ b/app/database.py @@ -648,6 +648,73 @@ class Database: (analyzer_id,) ) + # ================================================================ + # Observer brokers (MQTT targets for packet capture publishing) + # ================================================================ + + def create_observer_broker(self, name: str, host: str, port: int = 1883, + username: str = '', password: str = '', + use_tls: bool = False, tls_verify: bool = True) -> int: + """Insert a new observer broker. Raises sqlite3.IntegrityError on duplicate name.""" + with self._connect() as conn: + cursor = conn.execute( + """INSERT INTO observer_brokers + (name, host, port, username, password, use_tls, tls_verify) + VALUES (?, ?, ?, ?, ?, ?, ?)""", + (name, host, port, username, password, + 1 if use_tls else 0, 1 if tls_verify else 0) + ) + return cursor.lastrowid + + def list_observer_brokers(self) -> List[Dict]: + with self._connect() as conn: + rows = conn.execute( + "SELECT * FROM observer_brokers ORDER BY name COLLATE NOCASE" + ).fetchall() + return [dict(r) for r in rows] + + def get_observer_broker(self, broker_id: int) -> Optional[Dict]: + with self._connect() as conn: + row = conn.execute( + "SELECT * FROM observer_brokers WHERE id = ?", (broker_id,) + ).fetchone() + return dict(row) if row else None + + def update_observer_broker(self, broker_id: int, name: Optional[str] = None, + host: Optional[str] = None, port: Optional[int] = None, + username: Optional[str] = None, password: Optional[str] = None, + use_tls: Optional[bool] = None, + tls_verify: Optional[bool] = None, + is_disabled: Optional[bool] = None) -> bool: + """Update fields on an observer broker. Pass None to leave a field unchanged.""" + sets = [] + params: List[Any] = [] + for column, value in (("name", name), ("host", host), ("port", port), + ("username", username), ("password", password)): + if value is not None: + sets.append(f"{column} = ?") + params.append(value) + for column, value in (("use_tls", use_tls), ("tls_verify", tls_verify), + ("is_disabled", is_disabled)): + if value is not None: + sets.append(f"{column} = ?") + params.append(1 if value else 0) + if not sets: + return False + sets.append("updated_at = datetime('now')") + params.append(broker_id) + with self._connect() as conn: + cursor = conn.execute( + f"UPDATE observer_brokers SET {', '.join(sets)} WHERE id = ?", + params + ) + return cursor.rowcount > 0 + + def delete_observer_broker(self, broker_id: int) -> bool: + with self._connect() as conn: + cursor = conn.execute("DELETE FROM observer_brokers WHERE id = ?", (broker_id,)) + return cursor.rowcount > 0 + def set_channel_scope(self, channel_idx: int, region_id: Optional[int]) -> None: """Set or clear the region mapping for a channel. diff --git a/app/schema.sql b/app/schema.sql index 8cec069..4c5d775 100644 --- a/app/schema.sql +++ b/app/schema.sql @@ -66,6 +66,22 @@ CREATE TABLE IF NOT EXISTS analyzers ( CREATE UNIQUE INDEX IF NOT EXISTS idx_analyzers_one_default ON analyzers(is_default) WHERE is_default = 1; +-- Observer: MQTT brokers that receive captured mesh packets +-- (meshcore-packet-capture compatible publishing) +CREATE TABLE IF NOT EXISTS observer_brokers ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + name TEXT NOT NULL UNIQUE, + host TEXT NOT NULL, + port INTEGER NOT NULL DEFAULT 1883, + username TEXT NOT NULL DEFAULT '', + password TEXT NOT NULL DEFAULT '', -- plaintext (single-user LAN app) + use_tls INTEGER NOT NULL DEFAULT 0, + tls_verify INTEGER NOT NULL DEFAULT 1, -- 0 = accept self-signed certs + is_disabled INTEGER NOT NULL DEFAULT 0, + created_at TEXT NOT NULL DEFAULT (datetime('now')), + updated_at TEXT NOT NULL DEFAULT (datetime('now')) +); + -- Per-channel region mapping (absent row = no override; firmware default applies) CREATE TABLE IF NOT EXISTS channel_scopes ( channel_idx INTEGER PRIMARY KEY,