Files
elmer/web/raw_endpoint.py
T

351 lines
13 KiB
Python

"""Эндпоинт /api/v1/raw-obd — приём сырых ответов, парсинг, управление сессией.
Сервер — мозг. Знает протокол ELM327. Клиент — тупая труба.
Протокол:
Клиент → {"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 logging
import re
import threading
from datetime import datetime, timezone
logger = logging.getLogger("elmer.raw")
# ── Сессии ─────────────────────────────────────────────
# 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", "READY"):
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):
"""Регистрирует /api/v1/raw-obd на Flask-приложении."""
@app.route("/api/v1/raw-obd", methods=["POST"])
def raw_obd():
from flask import request, jsonify
data = request.get_json(silent=True)
if not data or "raw" not in data:
return jsonify({"error": "missing 'raw'"}), 400
raw = data["raw"].strip()
session_id = data.get("session")
# Новая сессия?
if not session_id:
session_id = _new_session()
logger.info(f"[{session_id}] NEW SESSION")
with sessions_lock:
sess = sessions.get(session_id)
if not sess:
session_id = _new_session()
sess = sessions[session_id]
sess["responses"].append(raw)
# Парсим ответ
parsed = _parse_response(raw, sess["data"])
logger.info(f"[{session_id}] {parsed}")
# SEARCHING/NO DATA — ждём, не продвигаем стейт
is_skip = raw.startswith("SEARCHING") or raw in ("NO DATA", "STOPPED", "?")
if is_skip:
sess["retries"] = sess.get("retries", 0) + 1
if sess["retries"] > 3:
logger.info(f"[{session_id}] Giving up after {sess['retries']} retries")
sess["retries"] = 0
# Продвигаем принудительно (ниже)
else:
# Ждём — не шлём команду, ELM327 сам ответит когда готов
return jsonify({"cmd": None, "session": session_id, "state": sess["state"], "msg": "Жду..."})
# Если есть очередь — отдаём следующую
if sess["queue"]:
cmd = sess["queue"].pop(0)
return jsonify({"cmd": cmd, "session": session_id, "state": sess["state"]})
sess["retries"] = 0
# Определяем что дальше
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":
from elmer.config import load
cfg = load()
api_key = cfg["deepseek"]["api_key"]
if not api_key:
# Без LLM — форматируем декодированные данные
data = sess["data"]
lines = []
if data.get("vin"): lines.append(f"VIN: {data['vin']}")
if data.get("dtc_stored"): lines.append(f"Ошибки: {', '.join(data['dtc_stored'])}")
if data.get("dtc_pending"): lines.append(f"Pending: {', '.join(data['dtc_pending'])}")
if data.get("pids"):
names = {"05": "ОЖ", "0C": "RPM", "0D": "Скорость", "11": "Дроссель",
"0B": "MAP", "0F": "IAT", "1F": "Время", "04": "Нагрузка",
"06": "STFT", "07": "LTFT"}
for p in data["pids"]:
n = names.get(p["pid"], p["pid"])
lines.append(f"{n}: {p['value']}")
sess["diagnosis"] = "\n".join(lines) if lines else "Нет данных"
sess["state"] = "DONE"
logger.info(f"[{session_id}] No LLM — returning raw data")
else:
sess["state"] = "LLM_WAIT"
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({
"cmd": None,
"session": session_id,
"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