Add MQTT heartbeat + remote commands (ping/reboot/test-shot)
- scale/device/status: публикует раз в 30с uptime, state, память, диск, temp SoC, версию - scale/device/cmd: слушает команды ping/reboot/test-shot от сервера - использует только stdlib (/proc, /sys) — без новых зависимостей
This commit is contained in:
+107
-1
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user