server: raw-obd state machine + session API + LLM integration

This commit is contained in:
“Naeel”
2026-05-25 22:43:33 +04:00
parent 0240322e35
commit 0b1f2f1d59
+285 -22
View File
@@ -1,26 +1,185 @@
"""Эндпоинт /api/v1/raw-obd — приём сырых OBD-ответов от тонкого клиента. """Эндпоинт /api/v1/raw-obd — приём сырых ответов, парсинг, управление сессией.
Формат входящих данных: Сервер — мозг. Знает протокол ELM327. Клиент — тупая труба.
{"raw": "41 05 5A", "timestamp": 1716652800000}
Сервер логирует, складывает в очередь, и по готовности парсит + отправляет в LLM. Протокол:
Клиент → {"raw": "..."} → Сервер
Сервер → {"cmd": "ATZ"} → Клиент пишет в BT
... итерации ...
Сервер → {"cmd": null} → сессия завершена, диагноз готов
Сессия:
1. INIT: ATZ, ATE0, ATL0, ATSP0, ATH1
2. VIN: 0902
3. DTC: 03 (stored), 07 (pending)
4. PIDS: 0105, 010C, 010D, ...
5. LLM: отправка в DeepSeek
6. LLM_LOOP: если LLM хочет ещё данных → ещё команды
7. DONE
""" """
import json
import logging import logging
import re
import threading import threading
from collections import deque
from datetime import datetime, timezone from datetime import datetime, timezone
logger = logging.getLogger("elmer.raw") logger = logging.getLogger("elmer.raw")
# Буфер сырых сообщений на сессию (в production — Redis очередь) # ── Сессии ─────────────────────────────────────────────
raw_buffer: deque[dict] = deque()
buffer_lock = threading.Lock()
# session_id → {state, commands[], responses[], data{}}
sessions: dict[str, dict] = {}
sessions_lock = threading.Lock()
def _new_session() -> str:
"""Создаёт новую сессию, возвращает ID."""
import uuid
sid = uuid.uuid4().hex[:12]
# Очередь команд для инициализации
init_cmds = ["ATZ", "ATE0", "ATL0", "ATSP0", "ATH1"]
with sessions_lock:
sessions[sid] = {
"state": "INIT",
"created": datetime.now(timezone.utc).isoformat(),
"queue": list(init_cmds), # команды для отправки
"responses": [], # сырые ответы
"data": { # распарсенные данные
"vin": None,
"dtc_stored": [],
"dtc_pending": [],
"pids": [],
},
"llm_history": [],
"diagnosis": None,
}
return sid
def _parse_response(raw: str, data: dict) -> str | None:
"""Парсит сырой ответ ELM327, заполняет data. Возвращает None или описание."""
raw = raw.strip().upper()
# OK / ?>
if raw in ("OK", "?", "NO DATA", "SEARCHING...", "STOPPED"):
return raw
# ELM327 version
if raw.startswith("ELM"):
data["elm_version"] = raw
return f"ELM: {raw}"
# VIN response (mode 09 PID 02): "014 0:49 02 01 57 56 57..."
if "49 02" in raw or "49:02" in raw:
hex_str = re.sub(r".*49.?02.?01", "", raw.replace("\n", " ").replace(":", " ")).strip()
hex_bytes = hex_str.split()
vin = ""
for h in hex_bytes:
try:
vin += chr(int(h, 16))
except (ValueError, OverflowError):
pass
if len(vin) == 17:
data["vin"] = vin
return f"VIN: {vin}"
return f"VIN partial: {vin}"
# DTC response (mode 43/47)
if raw.startswith("43") or raw.startswith("47"):
mode = "stored" if raw.startswith("43") else "pending"
# Парсим коды: 43 01 33 00... → P0301
hex_bytes = re.sub(r"^4[37]\s*", "", raw).split()
i = 1 # пропускаем байт количества
codes = []
while i + 1 < len(hex_bytes):
a, b = int(hex_bytes[i], 16), int(hex_bytes[i + 1], 16)
prefix = {0: "P", 1: "C", 2: "B", 3: "U"}.get(a >> 6, "?")
code = f"{prefix}{(a>>4)&3}{a&15}{b>>4:X}{b&15:X}"
if code != "P0000":
codes.append(code)
if mode == "stored":
data["dtc_stored"].append(code)
else:
data["dtc_pending"].append(code)
i += 2
return f"DTC {mode}: {codes}"
# PID response (mode 41): "41 05 5A"
if raw.startswith("41"):
parts = raw.split()
if len(parts) >= 3:
pid = parts[1]
hex_vals = parts[2:]
formulas = {
"05": lambda b: int(b, 16) - 40,
"0C": lambda b: (int(b[0], 16) * 256 + int(b[1], 16)) / 4,
"0D": lambda b: int(b[0], 16),
"11": lambda b: int(b[1], 16) * 100 / 255 if len(b) > 1 else int(b[0], 16) * 100 / 255,
"0B": lambda b: int(b[0], 16),
"0F": lambda b: int(b[0], 16) - 40,
"1F": lambda b: int(b[0], 16) * 256 + (int(b[1], 16) if len(b) > 1 else 0),
"04": lambda b: int(b[0], 16) * 100 / 255,
"06": lambda b: (int(b[0], 16) - 128) * 100 / 128,
"07": lambda b: (int(b[0], 16) - 128) * 100 / 128,
}
if pid in formulas:
try:
val = round(formulas[pid](hex_vals), 1)
data["pids"].append({"pid": pid, "value": val})
return f"PID {pid}: {val}"
except Exception:
pass
return f"PID {pid} raw: {hex_vals}"
# AT-ответы (протокол)
if raw.startswith("AUTO") or "ISO" in raw or "SAE" in raw:
data["protocol"] = raw
return f"Protocol: {raw}"
return raw
def _next_commands(sess: dict) -> list[str] | None:
"""Определяет следующие команды в зависимости от состояния сессии."""
state = sess["state"]
data = sess["data"]
if state == "INIT":
# Инициализация завершена → запрос VIN
sess["state"] = "VIN"
return ["0902"]
if state == "VIN":
sess["state"] = "DTC_STORED"
return ["03"]
if state == "DTC_STORED":
sess["state"] = "DTC_PENDING"
return ["07"]
if state == "DTC_PENDING":
sess["state"] = "PIDS"
# Стандартные PID
return ["0105", "010C", "010D", "0111", "010B", "010F", "011F", "0104", "0106", "0107"]
if state == "PIDS":
# Все данные собраны → LLM
sess["state"] = "LLM"
return None # Нет команд, вызываем LLM
if state == "LLM":
sess["state"] = "DONE"
return None
return None
# ── Flask endpoint ─────────────────────────────────────
def register(app): def register(app):
"""Регистрирует эндпоинт на Flask-приложении.""" """Регистрирует /api/v1/raw-obd на Flask-приложении."""
@app.route("/api/v1/raw-obd", methods=["POST"]) @app.route("/api/v1/raw-obd", methods=["POST"])
def raw_obd(): def raw_obd():
@@ -28,25 +187,129 @@ def register(app):
data = request.get_json(silent=True) data = request.get_json(silent=True)
if not data or "raw" not in data: if not data or "raw" not in data:
return jsonify({"error": "missing 'raw' field"}), 400 return jsonify({"error": "missing 'raw'"}), 400
raw = data["raw"].strip() raw = data["raw"].strip()
ts = data.get("timestamp", int(datetime.now(timezone.utc).timestamp() * 1000)) session_id = data.get("session")
# Логируем в консоль # Новая сессия?
logger.info(f"[RAW] {raw}") if not session_id:
session_id = _new_session()
logger.info(f"[{session_id}] NEW SESSION")
# Складываем в буфер with sessions_lock:
with buffer_lock: sess = sessions.get(session_id)
raw_buffer.append({ if not sess:
"raw": raw, session_id = _new_session()
"timestamp": ts, sess = sessions[session_id]
sess["responses"].append(raw)
# Парсим ответ
parsed = _parse_response(raw, sess["data"])
logger.info(f"[{session_id}] {parsed}")
# Если есть очередь команд — отдаём следующую
if sess["queue"]:
cmd = sess["queue"].pop(0)
return jsonify({"cmd": cmd, "session": session_id, "state": sess["state"]})
# Определяем что дальше
next_cmds = _next_commands(sess)
if next_cmds:
sess["queue"] = list(next_cmds)
cmd = sess["queue"].pop(0)
return jsonify({"cmd": cmd, "session": session_id, "state": sess["state"]})
# LLM фаза
if sess["state"] == "LLM":
sess["state"] = "LLM_WAIT"
# Запускаем LLM асинхронно (в потоке)
import threading as th
th.Thread(target=_call_llm, args=(session_id,), daemon=True).start()
return jsonify({
"cmd": None,
"session": session_id,
"state": "LLM",
"msg": "Анализирую..."
}) })
# Готово
return jsonify({ return jsonify({
"status": "ok", "cmd": None,
"received": raw, "session": session_id,
"buffer_size": len(raw_buffer), "state": sess["state"],
"diagnosis": sess.get("diagnosis"),
})
def _call_llm(session_id: str):
"""Вызывает DeepSeek с собранными данными."""
from elmer.config import load
from elmer.diagnose import Diagnoser
from elmer.prompts import SYSTEM_PROMPT
with sessions_lock:
sess = sessions.get(session_id)
if not sess:
return
data = sess["data"]
config = load()
api_key = config["deepseek"]["api_key"]
if not api_key:
logger.error(f"[{session_id}] No DEEPSEEK_API_KEY")
return
# Собираем промпт
parts = [f"VIN: {data.get('vin', 'неизвестен')}"]
if data.get("dtc_stored"):
parts.append(f"Ошибки (stored): {', '.join(data['dtc_stored'])}")
if data.get("dtc_pending"):
parts.append(f"Ошибки (pending): {', '.join(data['dtc_pending'])}")
if data.get("pids"):
parts.append("Параметры: " + ", ".join(
f"{p['pid']}={p['value']}" for p in data["pids"]))
user_prompt = "\n".join(parts)
try:
diagnoser = Diagnoser(api_key=api_key, model=config["deepseek"].get("model", "deepseek-chat"))
answer = diagnoser.diagnose(SYSTEM_PROMPT, user_prompt)
with sessions_lock:
if sess:
sess["diagnosis"] = answer
sess["state"] = "DONE"
logger.info(f"[{session_id}] Diagnosis ready ({len(answer)} chars)")
except Exception as e:
logger.error(f"[{session_id}] LLM error: {e}")
with sessions_lock:
if sess:
sess["diagnosis"] = f"Ошибка: {e}"
sess["state"] = "DONE"
@app.route("/api/v1/session/<session_id>", methods=["GET"])
def get_session(session_id):
"""Получить статус и диагноз сессии."""
from flask import jsonify
with sessions_lock:
sess = sessions.get(session_id)
if not sess:
return jsonify({"error": "session not found"}), 404
return jsonify({
"session": session_id,
"state": sess["state"],
"diagnosis": sess.get("diagnosis"),
"data": {
"vin": sess["data"].get("vin"),
"dtc_stored": sess["data"].get("dtc_stored", []),
"dtc_pending": sess["data"].get("dtc_pending", []),
"pids": sess["data"].get("pids", []),
},
"created": sess.get("created"),
}) })
return app return app