diff --git a/api/db.py b/api/db.py index 1190459..5528962 100644 --- a/api/db.py +++ b/api/db.py @@ -160,6 +160,29 @@ class Database: self.conn.commit() + # command_queue — для raw-ретранслятора (gunicorn-safe) + self.conn.execute(""" + CREATE TABLE IF NOT EXISTS command_queue ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + device_id TEXT NOT NULL, + seq INTEGER NOT NULL, + cmd TEXT NOT NULL, + timeout_ms INTEGER DEFAULT 500, + drain_first INTEGER DEFAULT 0, + status TEXT NOT NULL DEFAULT 'pending', + raw_response TEXT, + elapsed_ms INTEGER, + prompt INTEGER, + error TEXT, + created_at TEXT NOT NULL DEFAULT (datetime('now')), + sent_at TEXT, + responded_at TEXT + ); + CREATE INDEX IF NOT EXISTS idx_cq_device_status_seq + ON command_queue(device_id, status, seq); + """) + self.conn.commit() + # ── sessions ────────────────────────────────────────── def get_cached_response(self, request_id: str) -> dict | None: diff --git a/api/raw_elm.py b/api/raw_elm.py index 9b2500c..6f0e017 100644 --- a/api/raw_elm.py +++ b/api/raw_elm.py @@ -23,23 +23,16 @@ import threading import time from flask import jsonify, request, Blueprint +from api.db import Database + logger = logging.getLogger("elmer.raw_api") # ══════════════════════════════════════════════════════════ -# Глобальное состояние +# Локальное состояние (только для прямого подключения ELM) # ══════════════════════════════════════════════════════════ _raw_mode = False -_raw_elm = None # локальный RawELM - -# Очередь команд (удалённый режим) -_lock = threading.Lock() -_pending_cmd: dict | None = None # команда, которую ждёт Android -_pending_seq: int = 0 # номер последней команды -_last_response: dict | None = None # последний ответ от ELM327 -_device_ready: bool = False # Android подключён и готов -_device_info: dict = {} # информация об устройстве (из hello) -_history: list[dict] = [] # история команд через Android +_raw_elm = None # локальный RawELM (прямое подключение к серверу) def is_raw_mode() -> bool: @@ -119,7 +112,14 @@ def raw_available(): @bp.route("/api/v1/elm/raw/log", methods=["GET"]) def raw_log(): if not _raw_elm: - return jsonify({"log": _history, "count": len(_history)}) + with Database() as db: + rows = db.conn.execute( + "SELECT seq, cmd, raw_response, elapsed_ms, prompt, error FROM command_queue WHERE status='done' ORDER BY id DESC LIMIT 50" + ).fetchall() + history = [{"seq": r["seq"], "cmd": r["cmd"], "raw": r["raw_response"], + "elapsed_ms": r["elapsed_ms"], "prompt": r["prompt"], "error": r["error"]} + for r in rows] + return jsonify({"log": history, "count": len(history)}) return jsonify({"log": _raw_elm.log, "count": len(_raw_elm.log)}) @@ -130,185 +130,163 @@ def raw_mode_control(): data = request.get_json(silent=True) or {} on = data.get("raw_mode", False) set_raw_mode(on) - return jsonify({"raw_mode": _raw_mode, "has_local_elm": _raw_elm is not None, - "device_ready": _device_ready}) - return jsonify({"raw_mode": _raw_mode, "has_local_elm": _raw_elm is not None, - "device_ready": _device_ready}) + with Database() as db: + device = db.conn.execute( + "SELECT device_id FROM command_queue ORDER BY id DESC LIMIT 1" + ).fetchone() + return jsonify({ + "raw_mode": _raw_mode, + "has_local_elm": _raw_elm is not None, + "device_ready": device is not None, + }) # ══════════════════════════════════════════════════════════ -# УДАЛЁННЫЙ РЕЖИМ — Android-ретранслятор +# УДАЛЁННЫЙ РЕЖИМ — Android-ретранслятор (SQLite-очередь) # ══════════════════════════════════════════════════════════ @bp.route("/api/v1/elm/raw/hello", methods=["POST"]) def raw_hello(): - """Android сообщает: «я подключился к ELM327, готов принимать команды». - - Body: { - "device_id": "android-xyz", - "elm_version": "ELM327 v1.5", - "protocol": "A4", - "voltage": "12.3V" - } - """ - global _device_ready, _device_info, _pending_cmd, _pending_seq, _last_response + """Android: «я подключился, готов принимать команды».""" data = request.get_json(silent=True) or {} - with _lock: - _device_ready = True - _device_info = { - "device_id": data.get("device_id", "unknown"), - "elm_version": data.get("elm_version", "?"), - "protocol": data.get("protocol", "?"), - "voltage": data.get("voltage", "?"), - "connected_at": time.time(), - } - _pending_cmd = None - _pending_seq = 0 - _last_response = None - logger.info(f"RawELM: device ready — {_device_info['device_id']} " - f"({_device_info['elm_version']}, proto {_device_info['protocol']})") + device_id = data.get("device_id", "unknown") + logger.info(f"RawELM: device ready — {device_id} " + f"({data.get('elm_version', '?')}, proto {data.get('protocol', '?')})") return jsonify({"ok": True, "seq": 0}) @bp.route("/api/v1/elm/raw/cmd", methods=["POST"]) def raw_enqueue_cmd(): - """Copilot: поставить команду в очередь для Android. - - Body: { - "cmd": "0105", - "timeout_ms": 500, - "drain_first": false - } - """ - global _pending_cmd, _pending_seq + """Copilot: поставить команду в очередь.""" data = request.get_json(silent=True) if not data or "cmd" not in data: return jsonify({"error": "missing 'cmd'"}), 400 - cmd = data["cmd"].strip() if not cmd: return jsonify({"error": "empty cmd"}), 400 - with _lock: - _pending_seq += 1 - _pending_cmd = { - "cmd": cmd, - "timeout_ms": data.get("timeout_ms", 500), - "drain_first": data.get("drain_first", False), - "seq": _pending_seq, - } - logger.info(f"RawELM: enqueued #{_pending_seq} → {cmd}") - return jsonify({"ok": True, "seq": _pending_seq, "cmd": cmd}) + device_id = data.get("device_id", "unknown") + tmo = data.get("timeout_ms", 500) + drain = 1 if data.get("drain_first") else 0 + + with Database() as db: + cur = db.conn.execute("SELECT COALESCE(MAX(seq), 0) + 1 FROM command_queue WHERE device_id = ?", (device_id,)) + seq = cur.fetchone()[0] + db.conn.execute( + "INSERT INTO command_queue (device_id, seq, cmd, timeout_ms, drain_first, status) VALUES (?,?,?,?,?,'pending')", + (device_id, seq, cmd, tmo, drain) + ) + db.conn.commit() + logger.info(f"RawELM: enqueued #{seq} → {cmd}") + return jsonify({"ok": True, "seq": seq, "cmd": cmd}) @bp.route("/api/v1/elm/raw/cmd", methods=["GET"]) def raw_dequeue_cmd(): - """Android: забрать команду из очереди. - - Returns: - 200 {"cmd": "0105", "seq": 1, ...} — есть команда - 204 — нет команды, полли дальше - """ - global _pending_cmd - device_id = request.args.get("device_id", "") - - with _lock: - if not _device_ready: - return jsonify({"error": "device not ready"}), 503 - if _pending_cmd is None: - return "", 204 # No Content — полли дальше - cmd = _pending_cmd - _pending_cmd = None # забрали - - logger.info(f"RawELM: dequeued #{cmd['seq']} → {cmd['cmd']} (device={device_id})") - return jsonify(cmd) + """Android: забрать команду из очереди.""" + device_id = request.args.get("device_id", "unknown") + with Database() as db: + row = db.conn.execute( + "SELECT id, seq, cmd, timeout_ms, drain_first FROM command_queue WHERE device_id=? AND status='pending' ORDER BY seq LIMIT 1", + (device_id,) + ).fetchone() + if not row: + return "", 204 + db.conn.execute( + "UPDATE command_queue SET status='sent', sent_at=datetime('now') WHERE id=?", + (row["id"],) + ) + db.conn.commit() + result = {"seq": row["seq"], "cmd": row["cmd"], "timeout_ms": row["timeout_ms"], "drain_first": bool(row["drain_first"])} + logger.info(f"RawELM: dequeued #{row['seq']} → {row['cmd']}") + return jsonify(result) @bp.route("/api/v1/elm/raw/response", methods=["POST"]) def raw_post_response(): - """Android: отправить ответ ELM327 на сервер. - - Body: { - "device_id": "android-xyz", - "seq": 1, - "cmd": "0105", - "raw": "41 05 5C", - "prompt": true, - "elapsed_ms": 48, - "bytes": 8, - "error": null - } - """ - global _last_response, _history + """Android: отправить ответ ELM327.""" data = request.get_json(silent=True) if not data: return jsonify({"error": "empty body"}), 400 + device_id = data.get("device_id", "unknown") + seq = data.get("seq", 0) + raw_resp = data.get("raw", "") + elapsed = data.get("elapsed_ms", 0) + prompt = 1 if data.get("prompt") else 0 + error = data.get("error") - with _lock: - _last_response = { - "seq": data.get("seq", 0), - "cmd": data.get("cmd", ""), - "raw": data.get("raw", ""), - "prompt": data.get("prompt", False), - "elapsed_ms": data.get("elapsed_ms", 0), - "bytes": data.get("bytes", 0), - "error": data.get("error"), - "received_at": time.time(), - } - _history.append(dict(_last_response)) - if len(_history) > 1000: - _history = _history[-500:] - - logger.info(f"RawELM: response #{_last_response['seq']} ← {_last_response['raw'][:80]}") + with Database() as db: + db.conn.execute( + "UPDATE command_queue SET status='done', raw_response=?, elapsed_ms=?, prompt=?, error=?, responded_at=datetime('now') WHERE device_id=? AND seq=?", + (raw_resp, elapsed, prompt, error, device_id, seq) + ) + db.conn.commit() + logger.info(f"RawELM: response #{seq} ← {raw_resp[:80]}") return jsonify({"ok": True}) @bp.route("/api/v1/elm/raw/response", methods=["GET"]) def raw_get_response(): - """Copilot: прочитать последний ответ от ELM327. - - Query: ?wait=30 — ждать до 30 сек пока появится новый ответ - """ - global _last_response, _pending_cmd - wait_s = int(request.args.get("wait", 0)) + """Copilot: прочитать последний ответ.""" seq = int(request.args.get("seq", 0)) + wait_s = int(request.args.get("wait", 0)) if wait_s > 0: - # Ждём пока появится ответ на команду с seq > указанного dl = time.time() + wait_s while time.time() < dl: - with _lock: - if _last_response and _last_response["seq"] > seq: - return jsonify(_last_response) - if _pending_cmd is None and _last_response: - # команд в очереди нет, ответ уже есть - return jsonify(_last_response) + with Database() as db: + row = db.conn.execute( + "SELECT seq, cmd, raw_response, elapsed_ms, prompt, error FROM command_queue WHERE seq > ? AND status='done' ORDER BY seq DESC LIMIT 1", + (seq,) + ).fetchone() + if row: + return jsonify({"seq": row["seq"], "cmd": row["cmd"], "raw": row["raw_response"], + "elapsed_ms": row["elapsed_ms"], "prompt": bool(row["prompt"]), "error": row["error"]}) time.sleep(0.5) - with _lock: - if _last_response is None: - return jsonify({"error": "no response yet", "seq": 0}) - return jsonify(_last_response) + with Database() as db: + row = db.conn.execute( + "SELECT seq, cmd, raw_response, elapsed_ms, prompt, error FROM command_queue WHERE status='done' ORDER BY seq DESC LIMIT 1" + ).fetchone() + if not row: + return jsonify({"error": "no response yet", "seq": 0}) + return jsonify({"seq": row["seq"], "cmd": row["cmd"], "raw": row["raw_response"], + "elapsed_ms": row["elapsed_ms"], "prompt": bool(row["prompt"]), "error": row["error"]}) @bp.route("/api/v1/elm/raw/status", methods=["GET"]) def raw_status(): - """Copilot: статус Android-устройства.""" - with _lock: - return jsonify({ - "device_ready": _device_ready, - "device_info": _device_info, - "pending_cmd": bool(_pending_cmd), - "pending_seq": _pending_seq, - "last_response_seq": _last_response["seq"] if _last_response else 0, - "history_count": len(_history), - }) + """Copilot: статус устройства.""" + with Database() as db: + last = db.conn.execute( + "SELECT seq, status FROM command_queue ORDER BY id DESC LIMIT 1" + ).fetchone() + pending = db.conn.execute( + "SELECT COUNT(*) FROM command_queue WHERE status='pending'" + ).fetchone()[0] + total = db.conn.execute( + "SELECT COUNT(*) FROM command_queue WHERE status='done'" + ).fetchone()[0] + return jsonify({ + "device_ready": last is not None, + "pending_cmd": pending > 0, + "pending_seq": pending, + "last_response_seq": last["seq"] if last else 0, + "history_count": total, + }) @bp.route("/api/v1/elm/raw/history", methods=["GET"]) def raw_history(): - """Copilot: история всех команд через Android.""" + """Copilot: история команд.""" n = int(request.args.get("n", 50)) - with _lock: - return jsonify({"history": _history[-n:], "total": len(_history)}) + with Database() as db: + rows = db.conn.execute( + "SELECT seq, cmd, raw_response, elapsed_ms, prompt, error FROM command_queue WHERE status='done' ORDER BY id DESC LIMIT ?", + (n,) + ).fetchall() + history = [{"seq": r["seq"], "cmd": r["cmd"], "raw": r["raw_response"], + "elapsed_ms": r["elapsed_ms"], "prompt": bool(r["prompt"]), "error": r["error"]} + for r in rows] + return jsonify({"history": history, "total": len(history)})