From 0b1f2f1d59515e0ceb84693e87d31e4c2b141238 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E2=80=9CNaeel=E2=80=9D?= Date: Mon, 25 May 2026 22:43:33 +0400 Subject: [PATCH] server: raw-obd state machine + session API + LLM integration --- web/raw_endpoint.py | 307 ++++++++++++++++++++++++++++++++++++++++---- 1 file changed, 285 insertions(+), 22 deletions(-) diff --git a/web/raw_endpoint.py b/web/raw_endpoint.py index 3c05e55..0243e35 100644 --- a/web/raw_endpoint.py +++ b/web/raw_endpoint.py @@ -1,26 +1,185 @@ -"""Эндпоинт /api/v1/raw-obd — приём сырых OBD-ответов от тонкого клиента. +"""Эндпоинт /api/v1/raw-obd — приём сырых ответов, парсинг, управление сессией. -Формат входящих данных: - {"raw": "41 05 5A", "timestamp": 1716652800000} +Сервер — мозг. Знает протокол ELM327. Клиент — тупая труба. -Сервер логирует, складывает в очередь, и по готовности парсит + отправляет в 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 re import threading -from collections import deque from datetime import datetime, timezone 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): - """Регистрирует эндпоинт на Flask-приложении.""" + """Регистрирует /api/v1/raw-obd на Flask-приложении.""" @app.route("/api/v1/raw-obd", methods=["POST"]) def raw_obd(): @@ -28,25 +187,129 @@ def register(app): data = request.get_json(silent=True) 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() - 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 buffer_lock: - raw_buffer.append({ - "raw": raw, - "timestamp": ts, + 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}") + + # Если есть очередь команд — отдаём следующую + 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({ - "status": "ok", - "received": raw, - "buffer_size": len(raw_buffer), + "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/", 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