OPi: авто-реконнект MQTT + проверка rc у publish (heartbeat и вес молча терялись с 10.07)

This commit is contained in:
claude
2026-07-15 07:38:35 +03:00
parent 451a30f74d
commit 94044fde75
2 changed files with 103 additions and 27 deletions
+96 -27
View File
@@ -176,34 +176,98 @@ def on_message(client, userdata, msg):
mqtt1 = mqtt.Client(client_id="scales_opi_1")
mqtt2 = mqtt.Client(client_id="scales_opi_2")
# (client, имя, host, port, топик подписки)
_MQTT_DEFS = (
(mqtt1, 'MQTT1', MQTT1_HOST, MQTT1_PORT, MQTT_GO_TOPIC),
(mqtt2, 'MQTT2', MQTT2_HOST, MQTT2_PORT, MQTT_CMD_TOPIC),
)
_loop_started = {'MQTT1': False, 'MQTT2': False}
def _setup_client(client, name, sub_topic):
"""Колбэки + backoff. Подписку восстанавливаем в on_connect (после КАЖДОГО реконнекта)."""
client.on_message = on_message
try:
client.reconnect_delay_set(min_delay=1, max_delay=60)
except Exception:
pass
def _on_connect(c, userdata, flags, rc, *a):
if rc == 0:
print(f"[{name}] Подключён к брокеру, подписка на {sub_topic}")
try:
c.subscribe(sub_topic)
except Exception as e:
print(f"[{name}] Ошибка подписки: {e}")
else:
print(f"[{name}] Подключение отклонено, rc={rc}")
def _on_disconnect(c, userdata, *a):
rc = a[-1] if a else '?'
print(f"[{name}] ОТКЛЮЧЁН от брокера (rc={rc}) — будет переподключение")
client.on_connect = _on_connect
client.on_disconnect = _on_disconnect
def mqtt_connect():
if MQTT2_USER:
try:
mqtt2.username_pw_set(MQTT2_USER, MQTT2_PASS)
except Exception as e:
print(f"[MQTT2] Ошибка авторизации: {e}")
for client, name, host, port, sub in _MQTT_DEFS:
_setup_client(client, name, sub)
try:
client.connect(host, port, keepalive=60)
client.loop_start()
_loop_started[name] = True
print(f"[{name}] Подключение к {host}:{port} инициировано")
except Exception as e:
# сеть может быть не готова при старте — поднимем в ensure_mqtt_connected()
print(f"[{name}] Ошибка подключения к {host}:{port}: {e}")
if not _loop_started[name]:
try:
client.loop_start()
_loop_started[name] = True
except Exception:
pass
def ensure_mqtt_connected():
"""Самолечение: если клиент отвалился — поднимаем связь. Вызывается из heartbeat_loop."""
for client, name, host, port, sub in _MQTT_DEFS:
try:
if client.is_connected():
continue
print(f"[{name}] Нет соединения — переподключаюсь к {host}:{port}")
try:
client.reconnect()
except Exception:
client.connect(host, port, keepalive=60)
if not _loop_started[name]:
client.loop_start()
_loop_started[name] = True
except Exception as e:
print(f"[{name}] Переподключение не удалось: {e}")
def _pub(client, name, topic, payload, **kw):
"""publish с проверкой результата: paho НЕ бросает исключение, если клиента отключили."""
try:
mqtt1.on_message = on_message
mqtt1.connect(MQTT1_HOST, MQTT1_PORT, keepalive=60)
mqtt1.subscribe(MQTT_GO_TOPIC)
mqtt1.loop_start()
print(f"[MQTT1] Подключён к {MQTT1_HOST}:{MQTT1_PORT}, подписан на {MQTT_GO_TOPIC}")
info = client.publish(topic, payload, **kw)
rc = getattr(info, 'rc', 0)
if rc != 0:
print(f"[{name}] НЕ доставлено в {topic} (rc={rc}, нет связи с брокером)")
return False
return True
except Exception as e:
print(f"[MQTT1] Ошибка подключения: {e}")
try:
mqtt2.username_pw_set(MQTT2_USER, MQTT2_PASS)
mqtt2.connect(MQTT2_HOST, MQTT2_PORT, keepalive=60)
mqtt2.on_message = on_message
mqtt2.subscribe(MQTT_CMD_TOPIC)
mqtt2.loop_start()
print(f"[MQTT2] Подключён к {MQTT2_HOST}:{MQTT2_PORT}, подписан на {MQTT_CMD_TOPIC}")
except Exception as e:
print(f"[MQTT2] Ошибка подключения: {e}")
print(f"[{name}] Ошибка публикации в {topic}: {e}")
return False
def publish_weight(w):
ts = datetime.datetime.now().strftime('%Y-%m-%dT%H:%M:%S')
payload = f"{w:.1f}"
for client, name in ((mqtt1, 'MQTT1'), (mqtt2, 'MQTT2')):
try:
client.publish(MQTT_WEIGHT_TOPIC, payload)
ensure_mqtt_connected()
for client, name, _h, _p, _s in _MQTT_DEFS:
if _pub(client, name, MQTT_WEIGHT_TOPIC, payload):
print(f"[{name}] → {MQTT_WEIGHT_TOPIC}: {payload} кг [{ts}]")
except Exception as e:
print(f"[{name}] Ошибка публикации: {e}")
# ════════════════════════════════════════════════════════════════════════════════
@@ -390,11 +454,11 @@ def build_status(current_state):
def publish_status(current_state):
payload = json.dumps(build_status(current_state))
for client, name in ((mqtt1, 'MQTT1'), (mqtt2, 'MQTT2')):
try:
client.publish(MQTT_STATUS_TOPIC, payload, qos=0, retain=True)
except Exception as e:
print(f"[{name}] Ошибка публикации статуса: {e}")
ok = []
for client, name, _h, _p, _s in _MQTT_DEFS:
if _pub(client, name, MQTT_STATUS_TOPIC, payload, qos=0, retain=True):
ok.append(name)
return ok
_state_ref = {'value': 'EMPTY'} # шарится с heartbeat-потоком без блокировок (GIL спасает на простой записи)
@@ -406,9 +470,14 @@ def heartbeat_loop():
sent, pending = outbox.flush()
if sent > 0:
print(f"[OUTBOX] Досланных: {sent}, осталось: {pending}")
# Держим MQTT живым (самолечение после разрывов)
ensure_mqtt_connected()
# Публикуем статус (включая outbox_pending)
publish_status(_state_ref['value'])
print(f"[HEARTBEAT] state={_state_ref['value']}")
ok = publish_status(_state_ref['value'])
if ok:
print(f"[HEARTBEAT] state={_state_ref['value']}{', '.join(ok)}")
else:
print(f"[HEARTBEAT] state={_state_ref['value']} — НЕ ОТПРАВЛЕН (нет связи с брокерами)")
except Exception as e:
print(f"[HEARTBEAT] Ошибка: {e}")
time.sleep(HEARTBEAT_INTERVAL)