"""Живой канал с порталом: 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()