From 3c325cbfc194676cd1f2a5cc3443c8bfdc99fca6 Mon Sep 17 00:00:00 2001 From: Alexey Pavlov Date: Thu, 2 Jul 2026 16:00:21 +0300 Subject: [PATCH] =?UTF-8?q?OPi:=20store-and-forward=20outbox=20=D0=B4?= =?UTF-8?q?=D0=BB=D1=8F=20/api/capture=20(=D0=BF=D0=B5=D1=80=D0=B5=D0=B6?= =?UTF-8?q?=D0=B8=D0=B2=D0=B0=D0=B5=D1=82=20=D0=BE=D0=B1=D1=80=D1=8B=D0=B2?= =?UTF-8?q?=20=D1=81=D0=B5=D1=82=D0=B8)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - outbox.py: SQLite-очередь + JPEG на диске, at-least-once доставка - scales.py: capture_and_send кладёт в очередь вместо прямого POST, heartbeat_loop каждые 30с дренирует очередь, статус публикует outbox_pending - fw_version -> 2026-07-02-outbox --- orange_pi/outbox.py | 189 ++++++++++++++++++++++++++++++++++++++++++++ orange_pi/scales.py | 102 ++++++++++++++++++------ 2 files changed, 268 insertions(+), 23 deletions(-) create mode 100644 orange_pi/outbox.py diff --git a/orange_pi/outbox.py b/orange_pi/outbox.py new file mode 100644 index 0000000..b289cb2 --- /dev/null +++ b/orange_pi/outbox.py @@ -0,0 +1,189 @@ +""" +outbox.py — надёжная доставка взвешиваний на сервер (store-and-forward). + +При обрыве сети OPi↔сервер события не теряются: кладём JPEG на диск + метаданные +в SQLite, и досылаем при восстановлении связи. Сервер дедупит по event_id (upsert). +Семантика: at-least-once (повторная отправка безопасна). + +Зависимости: только stdlib (sqlite3, urllib, json). +""" + +import os +import time +import base64 +import sqlite3 +import threading +import logging +import json +import urllib.request +import urllib.error + +log = logging.getLogger("outbox") + + +class CaptureOutbox: + def __init__(self, server_url, device_id, + db_path="/root/scales_outbox/outbox.db", + photo_dir="/root/scales_outbox/photos", + timeout=15, max_attempts=20): + self.url = server_url.rstrip("/") + "/api/capture" + self.device_id = device_id + self.timeout = timeout + self.max_attempts = max_attempts + self.photo_dir = photo_dir + + os.makedirs(os.path.dirname(db_path), exist_ok=True) + os.makedirs(photo_dir, exist_ok=True) + + self._lock = threading.Lock() # защита БД (heartbeat + main) + self._flush_lock = threading.Lock() # одна flush за раз + self._db = sqlite3.connect(db_path, check_same_thread=False) + self._db.execute("PRAGMA journal_mode=WAL") + self._db.execute(""" + CREATE TABLE IF NOT EXISTS outbox ( + event_id TEXT PRIMARY KEY, + weight REAL, + weight_raw TEXT, + captured_at TEXT, + photo_path TEXT NOT NULL, + created_at TEXT NOT NULL, + attempts INTEGER NOT NULL DEFAULT 0, + status TEXT NOT NULL DEFAULT 'pending' -- pending | failed + )""") + self._db.commit() + self._reconcile() + + def _reconcile(self): + """После рестарта убрать осиротевшие JPEG.""" + with self._lock: + known = {r[0] for r in self._db.execute("SELECT photo_path FROM outbox")} + for name in os.listdir(self.photo_dir): + p = os.path.join(self.photo_dir, name) + if p not in known: + try: + os.remove(p) + except OSError: + pass + + # ── публичный API ───────────────────────────────────────────── + def enqueue(self, event_id, weight, weight_raw, captured_at, jpeg_bytes): + """Положить взвешивание в очередь. JPEG на диск, метаданные в БД.""" + photo_path = os.path.join(self.photo_dir, f"{event_id}.jpg") + with open(photo_path, "wb") as f: + f.write(jpeg_bytes) + with self._lock: + self._db.execute( + "INSERT OR REPLACE INTO outbox " + "(event_id, weight, weight_raw, captured_at, photo_path, created_at, attempts, status) " + "VALUES (?,?,?,?,?,?,0,'pending')", + (event_id, weight, weight_raw, captured_at, photo_path, + time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()))) + self._db.commit() + log.info("enqueued %s (%.1f kg)", event_id, weight) + + def pending_count(self): + """Количество событий в очереди.""" + with self._lock: + return self._db.execute( + "SELECT COUNT(*) FROM outbox WHERE status='pending'").fetchone()[0] + + def flush(self): + """Досылает pending-события (старые первыми). Вернёт (отправлено, осталось).""" + if not self._flush_lock.acquire(blocking=False): + return (0, self.pending_count()) + sent = 0 + try: + while True: + with self._lock: + row = self._db.execute( + "SELECT event_id, weight, weight_raw, captured_at, photo_path, attempts " + "FROM outbox WHERE status='pending' ORDER BY rowid LIMIT 1").fetchone() + if row is None: + break + event_id, weight, weight_raw, captured_at, photo_path, attempts = row + try: + with open(photo_path, "rb") as f: + b64 = base64.b64encode(f.read()).decode() + except OSError: + self._drop(event_id) + continue + + payload = json.dumps({ + "event_id": event_id, + "weight": weight, + "weight_raw": weight_raw, + "device_id": self.device_id, + "captured_at": captured_at, + "photo_b64": b64, + "media_type": "image/jpeg", + }).encode() + + try: + req = urllib.request.Request( + self.url, + data=payload, + headers={"Content-Type": "application/json"}) + with urllib.request.urlopen(req, timeout=self.timeout) as r: + resp = r.read().decode() + if _ok_response(resp): + self._delete(event_id) + sent += 1 + log.info("delivered %s", event_id) + else: + log.warning("server rejected %s: %s", event_id, resp[:100]) + self._park(event_id) + except urllib.error.HTTPError as e: + if 400 <= e.code < 500: + log.error("rejected %s: %d", event_id, e.code) + self._park(event_id) + else: + self._bump(event_id, attempts) + break + except (urllib.error.URLError, OSError, TimeoutError) as e: + log.warning("network unavailable (%s), pending=%d", + type(e).__name__, self.pending_count()) + break + finally: + self._flush_lock.release() + return (sent, self.pending_count()) + + # ── внутреннее ──────────────────────────────────────────────── + def _delete(self, event_id): + with self._lock: + row = self._db.execute("SELECT photo_path FROM outbox WHERE event_id=?", + (event_id,)).fetchone() + self._db.execute("DELETE FROM outbox WHERE event_id=?", (event_id,)) + self._db.commit() + if row: + try: + os.remove(row[0]) + except OSError: + pass + + def _drop(self, event_id): + """Удалить событие (фото потеряно).""" + with self._lock: + self._db.execute("DELETE FROM outbox WHERE event_id=?", (event_id,)) + self._db.commit() + + def _park(self, event_id): + """Паркировать событие (не отправлять больше).""" + with self._lock: + self._db.execute("UPDATE outbox SET status='failed' WHERE event_id=?", (event_id,)) + self._db.commit() + + def _bump(self, event_id, attempts): + """Увеличить счётчик попыток, паркировать если превышен лимит.""" + status = "failed" if attempts + 1 >= self.max_attempts else "pending" + with self._lock: + self._db.execute("UPDATE outbox SET attempts=attempts+1, status=? WHERE event_id=?", + (status, event_id)) + self._db.commit() + + +def _ok_response(text): + """Проверить, что в ответе есть 'ok': true.""" + try: + return bool(json.loads(text).get("ok")) + except Exception: + return False diff --git a/orange_pi/scales.py b/orange_pi/scales.py index 72e0ab2..14d5486 100644 --- a/orange_pi/scales.py +++ b/orange_pi/scales.py @@ -1,8 +1,9 @@ #!/usr/bin/env python3 """ -Весовой сервис для Orange Pi — с управлением светофором +Весовой сервис для Orange Pi — с управлением светофором + надёжной доставкой Читает данные с весового терминала A9 через USB-RS232 (WaveShare конвертер) -и публикует вес в MQTT топик scale/weight +и публикует вес в MQTT топик scale/weight. При обрыве сети события хранятся +локально и досылаются при восстановлении (outbox store-and-forward). Машина состояний: EMPTY → машины нет, оба реле выключены @@ -30,6 +31,13 @@ import base64 import json import urllib.request +# ── outbox (store-and-forward, встроенный) ────────────────────────────────── +try: + from outbox import CaptureOutbox +except ImportError: + print("[WARNING] outbox модуль не найден, очередь будет отключена") + CaptureOutbox = None + # ── Серийный порт ──────────────────────────────────────────────────────────── SERIAL_PORT = '/dev/ttyUSB0' SERIAL_BAUD = 9600 @@ -59,7 +67,7 @@ MQTT_STATUS_TOPIC = 'scale/device/status' # heartbeat: публикуем сю MQTT_CMD_TOPIC = 'scale/device/cmd' # команды от сервера: ping / reboot / test-shot HEARTBEAT_INTERVAL = 30 # сек между heartbeat -FIRMWARE_VERSION = '2026-06-30-heartbeat' +FIRMWARE_VERSION = '2026-07-02-outbox' BOOT_TIME = time.time() # ── Параметры взвешивания ───────────────────────────────────────────────────── @@ -80,6 +88,9 @@ BLINK_INTERVAL = 1.0 # сек между миганиями # ── Событие GO от MQTT ──────────────────────────────────────────────────────── go_event = threading.Event() +# ── Store-and-forward outbox (глобальный) ──────────────────────────────────── +outbox = None + # ════════════════════════════════════════════════════════════════════════════════ # GPIO через sysfs @@ -192,12 +203,11 @@ def publish_weight(w): # ════════════════════════════════════════════════════════════════════════════════ -# Захват кадра с камеры и отправка на сервер распознавания +# Захват кадра с камеры и отправка на сервер # ════════════════════════════════════════════════════════════════════════════════ def grab_frame(): """Снимает один кадр с RTSP-камеры через ffmpeg. Возвращает JPEG-байты или None.""" - # шлём НАТИВНЫЙ кадр (2560x1440) без ужатия — сервер сам вырежет зону номера и распознает cmd = ['ffmpeg', '-y', '-loglevel', 'error', '-rtsp_transport', 'tcp', '-i', RTSP_URL, '-frames:v', '1', '-update', '1', '-q:v', '2', '-f', 'image2', FRAME_PATH] @@ -210,27 +220,49 @@ def grab_frame(): return None def capture_and_send(weight, event_id, is_test=False): - """Кадр + вес → наш VPS. Запускать в отдельном потоке (не блокирует светофор).""" + """ + Кадр + вес → локальную очередь (outbox). + Очередь досылает когда сеть доступна. Запускать в отдельном потоке. + """ img = grab_frame() if not img: print(f"[CAPTURE] {event_id}: кадр не получен, пропуск") return - 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}") + + # Положить в 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}") # ════════════════════════════════════════════════════════════════════════════════ @@ -279,7 +311,7 @@ def get_load_avg(): def build_status(current_state): total_mb, free_mb = get_mem_mb() - return { + status = { "device_id": CAPTURE_DEVICE_ID, "ts": datetime.datetime.now().isoformat(), "uptime_s": round(time.time() - BOOT_TIME), @@ -292,6 +324,10 @@ def build_status(current_state): "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)) @@ -306,6 +342,12 @@ _state_ref = {'value': 'EMPTY'} # шарится с heartbeat-потоком def heartbeat_loop(): while True: try: + # Дренаж очереди каждый heartbeat + if outbox: + sent, pending = outbox.flush() + if sent > 0: + print(f"[OUTBOX] Досланных: {sent}, осталось: {pending}") + # Публикуем статус (включая outbox_pending) publish_status(_state_ref['value']) print(f"[HEARTBEAT] state={_state_ref['value']}") except Exception as e: @@ -347,10 +389,24 @@ def parse_weight(data): # ════════════════════════════════════════════════════════════════════════════════ 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}")