Files
weighing-controller/orange_pi/scales.py
T

670 lines
31 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#!/usr/bin/env python3
"""
Весовой сервис для Orange Pi — с управлением светофором + надёжной доставкой
Читает данные с весового терминала A9 через USB-RS232 (WaveShare конвертер)
и публикует вес в MQTT топик scale/weight. При обрыве сети события хранятся
локально и досылаются при восстановлении (outbox store-and-forward).
Машина состояний:
EMPTY → машины нет, оба реле выключены
LOADING → машина на платформе, красный ON, ждём стабилизации
WAIT_GO → вес отправлен, красный ON, ждём GO от сервера (макс 5 мин)
GREEN → GO получен, зелёный ON, ждём пока машина уедет
Установка зависимостей:
pip3 install pyserial paho-mqtt --break-system-packages
Автозапуск:
systemctl enable scales
systemctl start scales
"""
import subprocess
import paho.mqtt.client as mqtt
import re
import time
import datetime
import threading
import select
import os
import base64
import json
import urllib.request
import urllib.error
# ── outbox (store-and-forward, встроенный) ──────────────────────────────────
try:
from outbox import CaptureOutbox
except ImportError:
print("[WARNING] outbox модуль не найден, очередь будет отключена")
CaptureOutbox = None
# ── Серийный порт ────────────────────────────────────────────────────────────
SERIAL_PORT = '/dev/ttyUSB0'
SERIAL_BAUD = 9600
# ── MQTT локальный (сервер с ИИ системой) ────────────────────────────────────
MQTT1_HOST = '192.168.20.9'
MQTT1_PORT = 1883
MQTT1_USER = ''
MQTT1_PASS = ''
# ── MQTT наш VPS ─────────────────────────────────────────────────────────────
MQTT2_HOST = '77.222.43.248'
MQTT2_PORT = 1884
MQTT2_USER = 'esp32'
MQTT2_PASS = 'Esp32Scales#2026'
# ── Камера + распознавание (наш VPS) ─────────────────────────────────────────
CAPTURE_ENABLED = True
RTSP_URL = 'rtsp://iivideo:Vk9mT2xLp7nQ4wBj@192.168.20.3:9784/cameras/124/streaming/main'
SERVER_CAPTURE_URL = 'https://scales.zeroday.su/api/capture'
# ── офлайн-распознавание (Nomeroff shadow) ───────────────────────────────────
NOMEROFF_URL = 'http://192.168.20.4:5010/recognize?crop=1&enhance=0'
LOCAL_PLATE_URL = SERVER_CAPTURE_URL.rsplit('/api/capture', 1)[0] + '/api/local-plate'
CAPTURE_DEVICE_ID = 'scales_opi'
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-07-02-outbox'
BOOT_TIME = time.time()
# ── Параметры взвешивания ─────────────────────────────────────────────────────
MIN_WEIGHT = 2000 # кг — минимум для определения машины на платформе (поднят с 1600: отсекаем ложные срабатывания)
STABLE_DELTA = 50 # кг — допустимый разброс при стабилизации
STABLE_TIME = 10 # сек — время окна стабилизации
# ── Светофор (GPIO через sysfs) ───────────────────────────────────────────────
GPIO_RED = 0 # sysfs номер GPIO для PIN11 (красный)
GPIO_GREEN = 2 # sysfs номер GPIO для PIN13 (зелёный)
RELAY_ACTIVE_LOW = True # True = реле включается при LOW (стандартный модуль с оптопарой)
# ── Таймаут ожидания GO ───────────────────────────────────────────────────────
GO_TIMEOUT = 5 * 60 # 5 минут
BLINK_COUNT = 5 # количество миганий при таймауте
BLINK_INTERVAL = 1.0 # сек между миганиями
# ── Событие GO от MQTT ────────────────────────────────────────────────────────
go_event = threading.Event()
# ── Store-and-forward outbox (глобальный) ────────────────────────────────────
outbox = None
# ════════════════════════════════════════════════════════════════════════════════
# GPIO через sysfs
# ════════════════════════════════════════════════════════════════════════════════
gpio_ok = False
def _gpio_write_raw(pin, value):
with open(f'/sys/class/gpio/gpio{pin}/value', 'w') as f:
f.write('1' if value else '0')
def gpio_init():
global gpio_ok
try:
for pin in (GPIO_RED, GPIO_GREEN):
gpio_path = f'/sys/class/gpio/gpio{pin}'
if not os.path.exists(gpio_path):
with open('/sys/class/gpio/export', 'w') as f:
f.write(str(pin))
time.sleep(0.15)
with open(f'{gpio_path}/direction', 'w') as f:
f.write('out')
# Сразу гасим оба реле
_gpio_write_raw(GPIO_RED, not RELAY_ACTIVE_LOW)
_gpio_write_raw(GPIO_GREEN, not RELAY_ACTIVE_LOW)
gpio_ok = True
print("[GPIO] Инициализация OK (RED=pin11, GREEN=pin13)")
except Exception as e:
print(f"[GPIO] Недоступен, светофор работать не будет: {e}")
def _relay(pin, on):
if not gpio_ok:
return
try:
if RELAY_ACTIVE_LOW:
_gpio_write_raw(pin, not on) # LOW = включено
else:
_gpio_write_raw(pin, on)
except Exception as e:
print(f"[GPIO] Ошибка записи pin{pin}: {e}")
def set_red(on):
_relay(GPIO_RED, on)
print(f"[TRAFFIC] Красный {'ON' if on else 'OFF'}")
def set_green(on):
_relay(GPIO_GREEN, on)
print(f"[TRAFFIC] Зелёный {'ON' if on else 'OFF'}")
def lights_off():
_relay(GPIO_RED, False)
_relay(GPIO_GREEN, False)
print("[TRAFFIC] Оба реле OFF")
def blink_red_async():
"""Мигание красным при таймауте — в отдельном потоке"""
def _blink():
print(f"[TRAFFIC] Мигание красным ({BLINK_COUNT}x)")
for _ in range(BLINK_COUNT):
_relay(GPIO_RED, True)
time.sleep(BLINK_INTERVAL)
_relay(GPIO_RED, False)
time.sleep(BLINK_INTERVAL)
threading.Thread(target=_blink, daemon=True).start()
# ════════════════════════════════════════════════════════════════════════════════
# MQTT
# ════════════════════════════════════════════════════════════════════════════════
def on_message(client, userdata, msg):
print(f"[MQTT] ← {msg.topic}: {msg.payload.decode(errors='ignore')}")
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")
# (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:
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"[{name}] Ошибка публикации в {topic}: {e}")
return False
def publish_weight(w):
ts = datetime.datetime.now().strftime('%Y-%m-%dT%H:%M:%S')
payload = f"{w:.1f}"
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}]")
# ════════════════════════════════════════════════════════════════════════════════
# Захват кадра с камеры и отправка на сервер
# ════════════════════════════════════════════════════════════════════════════════
def grab_frame():
"""Снимает один кадр с RTSP-камеры через ffmpeg. Возвращает JPEG-байты или None."""
cmd = ['ffmpeg', '-y', '-loglevel', 'error', '-rtsp_transport', 'tcp',
'-i', RTSP_URL, '-frames:v', '1', '-update', '1',
'-q:v', '2', '-f', 'image2', FRAME_PATH]
try:
subprocess.run(cmd, timeout=15, check=True)
with open(FRAME_PATH, 'rb') as f:
return f.read()
except Exception as e:
print(f"[CAPTURE] Ошибка захвата кадра: {e}")
return None
def _offline_anpr(img, event_id):
"""Фоновый shadow-вызов Nomeroff + отправка результата на сервер.
Полностью изолирован: любая ошибка/таймаут НЕ влияет на боевой поток взвешивания."""
try:
boundary = '----opi' + str(int(time.time() * 1000))
body = (
('--%s\r\n' % boundary).encode()
+ b'Content-Disposition: form-data; name="photo"; filename="frame.jpg"\r\n'
+ b'Content-Type: image/jpeg\r\n\r\n' + img
+ ('\r\n--%s--\r\n' % boundary).encode()
)
req = urllib.request.Request(
NOMEROFF_URL, data=body,
headers={'Content-Type': 'multipart/form-data; boundary=%s' % boundary})
with urllib.request.urlopen(req, timeout=8) as r:
resp = json.loads(r.read().decode())
except Exception as e:
print(f"[ANPR] {event_id}: Nomeroff недоступен: {e}")
return
if not resp.get('found') or not resp.get('plate'):
print(f"[ANPR] {event_id}: офлайн не распознал")
return
payload = json.dumps({
"event_id": event_id,
"local_plate": resp.get('plate'),
"local_confidence": resp.get('confidence'),
"local_model": "nomeroff",
"local_elapsed": resp.get('elapsed_sec'),
}).encode()
# ретрай на 404: событие могло ещё не создаться на сервере (capture летит параллельно)
for attempt in range(3):
try:
req = urllib.request.Request(
LOCAL_PLATE_URL, data=payload,
headers={"Content-Type": "application/json"})
with urllib.request.urlopen(req, timeout=10) as r:
print(f"[ANPR] {event_id}: офлайн {resp.get('plate')} → сервер, ответ {r.status}")
return
except urllib.error.HTTPError as e:
if e.code == 404 and attempt < 2:
time.sleep(1.2)
continue
print(f"[ANPR] {event_id}: ошибка отправки: {e}")
return
except Exception as e:
print(f"[ANPR] {event_id}: ошибка отправки: {e}")
return
def capture_and_send(weight, event_id, is_test=False):
"""
Кадр + вес → локальную очередь (outbox).
Очередь досылает когда сеть доступна. Запускать в отдельном потоке.
"""
img = grab_frame()
if not img:
print(f"[CAPTURE] {event_id}: кадр не получен, пропуск")
return
# shadow-распознавание Nomeroff в фоне (не для тестовых снимков камеры)
if not is_test:
threading.Thread(target=_offline_anpr, args=(img, event_id), daemon=True).start()
# Положить в outbox
if outbox:
outbox.enqueue(
event_id=event_id,
weight=round(weight, 1),
weight_raw=str(round(weight, 1)),
captured_at=datetime.datetime.now().isoformat(),
jpeg_bytes=img
)
# Сразу попытаемся отправить (если сеть есть)
sent, pending = outbox.flush()
if sent > 0:
print(f"[CAPTURE] {event_id}: отправлено из очереди сразу")
elif pending > 0:
print(f"[CAPTURE] {event_id}: в очереди ({pending} ожидающих)")
else:
# fallback: прямой POST если outbox не доступен (на случай, если модуль не импортировался)
print("[CAPTURE] outbox недоступен, попытка прямой отправки...")
try:
payload = json.dumps({
"event_id": event_id,
"device_id": CAPTURE_DEVICE_ID,
"weight": round(weight, 1),
"captured_at": datetime.datetime.now().isoformat(),
"media_type": "image/jpeg",
"photo_b64": base64.b64encode(img).decode(),
"is_test": is_test,
}).encode()
req = urllib.request.Request(SERVER_CAPTURE_URL, data=payload,
headers={"Content-Type": "application/json"})
with urllib.request.urlopen(req, timeout=20) as r:
print(f"[CAPTURE] {event_id}: отправлено ({len(img)} б), ответ {r.status}")
except Exception as e:
print(f"[CAPTURE] {event_id}: ошибка отправки: {e}")
# ════════════════════════════════════════════════════════════════════════════════
# Парсер протокола 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()
status = {
"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,
}
# Добавить статус очереди, если outbox доступен
if outbox:
status["outbox_pending"] = outbox.pending_count()
return status
def publish_status(current_state):
payload = json.dumps(build_status(current_state))
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 спасает на простой записи)
def heartbeat_loop():
while True:
try:
# Дренаж очереди каждый heartbeat
if outbox:
sent, pending = outbox.flush()
if sent > 0:
print(f"[OUTBOX] Досланных: {sent}, осталось: {pending}")
# Держим MQTT живым (самолечение после разрывов)
ensure_mqtt_connected()
# Публикуем статус (включая outbox_pending)
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)
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)}", True), daemon=True).start()
else:
print(f"[CMD] Неизвестная команда: {cmd}")
def parse_weight(data):
"""
Формат A9: STX + знак + 7 цифр + контрольная сумма + ETX
Пример: \x02+00176001B\x03 → 1760.0 кг
"""
m = re.search(rb'\x02([+-]\d{8}[A-Z0-9])\x03', data)
if m:
s = m.group(1).decode('ascii')
digits = s[1:8]
return float(digits[:-1] + '.' + digits[-1])
return None
# ════════════════════════════════════════════════════════════════════════════════
# Главный цикл
# ════════════════════════════════════════════════════════════════════════════════
def main():
global outbox
gpio_init()
mqtt_connect()
lights_off()
# Инициализация outbox
if CaptureOutbox:
try:
outbox = CaptureOutbox(
server_url=SERVER_CAPTURE_URL.rsplit('/api/capture', 1)[0],
device_id=CAPTURE_DEVICE_ID
)
print("[OUTBOX] Инициализирована")
except Exception as e:
print(f"[OUTBOX] Ошибка инициализации: {e}")
outbox = None
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),
'cs8', '-cstopb', '-parenb'], check=False)
proc = subprocess.Popen(
['cat', SERIAL_PORT],
stdout=subprocess.PIPE,
stderr=subprocess.DEVNULL
)
state = 'EMPTY'
last_weights = []
buf = b""
go_sent_time = None
last_logged_w = None # последнее залогированное значение веса
while True:
# ── Неблокирующее чтение с таймаутом 100 мс ──────────────────────────
r, _, _ = select.select([proc.stdout], [], [], 0.1)
if r:
try:
chunk = proc.stdout.read(12)
if not chunk:
time.sleep(0.1)
continue
except Exception:
time.sleep(0.1)
continue
buf += chunk
# Извлекаем все полные фреймы из буфера
while b'\x02' in buf and b'\x03' in buf:
s = buf.find(b'\x02')
e = buf.find(b'\x03', s)
if s == -1 or e == -1:
break
frame = buf[s:e+1]
buf = buf[e+1:]
w = parse_weight(frame)
if w is None:
continue
now = time.time()
if w != last_logged_w:
print(f"[SCALE] {w:.1f} кг [{state}]")
last_logged_w = w
# ── EMPTY: нет машины ─────────────────────────────────────────
if state == 'EMPTY':
if w >= MIN_WEIGHT:
state = 'LOADING'
last_weights = []
set_red(True)
print(f"[STATE] → LOADING: машина на платформе ({w:.1f} кг)")
# ── LOADING: машина заехала, ждём стабилизации ─────────────────
elif state == 'LOADING':
if w < MIN_WEIGHT / 2:
state = 'EMPTY'
last_weights = []
lights_off()
print("[STATE] → EMPTY: машина уехала до стабилизации")
continue
last_weights.append((now, w))
last_weights = [(t, v) for t, v in last_weights
if now - t <= STABLE_TIME]
if len(last_weights) >= 5:
vals = [v for _, v in last_weights]
if max(vals) - min(vals) <= STABLE_DELTA:
avg = sum(vals) / len(vals)
print(f"[STATE] → WAIT_GO: вес стабилен {avg:.1f} кг, отправляем")
event_id = f"opi-{int(time.time()*1000)}"
publish_weight(avg)
if CAPTURE_ENABLED:
threading.Thread(target=capture_and_send,
args=(avg, event_id), daemon=True).start()
state = 'WAIT_GO'
go_event.clear()
go_sent_time = time.time()
# ── WAIT_GO: ждём команду GO от сервера ───────────────────────
elif state == 'WAIT_GO':
if w < MIN_WEIGHT / 2:
state = 'EMPTY'
lights_off()
go_event.clear()
go_sent_time = None
print("[STATE] → EMPTY: машина уехала без GO")
# ── GREEN: GO получен, ждём пока машина уедет ─────────────────
elif state == 'GREEN':
if w < MIN_WEIGHT / 2:
state = 'EMPTY'
lights_off()
print("[STATE] → EMPTY: машина проехала, сброс")
# ── Проверка GO и таймаута (независимо от данных весов) ───────────────
_state_ref['value'] = state
if state == 'WAIT_GO':
if go_event.is_set():
go_event.clear()
state = 'GREEN'
set_red(False)
set_green(True)
go_sent_time = None
print("[STATE] → GREEN: GO получен!")
elif go_sent_time and (time.time() - go_sent_time) > GO_TIMEOUT:
print("[STATE] → EMPTY: таймаут 5 мин, GO не получен")
state = 'EMPTY'
go_sent_time = None
go_event.clear()
blink_red_async() # мигаем красным для оператора
# После мигания оба реле погаснут внутри blink, но зелёный
# мог остаться — явно гасим после задержки
def _cleanup():
time.sleep(BLINK_COUNT * BLINK_INTERVAL * 2 + 0.5)
lights_off()
threading.Thread(target=_cleanup, daemon=True).start()
if __name__ == '__main__':
main()