From a4f8d5074ed5dddd299dd52cac6f90f6886c212c Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 30 Jun 2026 14:57:40 +0300 Subject: [PATCH] Add MQTT heartbeat + remote commands (ping/reboot/test-shot) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - scale/device/status: публикует раз в 30с uptime, state, память, диск, temp SoC, версию - scale/device/cmd: слушает команды ping/reboot/test-shot от сервера - использует только stdlib (/proc, /sys) — без новых зависимостей --- orange_pi/scales.py | 108 +++++++++++++++++++++++++++++++++++++++++++- 1 file changed, 107 insertions(+), 1 deletion(-) diff --git a/orange_pi/scales.py b/orange_pi/scales.py index baa7bdf..15adc59 100644 --- a/orange_pi/scales.py +++ b/orange_pi/scales.py @@ -55,6 +55,12 @@ FRAME_PATH = '/tmp/scale_frame.jpg' MQTT_WEIGHT_TOPIC = 'scale/weight' MQTT_GO_TOPIC = 'scale/traffic/go' +MQTT_STATUS_TOPIC = 'scale/device/status' # heartbeat: публикуем сюда раз в HEARTBEAT_INTERVAL сек +MQTT_CMD_TOPIC = 'scale/device/cmd' # команды от сервера: ping / reboot / test-shot + +HEARTBEAT_INTERVAL = 30 # сек между heartbeat +FIRMWARE_VERSION = '2026-06-30-heartbeat' +BOOT_TIME = time.time() # ── Параметры взвешивания ───────────────────────────────────────────────────── MIN_WEIGHT = 2000 # кг — минимум для определения машины на платформе (поднят с 1600: отсекаем ложные срабатывания) @@ -149,6 +155,8 @@ def on_message(client, userdata, msg): if msg.topic == MQTT_GO_TOPIC: print("[MQTT] GO получен — переключаем светофор") go_event.set() + elif msg.topic == MQTT_CMD_TOPIC: + handle_cmd(msg.payload.decode(errors='ignore')) mqtt1 = mqtt.Client(client_id="scales_opi_1") mqtt2 = mqtt.Client(client_id="scales_opi_2") @@ -165,8 +173,10 @@ def mqtt_connect(): 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}") + print(f"[MQTT2] Подключён к {MQTT2_HOST}:{MQTT2_PORT}, подписан на {MQTT_CMD_TOPIC}") except Exception as e: print(f"[MQTT2] Ошибка подключения: {e}") @@ -226,6 +236,98 @@ def capture_and_send(weight, event_id): # Парсер протокола A9 # ════════════════════════════════════════════════════════════════════════════════ +# ════════════════════════════════════════════════════════════════════════════════ +# Heartbeat / телеметрия устройства +# ════════════════════════════════════════════════════════════════════════════════ + +def get_mem_mb(): + try: + info = {} + with open('/proc/meminfo') as f: + for line in f: + k, v = line.split(':') + info[k] = int(v.strip().split()[0]) # kB + total = info.get('MemTotal', 0) // 1024 + free = info.get('MemAvailable', info.get('MemFree', 0)) // 1024 + return total, free + except Exception: + return None, None + +def get_disk_pct(): + try: + st = os.statvfs('/') + total = st.f_blocks * st.f_frsize + free = st.f_bavail * st.f_frsize + used_pct = round((1 - free / total) * 100, 1) if total else None + return used_pct + except Exception: + return None + +def get_cpu_temp_c(): + try: + with open('/sys/class/thermal/thermal_zone0/temp') as f: + return round(int(f.read().strip()) / 1000, 1) + except Exception: + return None + +def get_load_avg(): + try: + return os.getloadavg()[0] + except Exception: + return None + +def build_status(current_state): + total_mb, free_mb = get_mem_mb() + return { + "device_id": CAPTURE_DEVICE_ID, + "ts": datetime.datetime.now().isoformat(), + "uptime_s": round(time.time() - BOOT_TIME), + "fw_version": FIRMWARE_VERSION, + "state": current_state, + "mem_total_mb": total_mb, + "mem_free_mb": free_mb, + "disk_used_pct": get_disk_pct(), + "cpu_temp_c": get_cpu_temp_c(), + "load_avg_1m": get_load_avg(), + "gpio_ok": gpio_ok, + } + +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}") + +_state_ref = {'value': 'EMPTY'} # шарится с heartbeat-потоком без блокировок (GIL спасает на простой записи) + +def heartbeat_loop(): + while True: + try: + publish_status(_state_ref['value']) + print(f"[HEARTBEAT] state={_state_ref['value']}") + except Exception as e: + print(f"[HEARTBEAT] Ошибка: {e}") + time.sleep(HEARTBEAT_INTERVAL) + +def handle_cmd(raw): + try: + cmd = json.loads(raw).get('cmd', raw.strip()) + except Exception: + cmd = raw.strip() + print(f"[CMD] Получена команда: {cmd}") + if cmd == 'ping': + publish_status(_state_ref['value']) + elif cmd == 'reboot': + print("[CMD] Перезагрузка по команде сервера...") + threading.Thread(target=lambda: (time.sleep(1), subprocess.run(['reboot'])), daemon=True).start() + elif cmd == 'test-shot': + threading.Thread(target=capture_and_send, args=(0.0, f"test-{int(time.time()*1000)}"), daemon=True).start() + else: + print(f"[CMD] Неизвестная команда: {cmd}") + + def parse_weight(data): """ Формат A9: STX + знак + 7 цифр + контрольная сумма + ETX @@ -248,6 +350,9 @@ def main(): mqtt_connect() lights_off() + threading.Thread(target=heartbeat_loop, daemon=True).start() + print(f"[HEARTBEAT] Поток запущен, интервал {HEARTBEAT_INTERVAL}с, топик {MQTT_STATUS_TOPIC}") + print("[BOOT] Sistema gotova. Zhdu mashinu...") subprocess.run(['stty', '-F', SERIAL_PORT, 'raw', str(SERIAL_BAUD), @@ -351,6 +456,7 @@ def main(): print("[STATE] → EMPTY: машина проехала, сброс") # ── Проверка GO и таймаута (независимо от данных весов) ─────────────── + _state_ref['value'] = state if state == 'WAIT_GO': if go_event.is_set(): go_event.clear()