351 lines
13 KiB
Python
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
|