Исходный код soniks_client.mqtt

"""Живой канал с порталом: MQTT через WebSocket на ``sonik.space``.

Портал публикует документ состояния станции в момент сохранения формы, станция
применяет его сразу, а не на следующей сверке расписания, и тут же публикует
статус. Last Will даёт порталу мгновенный «offline». REST остаётся источником
истины и фолбэком: без брокера (закрытая сеть, старая база без ``paho``) всё то
же самое происходит раз в минуту через ``remote_config.sync()``.

Топики: ``stations/<id>/state`` (портал → станция, retained),
``stations/<id>/status`` (станция → портал, тело — как у ``POST status/``),
``stations/<id>/online`` (retained ``1``, Last Will ``0``).
"""

import json
import threading
from urllib.parse import urlsplit

from core.configs import logger, settings
from soniks_client import api, remote_config, sdr_survey

try:
    import paho.mqtt.client as paho
except ImportError:  # база образа собрана без paho — работаем на REST
    paho = None

_client = None
_publish_lock = threading.Lock()


def _topic(kind: str) -> str:
    return f"stations/{settings.station.ID}/{kind}"


[документация] def start_mqtt() -> bool: """Подключиться к брокеру в фоновом потоке. Returns: ``True``, если поток запущен. Ошибки гасятся — станция без MQTT живёт на REST. """ global _client if not settings.mqtt.ENABLED: return False if paho is None: logger.warning("paho-mqtt не установлен, живой канал с порталом выключен") return False url = urlsplit(settings.mqtt.URL) secure = url.scheme == "wss" username = f"station-{settings.station.ID}" try: client = paho.Client( paho.CallbackAPIVersion.VERSION2, client_id=username, transport="websockets", ) client.ws_set_options(path=url.path or "/mqtt") if secure: client.tls_set() client.username_pw_set(username, settings.station.TOKEN) client.will_set(_topic("online"), "0", qos=1, retain=True) client.reconnect_delay_set(min_delay=1, max_delay=120) client.on_connect = _on_connect client.on_message = _on_message client.on_disconnect = _on_disconnect # До запуска цикла: on_connect публикует статус через модульную ссылку. _client = client client.connect_async( url.hostname, url.port or (443 if secure else 80), keepalive=60 ) client.loop_start() except Exception as e: _client = None logger.error("MQTT не запущен (%s): %s", settings.mqtt.URL, e) return False logger.info("MQTT: подключение к %s", settings.mqtt.URL) return True
[документация] def publish_status() -> None: """Опубликовать статус станции — тот же документ, что уходит REST'ом.""" if _client is None or not _client.is_connected(): return with _publish_lock: _client.publish(_topic("status"), json.dumps(api.status_body()), qos=1)
def _on_connect(client, userdata, flags, reason_code, properties) -> None: if reason_code.is_failure: # Неверный токен или чужая станция — переподключение не поможет, # но и ломать ничего не надо: REST работает. logger.error("MQTT: брокер отказал в подключении: %s", reason_code) return logger.info("MQTT: подключено") client.publish(_topic("online"), "1", qos=1, retain=True) client.subscribe(_topic("state"), qos=1) publish_status() def _on_disconnect(client, userdata, flags, reason_code, properties) -> None: logger.warning("MQTT: соединение потеряно (%s), переподключение", reason_code) def _on_message(client, userdata, message) -> None: try: doc = json.loads(message.payload) except ValueError as e: logger.warning("MQTT: документ состояния не разобран: %s", e) return if not isinstance(doc, dict): return try: remote_config.apply_state(doc) except Exception as e: logger.exception("MQTT: ошибка применения состояния: %s", e) publish_status() # Команда с портала («пересканировать», «откалибровать») — в фон, не в # потоке брокера: обследование длится секунды, keepalive ждать не будет. sdr_survey.run_due()