"""Эндпоинт /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/", 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