From 94044fde752b0a048128c045ad67bc1b0be030f2 Mon Sep 17 00:00:00 2001 From: claude Date: Wed, 15 Jul 2026 07:38:35 +0300 Subject: [PATCH] =?UTF-8?q?OPi:=20=D0=B0=D0=B2=D1=82=D0=BE-=D1=80=D0=B5?= =?UTF-8?q?=D0=BA=D0=BE=D0=BD=D0=BD=D0=B5=D0=BA=D1=82=20MQTT=20+=20=D0=BF?= =?UTF-8?q?=D1=80=D0=BE=D0=B2=D0=B5=D1=80=D0=BA=D0=B0=20rc=20=D1=83=20publ?= =?UTF-8?q?ish=20(heartbeat=20=D0=B8=20=D0=B2=D0=B5=D1=81=20=D0=BC=D0=BE?= =?UTF-8?q?=D0=BB=D1=87=D0=B0=20=D1=82=D0=B5=D1=80=D1=8F=D0=BB=D0=B8=D1=81?= =?UTF-8?q?=D1=8C=20=D1=81=2010.07)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- CHANGELOG.md | 7 +++ orange_pi/scales.py | 123 ++++++++++++++++++++++++++++++++++---------- 2 files changed, 103 insertions(+), 27 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 4c1f9dc..5b24c5d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -93,3 +93,10 @@ - grab_frame: убран scale=1280, отдаём полный кадр; сервер сам кропает зону номера, делает автоконтраст и распознаёт жёлтые/белые номера ## 2026-06-26 — порог срабатывания MIN_WEIGHT 1600 -> 2000 кг (отсечь ложные заезды) + +## 2026-07-15 — MQTT: авто-реконнект и честная публикация +- publish() теперь проверяет rc (paho НЕ бросает исключение при отсутствии связи, qos=0 молча терялся) +- on_connect/on_disconnect + reconnect_delay_set(1..60с); подписки восстанавливаются после каждого реконнекта +- ensure_mqtt_connected() в heartbeat_loop и перед publish_weight — самолечение канала +- heartbeat пишет в лог, дошёл ли статус ("НЕ ОТПРАВЛЕН (нет связи с брокерами)") +- Причина: 10.07 оба MQTT-канала (к нам :1884 и к ребятам 192.168.20.9:1883) отвалились и не переподключились; фото шли по HTTP, поэтому проблема была не видна diff --git a/orange_pi/scales.py b/orange_pi/scales.py index db31a68..37a2fae 100644 --- a/orange_pi/scales.py +++ b/orange_pi/scales.py @@ -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)