fix: raw_elm.py на SQLite очередь (gunicorn-safe), таблица command_queue в db.py
This commit is contained in:
@@ -160,6 +160,29 @@ class Database:
|
|||||||
|
|
||||||
self.conn.commit()
|
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 ──────────────────────────────────────────
|
# ── sessions ──────────────────────────────────────────
|
||||||
|
|
||||||
def get_cached_response(self, request_id: str) -> dict | None:
|
def get_cached_response(self, request_id: str) -> dict | None:
|
||||||
|
|||||||
+118
-140
@@ -23,23 +23,16 @@ import threading
|
|||||||
import time
|
import time
|
||||||
from flask import jsonify, request, Blueprint
|
from flask import jsonify, request, Blueprint
|
||||||
|
|
||||||
|
from api.db import Database
|
||||||
|
|
||||||
logger = logging.getLogger("elmer.raw_api")
|
logger = logging.getLogger("elmer.raw_api")
|
||||||
|
|
||||||
# ══════════════════════════════════════════════════════════
|
# ══════════════════════════════════════════════════════════
|
||||||
# Глобальное состояние
|
# Локальное состояние (только для прямого подключения ELM)
|
||||||
# ══════════════════════════════════════════════════════════
|
# ══════════════════════════════════════════════════════════
|
||||||
|
|
||||||
_raw_mode = False
|
_raw_mode = False
|
||||||
_raw_elm = None # локальный RawELM
|
_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
|
|
||||||
|
|
||||||
|
|
||||||
def is_raw_mode() -> bool:
|
def is_raw_mode() -> bool:
|
||||||
@@ -119,7 +112,14 @@ def raw_available():
|
|||||||
@bp.route("/api/v1/elm/raw/log", methods=["GET"])
|
@bp.route("/api/v1/elm/raw/log", methods=["GET"])
|
||||||
def raw_log():
|
def raw_log():
|
||||||
if not _raw_elm:
|
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)})
|
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 {}
|
data = request.get_json(silent=True) or {}
|
||||||
on = data.get("raw_mode", False)
|
on = data.get("raw_mode", False)
|
||||||
set_raw_mode(on)
|
set_raw_mode(on)
|
||||||
return jsonify({"raw_mode": _raw_mode, "has_local_elm": _raw_elm is not None,
|
with Database() as db:
|
||||||
"device_ready": _device_ready})
|
device = db.conn.execute(
|
||||||
return jsonify({"raw_mode": _raw_mode, "has_local_elm": _raw_elm is not None,
|
"SELECT device_id FROM command_queue ORDER BY id DESC LIMIT 1"
|
||||||
"device_ready": _device_ready})
|
).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"])
|
@bp.route("/api/v1/elm/raw/hello", methods=["POST"])
|
||||||
def raw_hello():
|
def raw_hello():
|
||||||
"""Android сообщает: «я подключился к ELM327, готов принимать команды».
|
"""Android: «я подключился, готов принимать команды»."""
|
||||||
|
|
||||||
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
|
|
||||||
data = request.get_json(silent=True) or {}
|
data = request.get_json(silent=True) or {}
|
||||||
with _lock:
|
device_id = data.get("device_id", "unknown")
|
||||||
_device_ready = True
|
logger.info(f"RawELM: device ready — {device_id} "
|
||||||
_device_info = {
|
f"({data.get('elm_version', '?')}, proto {data.get('protocol', '?')})")
|
||||||
"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']})")
|
|
||||||
return jsonify({"ok": True, "seq": 0})
|
return jsonify({"ok": True, "seq": 0})
|
||||||
|
|
||||||
|
|
||||||
@bp.route("/api/v1/elm/raw/cmd", methods=["POST"])
|
@bp.route("/api/v1/elm/raw/cmd", methods=["POST"])
|
||||||
def raw_enqueue_cmd():
|
def raw_enqueue_cmd():
|
||||||
"""Copilot: поставить команду в очередь для Android.
|
"""Copilot: поставить команду в очередь."""
|
||||||
|
|
||||||
Body: {
|
|
||||||
"cmd": "0105",
|
|
||||||
"timeout_ms": 500,
|
|
||||||
"drain_first": false
|
|
||||||
}
|
|
||||||
"""
|
|
||||||
global _pending_cmd, _pending_seq
|
|
||||||
data = request.get_json(silent=True)
|
data = request.get_json(silent=True)
|
||||||
if not data or "cmd" not in data:
|
if not data or "cmd" not in data:
|
||||||
return jsonify({"error": "missing 'cmd'"}), 400
|
return jsonify({"error": "missing 'cmd'"}), 400
|
||||||
|
|
||||||
cmd = data["cmd"].strip()
|
cmd = data["cmd"].strip()
|
||||||
if not cmd:
|
if not cmd:
|
||||||
return jsonify({"error": "empty cmd"}), 400
|
return jsonify({"error": "empty cmd"}), 400
|
||||||
|
|
||||||
with _lock:
|
device_id = data.get("device_id", "unknown")
|
||||||
_pending_seq += 1
|
tmo = data.get("timeout_ms", 500)
|
||||||
_pending_cmd = {
|
drain = 1 if data.get("drain_first") else 0
|
||||||
"cmd": cmd,
|
|
||||||
"timeout_ms": data.get("timeout_ms", 500),
|
with Database() as db:
|
||||||
"drain_first": data.get("drain_first", False),
|
cur = db.conn.execute("SELECT COALESCE(MAX(seq), 0) + 1 FROM command_queue WHERE device_id = ?", (device_id,))
|
||||||
"seq": _pending_seq,
|
seq = cur.fetchone()[0]
|
||||||
}
|
db.conn.execute(
|
||||||
logger.info(f"RawELM: enqueued #{_pending_seq} → {cmd}")
|
"INSERT INTO command_queue (device_id, seq, cmd, timeout_ms, drain_first, status) VALUES (?,?,?,?,?,'pending')",
|
||||||
return jsonify({"ok": True, "seq": _pending_seq, "cmd": cmd})
|
(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"])
|
@bp.route("/api/v1/elm/raw/cmd", methods=["GET"])
|
||||||
def raw_dequeue_cmd():
|
def raw_dequeue_cmd():
|
||||||
"""Android: забрать команду из очереди.
|
"""Android: забрать команду из очереди."""
|
||||||
|
device_id = request.args.get("device_id", "unknown")
|
||||||
Returns:
|
with Database() as db:
|
||||||
200 {"cmd": "0105", "seq": 1, ...} — есть команда
|
row = db.conn.execute(
|
||||||
204 — нет команды, полли дальше
|
"SELECT id, seq, cmd, timeout_ms, drain_first FROM command_queue WHERE device_id=? AND status='pending' ORDER BY seq LIMIT 1",
|
||||||
"""
|
(device_id,)
|
||||||
global _pending_cmd
|
).fetchone()
|
||||||
device_id = request.args.get("device_id", "")
|
if not row:
|
||||||
|
return "", 204
|
||||||
with _lock:
|
db.conn.execute(
|
||||||
if not _device_ready:
|
"UPDATE command_queue SET status='sent', sent_at=datetime('now') WHERE id=?",
|
||||||
return jsonify({"error": "device not ready"}), 503
|
(row["id"],)
|
||||||
if _pending_cmd is None:
|
)
|
||||||
return "", 204 # No Content — полли дальше
|
db.conn.commit()
|
||||||
cmd = _pending_cmd
|
result = {"seq": row["seq"], "cmd": row["cmd"], "timeout_ms": row["timeout_ms"], "drain_first": bool(row["drain_first"])}
|
||||||
_pending_cmd = None # забрали
|
logger.info(f"RawELM: dequeued #{row['seq']} → {row['cmd']}")
|
||||||
|
return jsonify(result)
|
||||||
logger.info(f"RawELM: dequeued #{cmd['seq']} → {cmd['cmd']} (device={device_id})")
|
|
||||||
return jsonify(cmd)
|
|
||||||
|
|
||||||
|
|
||||||
@bp.route("/api/v1/elm/raw/response", methods=["POST"])
|
@bp.route("/api/v1/elm/raw/response", methods=["POST"])
|
||||||
def raw_post_response():
|
def raw_post_response():
|
||||||
"""Android: отправить ответ ELM327 на сервер.
|
"""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
|
|
||||||
data = request.get_json(silent=True)
|
data = request.get_json(silent=True)
|
||||||
if not data:
|
if not data:
|
||||||
return jsonify({"error": "empty body"}), 400
|
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:
|
with Database() as db:
|
||||||
_last_response = {
|
db.conn.execute(
|
||||||
"seq": data.get("seq", 0),
|
"UPDATE command_queue SET status='done', raw_response=?, elapsed_ms=?, prompt=?, error=?, responded_at=datetime('now') WHERE device_id=? AND seq=?",
|
||||||
"cmd": data.get("cmd", ""),
|
(raw_resp, elapsed, prompt, error, device_id, seq)
|
||||||
"raw": data.get("raw", ""),
|
)
|
||||||
"prompt": data.get("prompt", False),
|
db.conn.commit()
|
||||||
"elapsed_ms": data.get("elapsed_ms", 0),
|
logger.info(f"RawELM: response #{seq} ← {raw_resp[:80]}")
|
||||||
"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]}")
|
|
||||||
return jsonify({"ok": True})
|
return jsonify({"ok": True})
|
||||||
|
|
||||||
|
|
||||||
@bp.route("/api/v1/elm/raw/response", methods=["GET"])
|
@bp.route("/api/v1/elm/raw/response", methods=["GET"])
|
||||||
def raw_get_response():
|
def raw_get_response():
|
||||||
"""Copilot: прочитать последний ответ от ELM327.
|
"""Copilot: прочитать последний ответ."""
|
||||||
|
|
||||||
Query: ?wait=30 — ждать до 30 сек пока появится новый ответ
|
|
||||||
"""
|
|
||||||
global _last_response, _pending_cmd
|
|
||||||
wait_s = int(request.args.get("wait", 0))
|
|
||||||
seq = int(request.args.get("seq", 0))
|
seq = int(request.args.get("seq", 0))
|
||||||
|
wait_s = int(request.args.get("wait", 0))
|
||||||
|
|
||||||
if wait_s > 0:
|
if wait_s > 0:
|
||||||
# Ждём пока появится ответ на команду с seq > указанного
|
|
||||||
dl = time.time() + wait_s
|
dl = time.time() + wait_s
|
||||||
while time.time() < dl:
|
while time.time() < dl:
|
||||||
with _lock:
|
with Database() as db:
|
||||||
if _last_response and _last_response["seq"] > seq:
|
row = db.conn.execute(
|
||||||
return jsonify(_last_response)
|
"SELECT seq, cmd, raw_response, elapsed_ms, prompt, error FROM command_queue WHERE seq > ? AND status='done' ORDER BY seq DESC LIMIT 1",
|
||||||
if _pending_cmd is None and _last_response:
|
(seq,)
|
||||||
# команд в очереди нет, ответ уже есть
|
).fetchone()
|
||||||
return jsonify(_last_response)
|
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)
|
time.sleep(0.5)
|
||||||
|
|
||||||
with _lock:
|
with Database() as db:
|
||||||
if _last_response is None:
|
row = db.conn.execute(
|
||||||
return jsonify({"error": "no response yet", "seq": 0})
|
"SELECT seq, cmd, raw_response, elapsed_ms, prompt, error FROM command_queue WHERE status='done' ORDER BY seq DESC LIMIT 1"
|
||||||
return jsonify(_last_response)
|
).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"])
|
@bp.route("/api/v1/elm/raw/status", methods=["GET"])
|
||||||
def raw_status():
|
def raw_status():
|
||||||
"""Copilot: статус Android-устройства."""
|
"""Copilot: статус устройства."""
|
||||||
with _lock:
|
with Database() as db:
|
||||||
return jsonify({
|
last = db.conn.execute(
|
||||||
"device_ready": _device_ready,
|
"SELECT seq, status FROM command_queue ORDER BY id DESC LIMIT 1"
|
||||||
"device_info": _device_info,
|
).fetchone()
|
||||||
"pending_cmd": bool(_pending_cmd),
|
pending = db.conn.execute(
|
||||||
"pending_seq": _pending_seq,
|
"SELECT COUNT(*) FROM command_queue WHERE status='pending'"
|
||||||
"last_response_seq": _last_response["seq"] if _last_response else 0,
|
).fetchone()[0]
|
||||||
"history_count": len(_history),
|
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"])
|
@bp.route("/api/v1/elm/raw/history", methods=["GET"])
|
||||||
def raw_history():
|
def raw_history():
|
||||||
"""Copilot: история всех команд через Android."""
|
"""Copilot: история команд."""
|
||||||
n = int(request.args.get("n", 50))
|
n = int(request.args.get("n", 50))
|
||||||
with _lock:
|
with Database() as db:
|
||||||
return jsonify({"history": _history[-n:], "total": len(_history)})
|
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)})
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user