OPi: store-and-forward outbox для /api/capture (переживает обрыв сети)
- outbox.py: SQLite-очередь + JPEG на диске, at-least-once доставка - scales.py: capture_and_send кладёт в очередь вместо прямого POST, heartbeat_loop каждые 30с дренирует очередь, статус публикует outbox_pending - fw_version -> 2026-07-02-outbox
This commit is contained in:
@@ -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
|
||||||
+79
-23
@@ -1,8 +1,9 @@
|
|||||||
#!/usr/bin/env python3
|
#!/usr/bin/env python3
|
||||||
"""
|
"""
|
||||||
Весовой сервис для Orange Pi — с управлением светофором
|
Весовой сервис для Orange Pi — с управлением светофором + надёжной доставкой
|
||||||
Читает данные с весового терминала A9 через USB-RS232 (WaveShare конвертер)
|
Читает данные с весового терминала A9 через USB-RS232 (WaveShare конвертер)
|
||||||
и публикует вес в MQTT топик scale/weight
|
и публикует вес в MQTT топик scale/weight. При обрыве сети события хранятся
|
||||||
|
локально и досылаются при восстановлении (outbox store-and-forward).
|
||||||
|
|
||||||
Машина состояний:
|
Машина состояний:
|
||||||
EMPTY → машины нет, оба реле выключены
|
EMPTY → машины нет, оба реле выключены
|
||||||
@@ -30,6 +31,13 @@ import base64
|
|||||||
import json
|
import json
|
||||||
import urllib.request
|
import urllib.request
|
||||||
|
|
||||||
|
# ── outbox (store-and-forward, встроенный) ──────────────────────────────────
|
||||||
|
try:
|
||||||
|
from outbox import CaptureOutbox
|
||||||
|
except ImportError:
|
||||||
|
print("[WARNING] outbox модуль не найден, очередь будет отключена")
|
||||||
|
CaptureOutbox = None
|
||||||
|
|
||||||
# ── Серийный порт ────────────────────────────────────────────────────────────
|
# ── Серийный порт ────────────────────────────────────────────────────────────
|
||||||
SERIAL_PORT = '/dev/ttyUSB0'
|
SERIAL_PORT = '/dev/ttyUSB0'
|
||||||
SERIAL_BAUD = 9600
|
SERIAL_BAUD = 9600
|
||||||
@@ -59,7 +67,7 @@ MQTT_STATUS_TOPIC = 'scale/device/status' # heartbeat: публикуем сю
|
|||||||
MQTT_CMD_TOPIC = 'scale/device/cmd' # команды от сервера: ping / reboot / test-shot
|
MQTT_CMD_TOPIC = 'scale/device/cmd' # команды от сервера: ping / reboot / test-shot
|
||||||
|
|
||||||
HEARTBEAT_INTERVAL = 30 # сек между heartbeat
|
HEARTBEAT_INTERVAL = 30 # сек между heartbeat
|
||||||
FIRMWARE_VERSION = '2026-06-30-heartbeat'
|
FIRMWARE_VERSION = '2026-07-02-outbox'
|
||||||
BOOT_TIME = time.time()
|
BOOT_TIME = time.time()
|
||||||
|
|
||||||
# ── Параметры взвешивания ─────────────────────────────────────────────────────
|
# ── Параметры взвешивания ─────────────────────────────────────────────────────
|
||||||
@@ -80,6 +88,9 @@ BLINK_INTERVAL = 1.0 # сек между миганиями
|
|||||||
# ── Событие GO от MQTT ────────────────────────────────────────────────────────
|
# ── Событие GO от MQTT ────────────────────────────────────────────────────────
|
||||||
go_event = threading.Event()
|
go_event = threading.Event()
|
||||||
|
|
||||||
|
# ── Store-and-forward outbox (глобальный) ────────────────────────────────────
|
||||||
|
outbox = None
|
||||||
|
|
||||||
|
|
||||||
# ════════════════════════════════════════════════════════════════════════════════
|
# ════════════════════════════════════════════════════════════════════════════════
|
||||||
# GPIO через sysfs
|
# GPIO через sysfs
|
||||||
@@ -192,12 +203,11 @@ def publish_weight(w):
|
|||||||
|
|
||||||
|
|
||||||
# ════════════════════════════════════════════════════════════════════════════════
|
# ════════════════════════════════════════════════════════════════════════════════
|
||||||
# Захват кадра с камеры и отправка на сервер распознавания
|
# Захват кадра с камеры и отправка на сервер
|
||||||
# ════════════════════════════════════════════════════════════════════════════════
|
# ════════════════════════════════════════════════════════════════════════════════
|
||||||
|
|
||||||
def grab_frame():
|
def grab_frame():
|
||||||
"""Снимает один кадр с RTSP-камеры через ffmpeg. Возвращает JPEG-байты или None."""
|
"""Снимает один кадр с RTSP-камеры через ffmpeg. Возвращает JPEG-байты или None."""
|
||||||
# шлём НАТИВНЫЙ кадр (2560x1440) без ужатия — сервер сам вырежет зону номера и распознает
|
|
||||||
cmd = ['ffmpeg', '-y', '-loglevel', 'error', '-rtsp_transport', 'tcp',
|
cmd = ['ffmpeg', '-y', '-loglevel', 'error', '-rtsp_transport', 'tcp',
|
||||||
'-i', RTSP_URL, '-frames:v', '1', '-update', '1',
|
'-i', RTSP_URL, '-frames:v', '1', '-update', '1',
|
||||||
'-q:v', '2', '-f', 'image2', FRAME_PATH]
|
'-q:v', '2', '-f', 'image2', FRAME_PATH]
|
||||||
@@ -210,27 +220,49 @@ def grab_frame():
|
|||||||
return None
|
return None
|
||||||
|
|
||||||
def capture_and_send(weight, event_id, is_test=False):
|
def capture_and_send(weight, event_id, is_test=False):
|
||||||
"""Кадр + вес → наш VPS. Запускать в отдельном потоке (не блокирует светофор)."""
|
"""
|
||||||
|
Кадр + вес → локальную очередь (outbox).
|
||||||
|
Очередь досылает когда сеть доступна. Запускать в отдельном потоке.
|
||||||
|
"""
|
||||||
img = grab_frame()
|
img = grab_frame()
|
||||||
if not img:
|
if not img:
|
||||||
print(f"[CAPTURE] {event_id}: кадр не получен, пропуск")
|
print(f"[CAPTURE] {event_id}: кадр не получен, пропуск")
|
||||||
return
|
return
|
||||||
try:
|
|
||||||
payload = json.dumps({
|
# Положить в outbox
|
||||||
"event_id": event_id,
|
if outbox:
|
||||||
"device_id": CAPTURE_DEVICE_ID,
|
outbox.enqueue(
|
||||||
"weight": round(weight, 1),
|
event_id=event_id,
|
||||||
"captured_at": datetime.datetime.now().isoformat(),
|
weight=round(weight, 1),
|
||||||
"media_type": "image/jpeg",
|
weight_raw=str(round(weight, 1)),
|
||||||
"photo_b64": base64.b64encode(img).decode(),
|
captured_at=datetime.datetime.now().isoformat(),
|
||||||
"is_test": is_test,
|
jpeg_bytes=img
|
||||||
}).encode()
|
)
|
||||||
req = urllib.request.Request(SERVER_CAPTURE_URL, data=payload,
|
# Сразу попытаемся отправить (если сеть есть)
|
||||||
headers={"Content-Type": "application/json"})
|
sent, pending = outbox.flush()
|
||||||
with urllib.request.urlopen(req, timeout=20) as r:
|
if sent > 0:
|
||||||
print(f"[CAPTURE] {event_id}: отправлено ({len(img)} б), ответ {r.status}")
|
print(f"[CAPTURE] {event_id}: отправлено из очереди сразу")
|
||||||
except Exception as e:
|
elif pending > 0:
|
||||||
print(f"[CAPTURE] {event_id}: ошибка отправки: {e}")
|
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):
|
def build_status(current_state):
|
||||||
total_mb, free_mb = get_mem_mb()
|
total_mb, free_mb = get_mem_mb()
|
||||||
return {
|
status = {
|
||||||
"device_id": CAPTURE_DEVICE_ID,
|
"device_id": CAPTURE_DEVICE_ID,
|
||||||
"ts": datetime.datetime.now().isoformat(),
|
"ts": datetime.datetime.now().isoformat(),
|
||||||
"uptime_s": round(time.time() - BOOT_TIME),
|
"uptime_s": round(time.time() - BOOT_TIME),
|
||||||
@@ -292,6 +324,10 @@ def build_status(current_state):
|
|||||||
"load_avg_1m": get_load_avg(),
|
"load_avg_1m": get_load_avg(),
|
||||||
"gpio_ok": gpio_ok,
|
"gpio_ok": gpio_ok,
|
||||||
}
|
}
|
||||||
|
# Добавить статус очереди, если outbox доступен
|
||||||
|
if outbox:
|
||||||
|
status["outbox_pending"] = outbox.pending_count()
|
||||||
|
return status
|
||||||
|
|
||||||
def publish_status(current_state):
|
def publish_status(current_state):
|
||||||
payload = json.dumps(build_status(current_state))
|
payload = json.dumps(build_status(current_state))
|
||||||
@@ -306,6 +342,12 @@ _state_ref = {'value': 'EMPTY'} # шарится с heartbeat-потоком
|
|||||||
def heartbeat_loop():
|
def heartbeat_loop():
|
||||||
while True:
|
while True:
|
||||||
try:
|
try:
|
||||||
|
# Дренаж очереди каждый heartbeat
|
||||||
|
if outbox:
|
||||||
|
sent, pending = outbox.flush()
|
||||||
|
if sent > 0:
|
||||||
|
print(f"[OUTBOX] Досланных: {sent}, осталось: {pending}")
|
||||||
|
# Публикуем статус (включая outbox_pending)
|
||||||
publish_status(_state_ref['value'])
|
publish_status(_state_ref['value'])
|
||||||
print(f"[HEARTBEAT] state={_state_ref['value']}")
|
print(f"[HEARTBEAT] state={_state_ref['value']}")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
@@ -347,10 +389,24 @@ def parse_weight(data):
|
|||||||
# ════════════════════════════════════════════════════════════════════════════════
|
# ════════════════════════════════════════════════════════════════════════════════
|
||||||
|
|
||||||
def main():
|
def main():
|
||||||
|
global outbox
|
||||||
|
|
||||||
gpio_init()
|
gpio_init()
|
||||||
mqtt_connect()
|
mqtt_connect()
|
||||||
lights_off()
|
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()
|
threading.Thread(target=heartbeat_loop, daemon=True).start()
|
||||||
print(f"[HEARTBEAT] Поток запущен, интервал {HEARTBEAT_INTERVAL}с, топик {MQTT_STATUS_TOPIC}")
|
print(f"[HEARTBEAT] Поток запущен, интервал {HEARTBEAT_INTERVAL}с, топик {MQTT_STATUS_TOPIC}")
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user