From 27fa2381ea3a5c015529f726e681c3509934ef1d Mon Sep 17 00:00:00 2001 From: Rigear <278971+Rigear@users.noreply.github.com> Date: Wed, 15 Apr 2026 21:20:11 -0700 Subject: [PATCH] feat: * Added retain status message bool * Added back old templates * Added migration path from old mqtt and letsmesh configs to new mqtt_broker config --- repeater/config.py | 3 +- repeater/data_acquisition/mqtt_handler.py | 281 ++++++++++++++---- .../data_acquisition/storage_collector.py | 14 +- repeater/data_acquisition/storage_utils.py | 2 +- repeater/engine.py | 6 +- repeater/web/api_endpoints.py | 186 ++++++------ 6 files changed, 337 insertions(+), 155 deletions(-) diff --git a/repeater/config.py b/repeater/config.py index be8e095..d9bcc93 100644 --- a/repeater/config.py +++ b/repeater/config.py @@ -30,7 +30,8 @@ def get_node_info(config: Dict[str, Any]) -> Dict[str, Any]: radio_bw_khz = radio_bw / 1_000 radio_config_str = f"{radio_freq_mhz},{radio_bw_khz},{radio_sf},{radio_cr}" - mqtt_config = config.get("mqtt", {}) + # Handle getting the config from mqtt brokers, falling back to letsmesh if it doesn't exist + mqtt_config = config.get("mqtt_brokers", config.get("letsmesh", {})) return { "node_name": node_name, diff --git a/repeater/data_acquisition/mqtt_handler.py b/repeater/data_acquisition/mqtt_handler.py index 7b8c125..3cf89f6 100644 --- a/repeater/data_acquisition/mqtt_handler.py +++ b/repeater/data_acquisition/mqtt_handler.py @@ -17,7 +17,7 @@ except Exception: from datetime import timezone UTC = timezone.utc -from repeater import __version__ +from repeater import __version__, config # Try to import paho-mqtt error code mappings try: @@ -36,6 +36,23 @@ logger = logging.getLogger("MQTTHandler") def b64url(x: bytes) -> str: return base64.urlsafe_b64encode(x).rstrip(b"=").decode() +LETSMESH_BROKERS = [ + { + "name": "Europe (LetsMesh v1)", + "host": "mqtt-eu-v1.letsmesh.net", + "port": 443, + "audience": "mqtt-eu-v1.letsmesh.net", + "use_jwt_auth": True, + }, + { + "name": "US West (LetsMesh v1)", + "host": "mqtt-us-v1.letsmesh.net", + "port": 443, + "audience": "mqtt-us-v1.letsmesh.net", + "use_jwt_auth": True, + }, +] + # ==================================================================== # Single Broker Connection Manager @@ -53,10 +70,10 @@ class _BrokerConnection: public_key: str, iata_code: str, jwt_expiry_minutes: int, - use_tls: bool, email: str, owner: str, broker_index: int, + node_name: str, on_connect_callback: Optional[Callable] = None, on_disconnect_callback: Optional[Callable] = None ): @@ -65,27 +82,46 @@ class _BrokerConnection: self.public_key = public_key.upper() self.iata_code = iata_code self.jwt_expiry_minutes = jwt_expiry_minutes - self.use_tls = use_tls self.email = email self.owner = owner + self.node_name = node_name self.broker_index = broker_index self._on_connect_callback = on_connect_callback self._on_disconnect_callback = on_disconnect_callback self._connect_time = None - self._tls_verified = False self._running = False self._reconnect_attempts = 0 self._reconnect_timer = None self._max_reconnect_delay = 300 # 5 minutes max self._jwt_refresh_timer = None self.transport = broker.get('transport', 'websockets') - client_id = f"meshcore_{self.public_key}_{broker['host']}" - self.client = mqtt.Client(client_id=client_id, transport=self.transport) - self.client.on_connect = self._on_connect - self.client.on_disconnect = self._on_disconnect + self.use_jwt_auth = broker.get('use_jwt_auth', False) self.username = broker.get('username', None) self.password = broker.get('password', None) + + self.format=broker.get("format", "letsmesh") + self.tls=broker.get("tls", None) + + client_id = f"meshcore_{self.public_key}_{broker['host']}_{self.format}" + self.client = mqtt.Client(client_id=client_id, transport=self.transport) + self.client.on_connect = self._on_connect + self.client.on_disconnect = self._on_disconnect + + # If None, will be use defaults depending on the format value + self.base_topic=broker.get("base_topic", None) + + self.enabled = broker.get("enabled", False) + self.retain_status = broker.get("retain_status", False) + + if self.base_topic is None: + if self.format == "mqtt": + self.base_topic = f"meshcore/repeater/{self.node_name}" + elif self.format == "letsmesh": + self.base_topic = f"meshcore/{self.iata_code}/{self.public_key}" + else: + logger.warning(f"Unknown broker format '{self.format}' for {self.broker['name']}, using default base topic") + self.base_topic = f"meshcore/{self.iata_code}/{self.public_key}" from pymc_core.protocol.utils import PAYLOAD_TYPES @@ -280,16 +316,29 @@ class _BrokerConnection: self.client.disconnect() logger.info(f"Disconnected from {self.broker['name']}") - def publish(self, topic: str, payload: str, retain: bool = False, qos: int = 0): + def publish(self, subtopic: str, payload: str, retain: bool = False, qos: int = 0): """Publish message to broker""" - logger.debug(f"Publishing to topic '{topic}' with payload: {payload}: self._running={self._running}") + + # Legacy MQTT config uses singular "packet" topic, while LetsMesh uses "packets". Handle this for compatibility. + if self.format == "mqtt" and subtopic == "packets": + subtopic = "packet" + + if(subtopic == "status"): # Override the status topic retain and qos settings based on broker configuration + retain = self.retain_status + qos = 1 if self.retain_status else 0 + + logger.debug(f"Publishing to topic '{self.base_topic}/{subtopic}' with payload: [{payload}]. Running={self._running}. Retain={retain}, QoS={qos}") if self._running: - result = self.client.publish(topic, payload, retain=retain, qos=qos) + result = self.client.publish(f"{self.base_topic}/{subtopic}", payload, retain=retain, qos=qos) return result else: logger.warning(f"Cannot publish to {self.broker['name']} - not connected") return None + def is_enabled(self) -> bool: + """Check if connection is enabled""" + return self.enabled + def is_connected(self) -> bool: """Check if connection is active""" return self._running @@ -344,7 +393,7 @@ class _BrokerConnection: # ==================================================================== -# MeshCore → MQTT Publisher with Ed25519 auth token +# MeshCore → MQTT Publisher # ==================================================================== class MeshCoreToMqttPusher: @@ -365,51 +414,54 @@ class MeshCoreToMqttPusher: node_info = get_node_info(config) - iata_code = node_info["iata_code"] + self.iata_code = node_info["iata_code"] self.email = node_info.get("email", "") self.owner = node_info.get("owner", "") - status_interval = node_info["status_interval"] - node_name = node_info["node_name"] - radio_config = node_info["radio_config"] - - # Get additional brokers from config (optional) - mqtt_config = config.get("mqtt", {}) - brokers = mqtt_config.get("brokers", []) - - # Add additional brokers from config - self.brokers = [] - if brokers: - for broker_config in brokers: - if all(k in broker_config for k in ["name", "host", "port", "enabled"]): - if broker_config["enabled"]: - self.brokers.append(broker_config) - logger.info(f"Added broker: {broker_config['name']}") - else: - logger.info(f"Broker disabled in config, skipping: {broker_config['name']}") - else: - logger.warning(f"Skipping invalid broker config: {broker_config}") - - # Validate that we have at least one broker - # if not self.brokers: - # raise ValueError( - # "No brokers configured. Either set broker_index to a valid value " - # "or provide additional_brokers in config." - # ) - + self.status_interval = node_info["status_interval"] + self.node_name = node_info["node_name"] self.local_identity = local_identity self.public_key = public_key - self.iata_code = iata_code self.jwt_expiry_minutes = jwt_expiry_minutes - self.use_tls = use_tls - self.status_interval = status_interval self.app_version = __version__ - self.node_name = node_name - self.radio_config = radio_config + self.radio_config = node_info["radio_config"] self.stats_provider = stats_provider self._status_task = None self._running = False self._lock = threading.Lock() + # Initialize brokers list + mqtt_brokers_config = config.get("mqtt_brokers", {}) + letsmesh_config = config.get("letsmesh", {}) + mqtt_config = config.get("mqtt", {}) + + brokers = [] + if mqtt_brokers_config: + # Pull in brokers from mqtt_brokers config + brokers.extend(mqtt_brokers_config.get("brokers", [])) + + if letsmesh_config or mqtt_config: + logger.warning("Multiple MQTT broker configurations found (mqtt_brokers, letsmesh, mqtt). Only mqtt_brokers will be used") + + else: + if mqtt_config: + imported_mqtt_config = self.convert_mqtt_to_broker_config(mqtt_config) + brokers.append(imported_mqtt_config) + + if letsmesh_config: + imported_letsmesh_configs = self.convert_letsmesh_to_broker_config(letsmesh_config) + brokers.extend(imported_letsmesh_configs) + + self.brokers = [] + if brokers: + for broker_config in brokers: + if all(k in broker_config for k in ["name", "host", "port", "enabled"]): + self.brokers.append(broker_config) + logger.info(f"Added broker: {broker_config['name']}") + else: + logger.warning(f"Skipping invalid broker config: {broker_config}") + + + # Create broker connections self.connections: List[_BrokerConnection] = [] for idx, broker in enumerate(self.brokers): @@ -419,10 +471,10 @@ class MeshCoreToMqttPusher: public_key=self.public_key, iata_code=self.iata_code, jwt_expiry_minutes=self.jwt_expiry_minutes, - use_tls=self.use_tls, email=self.email, owner=self.owner, broker_index=idx, + node_name=self.node_name, on_connect_callback=self._on_broker_connected, on_disconnect_callback=self._on_broker_disconnected, ) @@ -430,6 +482,100 @@ class MeshCoreToMqttPusher: logger.info(f"Initialized with {len(self.connections)} broker connection(s)") + # Convert legacy configration to new one + if not mqtt_brokers_config: + logger.info("Storing mqtt_brokers config from legacy mqtt/letsmesh configuration") + mqtt_brokers_config = { + "iata_code": self.iata_code, + "status_interval": self.status_interval, + "owner": self.owner, + "email": self.email, + "brokers": brokers + } + + # Update the configuration with the new configuration + config["mqtt_brokers"] = mqtt_brokers_config + + def convert_mqtt_to_broker_config(self, mqtt_cfg: dict) -> dict: + """Convert legacy MQTT config format to internal broker config format""" + logger.info(f"Imported MQTT broker from 'mqtt' config: {mqtt_cfg['broker']}") + transport = "websockets" if mqtt_cfg.get("use_websockets", False) else "tcp" + return { + "enabled": mqtt_cfg.get("enabled", False), + "name": mqtt_cfg["broker"], + "host": mqtt_cfg["broker"], + "port": mqtt_cfg["port"], + "use_jwt_auth": False, # The legacy MQTT config does not support JWT auth, so we set this to False + "username": mqtt_cfg.get("username", None), + "password": mqtt_cfg.get("password", None), + "transport": transport, + "tls": mqtt_cfg.get("tls", None), + "format": "mqtt", + "base_topic": mqtt_cfg.get("base_topic", None), + } + + def convert_letsmesh_to_broker_config(self, letsmesh_cfg: dict) -> List[dict]: + """Convert LetsMesh config format to internal broker config format""" + + brokers = [] + + enabled = letsmesh_cfg.get("enabled", False) + + idx = letsmesh_cfg.get("broker_index", None) + if idx == 0 or idx == 1: + broker_info = LETSMESH_BROKERS[idx] + logger.info(f"Imported LetsMesh broker from 'letsmesh' config: {broker_info['name']}") + brokers.append({ + "enabled": enabled, + "name": broker_info["name"], + "host": broker_info["host"], + "port": broker_info["port"], + "audience": broker_info["audience"], + "use_jwt_auth": True, + "transport": "websockets", + "tls": None, + "format": "letsmesh", + "base_topic": None, + "retain_status": False + }) + elif idx < 0: + if idx == -1: + brokers.extend({ + "enabled": enabled, + "name": broker_info["name"], + "host": broker_info["host"], + "port": broker_info["port"], + "audience": broker_info["audience"], + "use_jwt_auth": True, + "transport": "websockets", + "tls": None, + "format": "letsmesh", + "base_topic": None, + "retain_status": False + } for broker_info in LETSMESH_BROKERS) + + additional = letsmesh_cfg.get("additional_brokers", []) + for add_broker in additional: + logger.info(f"Imported additional LetsMesh broker from 'letsmesh' config: {add_broker['name']}") + brokers.append({ + "enabled": enabled, + "name": add_broker["name"], + "host": add_broker["host"], + "port": add_broker["port"], + "audience": add_broker["audience"], + "use_jwt_auth": True, + "transport": "websockets", + "use_jwt_auth": add_broker.get("use_jwt_auth", True), + "transport": add_broker.get("transport", "websockets"), + "tls": None, + "format": "letsmesh", + "base_topic": None, + "retain_status": False + }) + + + return brokers # Placeholder for now - we will implement this if we need to support the old letsmesh config format + def _on_broker_connected(self, broker_name: str): """Callback when a broker connects""" # Publish initial status on first connection @@ -523,9 +669,6 @@ class MeshCoreToMqttPusher: def _process_packet(self, pkt: dict) -> dict: return {"timestamp": datetime.now(UTC).isoformat(), "origin_id": self.public_key, **pkt} - def _topic(self, subtopic: str) -> str: - return f"meshcore/{self.iata_code}/{self.public_key}/{subtopic}" - def publish_packet(self, pkt: dict, subtopic="packets", retain=False): return self.publish(subtopic, self._process_packet(pkt), retain) @@ -576,26 +719,50 @@ class MeshCoreToMqttPusher: def publish(self, subtopic: str, payload: dict, retain: bool = False, qos: int = 0): """Publish message to all connected brokers""" - topic = self._topic(subtopic) message = json.dumps(payload) - logger.debug(f"Publishing to topic '{topic}' with payload: {message}") + # _BrokerConnection now handles topic prefixing, so we only log the subtopic here + logger.debug(f"Publishing to topic '{subtopic}' with payload: {message}") packet_type = payload.get("type") results = [] with self._lock: for conn in self.connections: - if conn.is_connected(): + if conn.enabled and conn.is_connected(): if packet_type in conn.disallowed_types: logger.debug(f"Skipped publishing packet type 0x{packet_type:02X} (disallowed)") - return - result = conn.publish(topic, message, retain=retain, qos=qos) + continue + result = conn.publish(subtopic, message, retain=retain, qos=qos) results.append((conn.broker["name"], result)) - logger.debug(f"Published to {conn.broker['name']}/{topic}") + logger.debug(f"Published to {conn.broker['name']} -- {subtopic}") if not results: - logger.warning(f"No active broker connections for publishing to {topic}") + logger.warning(f"No active broker connections for publishing to {subtopic}") + + return results + + + def publish_mqtt(self, subtopic: str, payload: dict, retain: bool = False, qos: int = 0): + """Publish message to all connected brokers""" + message = json.dumps(payload) + + # _BrokerConnection now handles topic prefixing, so we only log the subtopic here + logger.debug(f"Publishing to topic '{subtopic}' with payload: {message}") + + results = [] + with self._lock: + for conn in self.connections: + if conn.enabled and conn.is_connected(): + if conn.format != "mqtt": + logger.debug(f"Skipped publishing to {conn.broker['name']} (wrong format)") + continue + result = conn.publish(subtopic, message, retain=retain, qos=qos) + results.append((conn.broker["name"], result)) + logger.debug(f"Published to {conn.broker['name']} -- {subtopic}") + + if not results: + logger.warning(f"No active broker connections for publishing to {subtopic}") return results diff --git a/repeater/data_acquisition/storage_collector.py b/repeater/data_acquisition/storage_collector.py index 517e3b4..8df6ec4 100644 --- a/repeater/data_acquisition/storage_collector.py +++ b/repeater/data_acquisition/storage_collector.py @@ -32,7 +32,7 @@ class StorageCollector: # Initialize MQTT handler if configured self.mqtt_handler = None - if config.get("mqtt", {}) and local_identity: + if (config.get("mqtt_brokers", {}) or config.get("letsmesh", {}) or config.get("mqtt", {})) and local_identity: try: # Pass local_identity directly (supports both standard and firmware keys) self.mqtt_handler = MeshCoreToMqttPusher( @@ -159,9 +159,9 @@ class StorageCollector: # f"Skipping mqtt publish for packet with drop_reason: {packet_record.get('drop_reason')}" # ) # else: - self._publish_to_mqtt(packet_record) + self._publish_packet_to_mqtt(packet_record) - def _publish_to_mqtt(self, packet_record: dict): + def _publish_packet_to_mqtt(self, packet_record: dict): """Publish packet to mqtt broker if enabled and allowed""" if not self.mqtt_handler: return @@ -172,10 +172,6 @@ class StorageCollector: logger.error("Cannot publish to mqtt: packet_record missing 'type' field") return - # if packet_type in self.disallowed_packet_types: - # logger.debug(f"Skipped publishing packet type 0x{packet_type:02X} (disallowed)") - # return - node_name = self.config.get("repeater", {}).get("node_name", "Unknown") packet = PacketRecord.from_packet_record( packet_record, origin=node_name, origin_id=self.mqtt_handler.public_key @@ -192,12 +188,12 @@ class StorageCollector: def record_advert(self, advert_record: dict): self.sqlite_handler.store_advert(advert_record) - #self.old_mqtt_handler.publish(advert_record, "advert") + self.mqtt_handler.publish_mqtt("advert", advert_record) def record_noise_floor(self, noise_floor_dbm: float): noise_record = {"timestamp": time.time(), "noise_floor_dbm": noise_floor_dbm} self.sqlite_handler.store_noise_floor(noise_record) - #self.old_mqtt_handler.publish(noise_record, "noise_floor") + self.mqtt_handler.publish_mqtt("noise_floor", noise_record) def record_crc_errors(self, count: int): """Record a batch of CRC errors detected since last poll.""" diff --git a/repeater/data_acquisition/storage_utils.py b/repeater/data_acquisition/storage_utils.py index bde938e..f0d0751 100644 --- a/repeater/data_acquisition/storage_utils.py +++ b/repeater/data_acquisition/storage_utils.py @@ -10,7 +10,7 @@ class PacketRecord: """ Data class for packet record format. Converts internal packet_record format to standardized publish format. - Reusable across MQTT, LetsMesh, and other handlers. + Reusable across MQTT and other handlers. """ origin: str diff --git a/repeater/engine.py b/repeater/engine.py index 38f47df..8ca0a18 100644 --- a/repeater/engine.py +++ b/repeater/engine.py @@ -359,10 +359,10 @@ class RepeaterHandler(BaseHandler): ) # Store packet record to persistent storage - # Skip LetsMesh only for invalid packets (not duplicates or operational drops) + # Skip mqtt only for invalid packets (not duplicates or operational drops) if self.storage: try: - # Only skip LetsMesh for actual invalid/bad packets + # Only skip mqtt for actual invalid/bad packets invalid_reasons = ["Invalid advert packet", "Empty payload", "Path too long"] skip_mqtt = drop_reason in invalid_reasons if drop_reason else False self.storage.record_packet(packet_record, skip_mqtt_if_invalid=skip_mqtt) @@ -1138,7 +1138,7 @@ class RepeaterHandler(BaseHandler): "unscoped_flood_allow": self.config.get("mesh", {}).get("unscoped_flood_allow", self.config.get("mesh", {}).get("global_flood_allow", True)), "path_hash_mode": self.config.get("mesh", {}).get("path_hash_mode", 0), }, - #"mqtt": self.config.get("mqtt", {}), + "mqtt_brokers": self.config.get("mqtt_brokers", {}), }, "public_key": None, } diff --git a/repeater/web/api_endpoints.py b/repeater/web/api_endpoints.py index 03a0ba1..7799f8b 100644 --- a/repeater/web/api_endpoints.py +++ b/repeater/web/api_endpoints.py @@ -1003,8 +1003,7 @@ class APIEndpoints: """Get MQTT connection status and configuration.""" self._set_cors_headers() try: - mqtt_cfg = self.config.get("mqtt", {}) - enabled = mqtt_cfg.get("enabled", False) + # mqtt_cfg = self.config.get("mqtt_brokers", {}) # Walk the chain to the mqtt_handler handler = None @@ -1018,14 +1017,17 @@ class APIEndpoints: if handler: for conn in getattr(handler, "connections", []): connected_brokers.append({ + "enabled": conn.enabled, "name": conn.broker.get("name", ""), "host": conn.broker.get("host", ""), - "connected": conn.is_connected(), - "reconnecting": conn.has_pending_reconnect(), + "status": { + "connected": conn.is_connected(), + "reconnecting": conn.has_pending_reconnect(), + }, + "format": conn.format }) return self._success({ - "enabled": enabled, "handler_active": handler is not None, "brokers": connected_brokers, }) @@ -1033,95 +1035,111 @@ class APIEndpoints: logger.error(f"Error getting MQTT status: {e}") return self._error(str(e)) - # @cherrypy.expose - # @cherrypy.tools.json_out() - # @cherrypy.tools.json_in() - # def update_mqtt_config(self): - # """Update MQTT Observer configuration. + @cherrypy.expose + @cherrypy.tools.json_out() + @cherrypy.tools.json_in() + def update_mqtt_config(self): + """Update MQTT Observer configuration. - # POST /api/update_mqtt_config - # Body: { - # "iata_code": "SFO", - # "status_interval": 300, - # "owner": "Callsign", - # "email": "user@example.com", - # "disallowed_packet_types": ["ACK"] - # } - # """ - # self._set_cors_headers() + POST /api/update_mqtt_config + Body: { + "iata_code": "SFO", + "status_interval": 300, + "owner": "Callsign", + "email": "user@example.com", + "brokers": [ + { + + }] + } + """ + self._set_cors_headers() - # if cherrypy.request.method == "OPTIONS": - # return "" + if cherrypy.request.method == "OPTIONS": + return "" - # try: - # self._require_post() - # data = cherrypy.request.json or {} + try: + self._require_post() + data = cherrypy.request.json or {} - # if not data: - # return self._error("No configuration updates provided") + if not data: + return self._error("No configuration updates provided") - # letsmesh_updates = {} + mqtt_updates = {} - # if "enabled" in data: - # letsmesh_updates["enabled"] = bool(data["enabled"]) - # if "iata_code" in data: - # letsmesh_updates["iata_code"] = str(data["iata_code"]).strip() - # if "broker_index" in data: - # letsmesh_updates["broker_index"] = int(data["broker_index"]) - # if "status_interval" in data: - # letsmesh_updates["status_interval"] = max(60, int(data["status_interval"])) - # if "owner" in data: - # letsmesh_updates["owner"] = str(data["owner"]).strip() - # if "email" in data: - # letsmesh_updates["email"] = str(data["email"]).strip() - # if "disallowed_packet_types" in data: - # letsmesh_updates["disallowed_packet_types"] = list(data["disallowed_packet_types"]) - # if "additional_brokers" in data: - # brokers = data["additional_brokers"] - # if not isinstance(brokers, list): - # return self._error("additional_brokers must be a list") - # validated = [] - # for i, b in enumerate(brokers): - # if not isinstance(b, dict): - # return self._error(f"Broker at index {i} must be an object") - # for field in ("name", "host", "audience"): - # if not b.get(field, "").strip(): - # return self._error(f"Broker at index {i} missing required field: {field}") - # try: - # port = int(b.get("port", 443)) - # except (ValueError, TypeError): - # return self._error(f"Broker at index {i} has invalid port") - # validated.append({ - # "name": str(b["name"]).strip(), - # "host": str(b["host"]).strip(), - # "port": port, - # "audience": str(b["audience"]).strip(), - # }) - # letsmesh_updates["additional_brokers"] = validated + if "iata_code" in data: + mqtt_updates["iata_code"] = str(data["iata_code"]).strip() + if "status_interval" in data: + mqtt_updates["status_interval"] = max(60, int(data["status_interval"])) + if "owner" in data: + mqtt_updates["owner"] = str(data["owner"]).strip() + if "email" in data: + mqtt_updates["email"] = str(data["email"]).strip() + # if "disallowed_packet_types" in data: + # mqtt_updates["disallowed_packet_types"] = list(data["disallowed_packet_types"]) + if "brokers" in data: + brokers = data["brokers"] + if not isinstance(brokers, list): + return self._error("brokers must be a list") + validated = [] + for i, b in enumerate(brokers): + if not isinstance(b, dict): + return self._error(f"Broker at index {i} must be an object") + for field in ("name", "host", "port", "format"): + if not b.get(field, ""): + return self._error(f"Broker at index {i} missing required field: {field}") + + try: + port = int(b.get("port", 443)) + except (ValueError, TypeError): + return self._error(f"Broker at index {i} has invalid port") + + new_broker = { + "enabled": b.get("enabled", False), + "name": str(b["name"]).strip(), + "transport": str(b.get("transport", "websockets")).strip(), + "host": str(b["host"]).strip(), + "port": port, + "format": str(b["format"]).strip(), + "disallowed_packet_types": list(b.get("disallowed_packet_types", [])), + "retain_status": bool(b.get("retain_status", False)), + } + + if b.get("use_jwt_auth", False): + new_broker["use_jwt_auth"] = True + new_broker["audience"] = str(b["audience"]).strip() + else: + new_broker["use_jwt_auth"] = False + new_broker["username"] = b.get("username", None) + new_broker["password"] = b.get("password", None) - # if not letsmesh_updates: - # return self._error("No valid settings provided") + validated.append(new_broker) - # result = self.config_manager.update_and_save( - # updates={"letsmesh": letsmesh_updates}, - # live_update=False, # Restart required for LetsMesh handler changes - # ) + mqtt_updates["brokers"] = validated - # if result.get("success"): - # logger.info(f"LetsMesh config updated: {list(letsmesh_updates.keys())}") - # return self._success({ - # "persisted": result.get("saved", False), - # "restart_required": True, - # "message": "Observer settings saved. Restart the service for changes to take effect.", - # }) - # else: - # return self._error(result.get("error", "Failed to update LetsMesh configuration")) + if not mqtt_updates: + return self._error("No valid settings provided") - # except cherrypy.HTTPError: - # raise - # except Exception as e: - # logger.error(f"Error updating LetsMesh config: {e}") - # return self._error(str(e)) + result = self.config_manager.update_and_save( + updates={"mqtt_brokers": mqtt_updates, "mqtt": None, "letsmesh": None}, + live_update=False, # Restart required for MQTT handler changes + ) + + if result.get("success"): + logger.info(f"MQTT config updated: {list(mqtt_updates.keys())}") + return self._success({ + "persisted": result.get("saved", False), + "restart_required": True, + "message": "Observer settings saved. Restart the service for changes to take effect.", + }) + else: + return self._error(result.get("error", "Failed to update LetsMesh configuration")) + + except cherrypy.HTTPError: + raise + except Exception as e: + logger.error(f"Error updating LetsMesh config: {e}") + return self._error(str(e)) @cherrypy.expose @cherrypy.tools.json_out()