diff --git a/api/__init__.py b/api/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/api/config.py b/api/config.py new file mode 100644 index 0000000..fa27765 --- /dev/null +++ b/api/config.py @@ -0,0 +1,39 @@ +"""Загрузка конфигурации из config.yaml.""" + +import os +from pathlib import Path + +import yaml + +CONFIG_PATH = Path(os.environ.get("ELMER_CONFIG", Path(__file__).parent.parent / "config.yaml")) + + +def load() -> dict: + """Читает config.yaml, подставляет переменные окружения в значения.""" + if not CONFIG_PATH.exists(): + raise FileNotFoundError(f"Конфиг не найден: {CONFIG_PATH}") + + with open(CONFIG_PATH) as f: + config = yaml.safe_load(f) + + # Подстановка ${VAR} из переменных окружения + _resolve_env(config) + return config + + +def _resolve_env(obj): + """Рекурсивно заменяет ${VAR} на os.environ['VAR'].""" + if isinstance(obj, dict): + for k, v in obj.items(): + if isinstance(v, str) and v.startswith("${") and v.endswith("}"): + env_var = v[2:-1] + obj[k] = os.environ.get(env_var, "") + else: + _resolve_env(v) + elif isinstance(obj, list): + for i, v in enumerate(obj): + if isinstance(v, str) and v.startswith("${") and v.endswith("}"): + env_var = v[2:-1] + obj[i] = os.environ.get(env_var, "") + else: + _resolve_env(v) diff --git a/api/db.py b/api/db.py new file mode 100644 index 0000000..de1e4cd --- /dev/null +++ b/api/db.py @@ -0,0 +1,275 @@ +"""SQLite — сохранение сессий диагностики. + +Схема: + cars — VIN, марка, модель, год, двигатель + diagnostic_tokens — id (PK), car_id (FK), created_at + llm_messages — token_id (FK), role, content, timestamp + ecu_parameters — token_id (FK), pid_code, value, unit, timestamp + dtc_codes — token_id (FK), code, description, status + sessions — сводная таблица всех сессий (клиент, ELM, авто, LLM) +""" + +import json +import sqlite3 +from datetime import datetime, timezone +from pathlib import Path + + +class Database: + def __init__(self, path: str | Path = "elmer.db"): + self.path = Path(path) + self.conn = sqlite3.connect(str(self.path)) + self.conn.row_factory = sqlite3.Row + self._init_schema() + + def _init_schema(self): + self.conn.executescript(""" + CREATE TABLE IF NOT EXISTS cars ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + vin TEXT NOT NULL UNIQUE, + make TEXT, + model TEXT, + year INTEGER, + engine TEXT, + created_at TEXT NOT NULL DEFAULT (datetime('now')) + ); + + CREATE TABLE IF NOT EXISTS diagnostic_tokens ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + car_id INTEGER NOT NULL REFERENCES cars(id), + created_at TEXT NOT NULL DEFAULT (datetime('now')) + ); + + CREATE TABLE IF NOT EXISTS llm_messages ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + token_id INTEGER NOT NULL REFERENCES diagnostic_tokens(id), + role TEXT NOT NULL, -- 'system' | 'user' | 'assistant' + content TEXT NOT NULL, + created_at TEXT NOT NULL DEFAULT (datetime('now')) + ); + + CREATE TABLE IF NOT EXISTS ecu_parameters ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + token_id INTEGER NOT NULL REFERENCES diagnostic_tokens(id), + pid_code TEXT NOT NULL, -- напр. '0105', '010C' + name TEXT, -- напр. 'coolant_temp', 'rpm' + value REAL, + unit TEXT, + created_at TEXT NOT NULL DEFAULT (datetime('now')) + ); + + CREATE TABLE IF NOT EXISTS dtc_codes ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + token_id INTEGER NOT NULL REFERENCES diagnostic_tokens(id), + code TEXT NOT NULL, -- напр. 'P0301' + description TEXT, + status TEXT, -- 'stored' | 'pending' + created_at TEXT NOT NULL DEFAULT (datetime('now')) + ); + + CREATE TABLE IF NOT EXISTS sessions ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + + -- Сервер + client_ip TEXT, + real_ip TEXT, + user_agent TEXT, + content_length INTEGER, + created_at TEXT NOT NULL DEFAULT (datetime('now')), + + -- Телефон + phone_model TEXT, + phone_maker TEXT, + android_version TEXT, + android_sdk INTEGER, + app_version TEXT, + android_id TEXT, + + -- ELM327 + elm_mac TEXT, + elm_bt_name TEXT, + obd_protocol TEXT, + + -- Авто + vin TEXT, + dtc_count INTEGER DEFAULT 0, + pid_count INTEGER DEFAULT 0, + + -- Сессия + duration_ms INTEGER, + response_count INTEGER DEFAULT 0, + error_count INTEGER DEFAULT 0, + retry_count INTEGER DEFAULT 0, + timeout_count INTEGER DEFAULT 0, + script_mode TEXT, + transport TEXT, -- 'bt' | 'tcp' + mock_mode INTEGER DEFAULT 0, + + -- LLM + diagnosis_text TEXT, + diagnosis_len INTEGER, + llm_model TEXT, + llm_duration_ms INTEGER, + llm_success INTEGER DEFAULT 0, + + -- Сырые данные (JSON) + raw_responses TEXT + ); + + CREATE INDEX IF NOT EXISTS idx_sessions_created ON sessions(created_at); + CREATE INDEX IF NOT EXISTS idx_sessions_vin ON sessions(vin); + CREATE INDEX IF NOT EXISTS idx_sessions_mac ON sessions(elm_mac); + CREATE INDEX IF NOT EXISTS idx_sessions_aid ON sessions(android_id); + """) + self.conn.commit() + + # ── sessions ────────────────────────────────────────── + + def save_session(self, client_info: dict, responses: list[dict], + diagnosis: str = "", llm_model: str = "", + llm_duration_ms: int = 0, llm_success: bool = False): + """Сохраняет сводную запись о сессии.""" + ci = client_info + + # Подсчёт DTC/PID из ответов + dtc_count = 0 + pid_count = 0 + for r in responses: + dec = (r.get("decoded") or "").lower() + if dec.startswith("dtc"): + dtc_count += 1 + elif ":" in dec and not dec.startswith(("vin", "dtc", "elm", "protocol")): + pid_count += 1 + + # VIN из ответов + vin = None + for r in responses: + dec = (r.get("decoded") or "") + if dec.startswith("VIN:"): + vin = dec[4:].strip() + if len(vin) != 17: + vin = None + break + + self.conn.execute(""" + INSERT INTO sessions ( + client_ip, real_ip, user_agent, content_length, + phone_model, phone_maker, android_version, android_sdk, + app_version, android_id, + elm_mac, elm_bt_name, obd_protocol, + vin, dtc_count, pid_count, + duration_ms, response_count, error_count, + retry_count, timeout_count, script_mode, + transport, mock_mode, + diagnosis_text, diagnosis_len, llm_model, + llm_duration_ms, llm_success, + raw_responses + ) VALUES (?,?,?,?, ?,?,?,?, ?,?, ?,?,?, ?,?,?, ?,?,?, ?,?,?, ?,?, + ?,?,?, ?,?, ?) + """, ( + ci.get("client_ip"), ci.get("real_ip"), ci.get("user_agent"), + ci.get("content_length"), + ci.get("phone_model"), ci.get("phone_maker"), ci.get("android_version"), + ci.get("android_sdk"), ci.get("app_version"), ci.get("android_id"), + ci.get("elm_mac"), ci.get("elm_bt_name"), ci.get("obd_protocol"), + vin, dtc_count, pid_count, + ci.get("duration_ms"), len(responses), ci.get("error_count", 0), + ci.get("retry_count", 0), ci.get("timeout_count", 0), + ci.get("script_mode"), ci.get("transport"), ci.get("mock_mode", 0), + diagnosis, len(diagnosis), llm_model, + llm_duration_ms, 1 if llm_success else 0, + json.dumps(responses, ensure_ascii=False) if responses else None, + )) + self.conn.commit() + + def get_recent_sessions(self, limit: int = 50) -> list[dict]: + """Последние N сессий.""" + rows = self.conn.execute( + "SELECT * FROM sessions ORDER BY created_at DESC LIMIT ?", (limit,) + ).fetchall() + return [dict(r) for r in rows] + + # ── cars ────────────────────────────────────────────── + + def get_or_create_car(self, vin: str) -> int: + """Возвращает car_id по VIN, создаёт запись если нет.""" + row = self.conn.execute("SELECT id FROM cars WHERE vin = ?", (vin,)).fetchone() + if row: + return row["id"] + cur = self.conn.execute("INSERT INTO cars (vin) VALUES (?)", (vin,)) + self.conn.commit() + return cur.lastrowid + + def update_car_info(self, car_id: int, make: str, model: str, year: int, engine: str): + self.conn.execute( + "UPDATE cars SET make=?, model=?, year=?, engine=? WHERE id=?", + (make, model, year, engine, car_id), + ) + self.conn.commit() + + # ── tokens ──────────────────────────────────────────── + + def create_token(self, car_id: int) -> int: + """Создаёт новую сессию диагностики, возвращает token_id.""" + cur = self.conn.execute( + "INSERT INTO diagnostic_tokens (car_id) VALUES (?)", (car_id,) + ) + self.conn.commit() + return cur.lastrowid + + def last_token_for_car(self, car_id: int) -> int | None: + """Последняя сессия для VIN (для продолжения диалога), или None.""" + row = self.conn.execute( + "SELECT id FROM diagnostic_tokens WHERE car_id=? ORDER BY created_at DESC LIMIT 1", + (car_id,), + ).fetchone() + return row["id"] if row else None + + # ── llm_messages ────────────────────────────────────── + + def add_llm_message(self, token_id: int, role: str, content: str): + self.conn.execute( + "INSERT INTO llm_messages (token_id, role, content) VALUES (?, ?, ?)", + (token_id, role, content), + ) + self.conn.commit() + + def get_llm_messages(self, token_id: int) -> list[dict]: + """Возвращает историю диалога для токена.""" + rows = self.conn.execute( + "SELECT role, content FROM llm_messages WHERE token_id=? ORDER BY id", + (token_id,), + ).fetchall() + return [{"role": r["role"], "content": r["content"]} for r in rows] + + # ── ecu_parameters ──────────────────────────────────── + + def add_parameter(self, token_id: int, pid_code: str, name: str, value: float, unit: str): + self.conn.execute( + "INSERT INTO ecu_parameters (token_id, pid_code, name, value, unit) VALUES (?, ?, ?, ?, ?)", + (token_id, pid_code, name, value, unit), + ) + self.conn.commit() + + def get_parameters(self, token_id: int) -> list[dict]: + rows = self.conn.execute( + "SELECT pid_code, name, value, unit FROM ecu_parameters WHERE token_id=? ORDER BY id", + (token_id,), + ).fetchall() + return [dict(r) for r in rows] + + # ── dtc_codes ───────────────────────────────────────── + + def add_dtc(self, token_id: int, code: str, description: str = "", status: str = "stored"): + self.conn.execute( + "INSERT INTO dtc_codes (token_id, code, description, status) VALUES (?, ?, ?, ?)", + (token_id, code, description, status), + ) + self.conn.commit() + + def get_dtcs(self, token_id: int) -> list[dict]: + rows = self.conn.execute( + "SELECT code, description, status FROM dtc_codes WHERE token_id=? ORDER BY id", + (token_id,), + ).fetchall() + return [dict(r) for r in rows] diff --git a/api/parser.py b/api/parser.py new file mode 100644 index 0000000..f92ec09 --- /dev/null +++ b/api/parser.py @@ -0,0 +1,130 @@ +"""Парсинг батча ответов ELM327 в структуру для LLM. + +Поддерживает: + - VIN (mode 09 PID 02) — из decoded и fallback из raw HEX + - DTC stored/pending (mode 03/07) — из decoded и fallback из raw HEX + - PID параметры (mode 01) — из decoded +""" + +import logging + +logger = logging.getLogger("elmer.parser") + + +def parse_batch(responses: list[dict]) -> dict: + """Парсит батч ответов в структуру для LLM. + + Returns: + {"vin": str|None, "dtc_stored": [str], "dtc_pending": [str], + "parameters": [{"name": str, "value": str}], "raw_log": [str]} + """ + result = { + "vin": None, + "dtc_stored": [], + "dtc_pending": [], + "parameters": [], + "raw_log": [], + } + + for r in responses: + cmd = (r.get("cmd") or "").strip() + raw = (r.get("raw") or "").strip() + decoded = (r.get("decoded") or "").strip() + + result["raw_log"].append(f"→ {cmd}\n← {raw}") + + _parse_vin(result, raw, decoded) + _parse_dtc(result, raw, decoded, mode="03", key="dtc_stored", prefix="DTC stored:") + _parse_dtc(result, raw, decoded, mode="47", key="dtc_pending", prefix="DTC pending:") + _parse_pid(result, decoded, cmd) + + logger.info(f"[{cmd}] decoded={decoded[:60]}") + + return result + + +def _parse_vin(result: dict, raw: str, decoded: str): + if decoded.startswith("VIN:"): + vin = decoded.replace("VIN:", "").strip() + if len(vin) == 17: + result["vin"] = vin + return + + # Fallback: парсим VIN из raw HEX + if "49" in raw and ("02" in raw or "4902" in raw.replace(" ", "")): + clean = raw.replace(":", "").replace(" ", "").upper() + if "490201" in clean: + hex_str = clean.split("490201")[-1][:34] + vin = "" + for i in range(0, len(hex_str) - 1, 2): + try: + vin += chr(int(hex_str[i:i+2], 16)) + except (ValueError, OverflowError): + pass + if len(vin) == 17: + result["vin"] = vin + + +def _parse_dtc(result: dict, raw: str, decoded: str, *, mode: str, key: str, prefix: str): + if decoded.startswith(prefix): + codes = decoded.replace(prefix, "").strip() + if codes != "none": + result[key] = [c.strip() for c in codes.split()] + return + + # Fallback: парсим DTC из raw HEX (43XX... или 47XX...) + clean = raw.replace(" ", "").upper() + if clean.startswith(mode) and len(clean) >= 4: + codes = _decode_dtc_bytes(clean[2:]) + if codes: + result[key] = codes + + +def _decode_dtc_bytes(hex_str: str) -> list[str]: + """Декодирует HEX-строку DTC (после 43/47) в коды.""" + codes = [] + i = 2 # skip byte count + while i + 3 < len(hex_str): + try: + a = int(hex_str[i:i+2], 16) + b = int(hex_str[i+2:i+4], 16) + p = {0: "P", 1: "C", 2: "B", 3: "U"}.get(a >> 6, "?") + code = f"{p}{(a>>4)&3}{a&15}{b>>4:X}{b&15:X}" + if code != "P0000": + codes.append(code) + except Exception: + pass + i += 4 + return codes + + +def _parse_pid(result: dict, decoded: str, cmd: str): + if ":" not in decoded: + return + if decoded.startswith(("VIN", "DTC", "ELM", "Protocol")): + return + parts = decoded.split(":", 1) + if len(parts) == 2: + result["parameters"].append({ + "name": parts[0].strip(), + "value": parts[1].strip(), + }) + + +def format_no_llm(parsed: dict) -> str: + """Форматирует ответ без LLM.""" + lines = [] + if parsed["vin"]: + lines.append(f"VIN: {parsed['vin']}") + if parsed["dtc_stored"]: + lines.append(f"Ошибки: {', '.join(parsed['dtc_stored'])}") + if parsed["dtc_pending"]: + lines.append(f"Pending: {', '.join(parsed['dtc_pending'])}") + if parsed["parameters"]: + lines.append("Параметры:") + for p in parsed["parameters"]: + lines.append(f" {p['name']}: {p['value']}") + if not lines: + lines.append("Данные не распознаны.") + lines.append("\n(LLM не настроен — только сырые данные)") + return "\n".join(lines) diff --git a/api/routes.py b/api/routes.py new file mode 100644 index 0000000..b856d4c --- /dev/null +++ b/api/routes.py @@ -0,0 +1,224 @@ +"""Эндпоинты для толстого клиента: скрипты и батчевая загрузка. + +GET /api/v1/script — выдача скрипта диагностики +POST /api/v1/session/upload — приём батча, LLM-анализ, возврат диагноза +""" + +import logging +import time +from api.scripts import build_default_script, build_full_script +from api.parser import parse_batch, format_no_llm + +logger = logging.getLogger("elmer.script") + + +def _build_diagnosis_prompt(data: dict) -> str: + """Строит промпт для LLM из распарсенных данных.""" + parts = ["## Данные диагностики\n"] + + if data["vin"]: + parts.append(f"**VIN:** {data['vin']}") + + if data["dtc_stored"]: + parts.append(f"\n**Сохранённые ошибки (mode 03):** {', '.join(data['dtc_stored'])}") + if data["dtc_pending"]: + parts.append(f"**Ожидающие ошибки (mode 07):** {', '.join(data['dtc_pending'])}") + + if data["parameters"]: + parts.append("\n**Параметры в реальном времени:**") + for p in data["parameters"]: + parts.append(f"- {p['name']}: {p['value']}") + + if not data["vin"] and not data["dtc_stored"] and not data["parameters"]: + parts.append("\n(данные не распознаны)") + + parts.append("\n**Сырые ответы ЭБУ:**") + parts.extend(data["raw_log"]) + + parts.append("\n---") + parts.append("## Запрос на анализ") + parts.append( + "Дай ГЛУБОКИЙ, РАЗВЁРНУТЫЙ анализ на основе этих данных. " + "Не ограничивайся кратким резюме — мне нужен полный технический разбор.\n" + "1. Разбери КАЖДУЮ ошибку: все возможные причины, от частых к редким.\n" + "2. Проанализируй КАЖДЫЙ параметр: норма/отклонение, на что влияет.\n" + "3. Найди ВЗАИМОСВЯЗИ между ошибками и параметрами.\n" + "4. Предложи КОНКРЕТНЫЙ план проверок: что сначала, что потом, как проверять.\n" + "5. Если данных мало — скажи какие PID'ы досчитать и зачем.\n" + "6. Дай оценки уверенности в ПРОЦЕНТАХ для каждой версии." + ) + return "\n".join(parts) + + +def register(app): + """Регистрирует эндпоинты скриптов на Flask-приложении.""" + + @app.route("/api/v1/script", methods=["GET"]) + def get_script(): + from flask import jsonify, request + mode = request.args.get("mode", "full") + script = build_full_script() if mode == "full" else build_default_script() + return jsonify(script) + + @app.route("/api/v1/session/upload", methods=["POST"]) + def upload_session(): + from flask import request, jsonify + from api.config import load + from api.db import Database + from brain.client import Diagnoser + from brain.prompts import SYSTEM_PROMPT + + data = request.get_json(silent=True) + if not data or "responses" not in data: + return jsonify({"error": "missing 'responses'"}), 400 + + responses = data["responses"] + logger.info(f"Upload: {len(responses)} responses") + + # ── Информация о клиенте ────────────────────── + client_info = data.get("client_info", {}) + client_info["client_ip"] = request.remote_addr + client_info["real_ip"] = request.headers.get("X-Real-IP", "") + client_info["user_agent"] = request.headers.get("User-Agent", "") + client_info["content_length"] = request.content_length + + parsed = parse_batch(responses) + + cfg = load() + api_key = cfg["llm"]["api_key"] + model = cfg["llm"].get("model", "gpt-oss-120b") + llm_available = bool(api_key) + + llm_start = time.time() + llm_success = False + diagnosis = "" + + if not api_key: + diagnosis = format_no_llm(parsed) + else: + diagnoser = Diagnoser( + api_key=api_key, + model=model, + base_url=cfg["llm"].get("base_url", "https://api.aillm.ru/v1"), + ) + try: + diagnosis = diagnoser.diagnose(SYSTEM_PROMPT, _build_diagnosis_prompt(parsed)) + llm_success = True + except Exception as e: + logger.warning(f"LLM failed: {e}") + diagnosis = format_no_llm(parsed) + f"\n\n(LLM недоступен: {e})" + + llm_duration_ms = int((time.time() - llm_start) * 1000) + + # ── Сохранение в БД ─────────────────────────── + try: + db = Database() + db.save_session( + client_info=client_info, + responses=responses, + diagnosis=diagnosis, + llm_model=model, + llm_duration_ms=llm_duration_ms, + llm_success=llm_success, + ) + except Exception as e: + logger.error(f"DB save failed: {e}") + + return jsonify({ + "diagnosis": diagnosis, + "parsed": _summary(parsed), + "llm_available": llm_available, + "llm_success": llm_success, + }) + + @app.route("/api/v1/chat", methods=["POST"]) + def chat(): + """Свободный вопрос к LLM (без ELM).""" + from flask import request, jsonify + from api.config import load + from brain.client import Diagnoser + + data = request.get_json(silent=True) + if not data or "question" not in data: + return jsonify({"error": "missing 'question'"}), 400 + + question = data["question"].strip() + if not question: + return jsonify({"answer": "Пустой вопрос."}) + + # История диалога + history = data.get("history", []) + history_text = "" + if history: + history_text = "## История диалога\n" + for m in history[-10:]: # последние 10 сообщений + role = "Водитель" if m.get("role") == "user" else "Автоэксперт" + history_text += f"{role}: {m.get('content', '')}\n" + history_text += "\n" + + cfg = load() + api_key = cfg["llm"]["api_key"] + if not api_key: + return jsonify({"answer": "LLM не настроен."}) + + prompt = ( + f"{history_text}" + f"Ты — автоэксперт. Помни контекст диалога выше. " + f"Отвечай КРАТКО, не более 20 строк. Без воды, только по делу.\n\n" + f"Вопрос: {question}" + ) + + try: + diagnoser = Diagnoser( + api_key=api_key, + model=cfg["llm"].get("model", "gpt-oss-120b"), + base_url=cfg["llm"].get("base_url", "https://api.aillm.ru/v1"), + ) + answer = diagnoser.diagnose( + "Ты — лаконичный автоэксперт. Помни контекст диалога. Отвечай кратко, максимум 20 строк.", + prompt, + ) + except Exception as e: + answer = f"LLM недоступен: {e}" + + return jsonify({"answer": answer}) + + @app.route("/api/v1/ping", methods=["GET"]) + def ping(): + """Быстрая проверка доступности сервера.""" + return {"ok": True} + + @app.route("/api/v1/ping-llm", methods=["GET"]) + def ping_llm(): + """Быстрая проверка доступности LLM.""" + from flask import jsonify + from api.config import load + from brain.client import Diagnoser + + cfg = load() + api_key = cfg["llm"]["api_key"] + if not api_key: + return jsonify({"ok": False, "error": "no API key"}) + + t0 = time.time() + try: + diagnoser = Diagnoser( + api_key=api_key, + model=cfg["llm"].get("model", "gpt-oss-120b"), + base_url=cfg["llm"].get("base_url", "https://api.aillm.ru/v1"), + ) + diagnoser.diagnose("Отвечай одним словом.", "OK") + ms = int((time.time() - t0) * 1000) + return jsonify({"ok": True, "ms": ms}) + except Exception as e: + ms = int((time.time() - t0) * 1000) + return jsonify({"ok": False, "ms": ms, "error": str(e)[:100]}) + + +def _summary(p: dict) -> dict: + return { + "vin": p["vin"], + "dtc_stored": p["dtc_stored"], + "dtc_pending": p["dtc_pending"], + "parameters": p["parameters"], + } diff --git a/api/scripts.py b/api/scripts.py new file mode 100644 index 0000000..28a8ccb --- /dev/null +++ b/api/scripts.py @@ -0,0 +1,29 @@ +"""Сборка диагностических скриптов.""" + + +def build_default_script() -> dict: + """Минимальный скрипт для отладки: 1 PID → LLM.""" + return { + "version": 1, + "title": "Экспресс-диагностика", + "steps": [ + {"id": "pid_05", "cmd": "0105", "desc": "Температура ОЖ"}, + ], + } + + +def build_full_script() -> dict: + """Полный скрипт диагностики.""" + return { + "version": 1, + "title": "Полная диагностика", + "steps": [ + {"id": "pid_05", "cmd": "0105", "desc": "Температура ОЖ"}, + {"id": "pid_0C", "cmd": "010C", "desc": "Обороты"}, + {"id": "pid_0D", "cmd": "010D", "desc": "Скорость"}, + {"id": "pid_11", "cmd": "0111", "desc": "Дроссель"}, + {"id": "pid_04", "cmd": "0104", "desc": "Нагрузка"}, + {"id": "pid_06", "cmd": "0106", "desc": "STFT"}, + {"id": "pid_07", "cmd": "0107", "desc": "LTFT"}, + ], + } diff --git a/brain/__init__.py b/brain/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/brain/client.py b/brain/client.py new file mode 100644 index 0000000..3c0e6e1 --- /dev/null +++ b/brain/client.py @@ -0,0 +1,62 @@ +""" +brain/client.py — LLM-клиент для диагностики авто. + +Отправляет запросы к OpenAI-совместимому API (api.aillm.ru). +Используется сервером для анализа данных диагностики. + +## Использование + from brain.client import Diagnoser + d = Diagnoser(api_key="sk-...", model="gpt-oss-120b") + answer = d.diagnose(system_prompt, user_prompt) + +## Модель + gpt-oss-120b — основная (через api.aillm.ru) + qwen3-6-27b-fp8 — быстрая (но CoT leak bug) +""" + +import requests + +DEFAULT_BASE = "https://api.aillm.ru/v1" +DEFAULT_MODEL = "gpt-oss-20b" + + +class Diagnoser: + """Отправляет данные в DeepSeek и возвращает диагноз.""" + + def __init__(self, api_key: str, model: str = DEFAULT_MODEL, base_url: str = DEFAULT_BASE): + self.api_key = api_key + self.model = model + self.base_url = base_url.rstrip("/") + + def ask(self, messages: list[dict]) -> str: + """Отправляет сообщения в DeepSeek, возвращает текст ответа.""" + resp = requests.post( + f"{self.base_url}/chat/completions", + headers={ + "Authorization": f"Bearer {self.api_key}", + "Content-Type": "application/json", + }, + json={ + "model": self.model, + "messages": messages, + "temperature": 0.3, # пониже — меньше фантазий + "max_tokens": 4096, + }, + timeout=120, # api.aillm.ru бывает медленным + ) + resp.raise_for_status() + data = resp.json() + return data["choices"][0]["message"]["content"] + + def diagnose( + self, + system: str, + user_prompt: str, + history: list[dict] | None = None, + ) -> str: + """Полный цикл: system + история + user_prompt → ответ.""" + messages = [{"role": "system", "content": system}] + if history: + messages.extend(history) + messages.append({"role": "user", "content": user_prompt}) + return self.ask(messages) diff --git a/brain/prompts.py b/brain/prompts.py new file mode 100644 index 0000000..6a92585 --- /dev/null +++ b/brain/prompts.py @@ -0,0 +1,87 @@ +""" +brain/prompts.py — промпты для LLM. + +SYSTEM_PROMPT — системный промпт для диагностики (10 правил). +Используется в api/routes.py при формировании запроса к LLM. +""" + +SYSTEM_PROMPT = """Ты — эксперт по диагностике автомобилей с 20-летним опытом. Ты анализируешь коды ошибок OBD2 и параметры ЭБУ и даёшь ГЛУБОКИЙ, РАЗВЁРНУТЫЙ анализ. + +ПРАВИЛА ОТВЕТА: +1. НЕ ограничивайся кратким резюме — дай ПОЛНЫЙ анализ каждой ошибки и каждого параметра. +2. Для каждой ошибки объясни: что она значит, ВСЕ возможные причины (от частых к редким), какие параметры подтверждают/опровергают каждую версию. +3. Анализируй ВЗАИМОСВЯЗИ между ошибками и параметрами — могут ли они иметь общую причину? +4. Указывай степень уверенности в процентах для КАЖДОГО вывода. +5. Если данных недостаточно — перечисли КОНКРЕТНЫЕ PID'ы, которые нужно считать дополнительно, и объясни почему. +6. Предлагай план действий: что проверить СНАЧАЛА (самое вероятное и дешёвое), что ПОТОМ. +7. Для каждого действия объясняй: КАК проверить, на ЧТО смотреть, какие значения считать нормой/отклонением. +8. Добавляй секцию «Если не поможет» — план Б для каждого пункта. +9. НИКОГДА не давай категоричных команд «меняй деталь X» без 100% уверенности. Пиши «проверь X перед заменой Y». +10. Пиши на русском языке, доступно, но ТЕХНИЧЕСКИ ТОЧНО. Используй таблицы где уместно. + +ФОРМАТ ОТВЕТА: +## Диагноз (развёрнутый) +(полный анализ ситуации, 3-5 абзацев) + +## Анализ ошибок +| Код | Описание | Вероятные причины | Подтверждающие параметры | Уверенность | +|-----|----------|-------------------|--------------------------|-------------| +... + +## Анализ параметров +| Параметр | Значение | Норма | Отклонение | На что влияет | +|----------|----------|-------|------------|---------------| +... + +## Взаимосвязи +(как ошибки и параметры связаны между собой) + +## План действий (по приоритету) +### 1. Проверить ... (самое вероятное) +- КАК проверить: ... +- На что смотреть: ... +- Норма: ... + +### 2. Если не помогло — проверить ... +... + +## Каких данных не хватает +- PID XX (название) — потому что ... +- ... + +## Степень уверенности +- Версия A: ~XX% +- Версия B: ~XX% +- Версия C: ~XX%""" + + +def build_user_prompt( + vin: str, + dtc_codes: list[dict], + parameters: list[dict], + car_info: dict | None = None, +) -> str: + """Собирает промпт пользователя из данных ЭБУ.""" + + parts = [f"## Данные диагностики\n"] + parts.append(f"**VIN:** {vin}") + + if car_info: + parts.append(f"**Автомобиль:** {car_info.get('make', '?')} {car_info.get('model', '?')} " + f"({car_info.get('year', '?')}), двигатель: {car_info.get('engine', '?')}") + + if dtc_codes: + parts.append("\n### Коды ошибок") + for dtc in dtc_codes: + parts.append(f"- **{dtc['code']}** ({dtc.get('status', 'stored')}): {dtc.get('description', '')}") + + if parameters: + parts.append("\n### Параметры ЭБУ") + for p in parameters: + parts.append(f"- {p['name']} ({p['pid_code']}): {p['value']} {p['unit']}") + + parts.append("\n## Запрос") + parts.append("Дай диагноз на основе этих данных. Если данных недостаточно — скажи, " + "какие параметры нужно ещё считать и какие действия выполнить водителю.") + + return "\n".join(parts) diff --git a/elmer/elm.py b/elmer/elm.py deleted file mode 100644 index 970dd35..0000000 --- a/elmer/elm.py +++ /dev/null @@ -1,208 +0,0 @@ -"""Связь с ELM327 через Bluetooth SPP (pyserial).""" - -import re -import time - -import serial - - -class ELM327: - """Работа с ELM327 по последовательному порту (Bluetooth SPP).""" - - # Стандартные PID'ы для чтения - DEFAULT_PIDS = { - "0105": ("coolant_temp", "°C"), # температура ОЖ - "010C": ("rpm", "об/мин"), # обороты двигателя - "010D": ("speed", "км/ч"), # скорость - "0111": ("throttle_pos", "%"), # положение дросселя - "010B": ("map", "кПа"), # давление впуска (MAP) - "010F": ("iat", "°C"), # температура впуска - "011F": ("runtime_since_start", "с"), # время с запуска - "0104": ("engine_load", "%"), # нагрузка двигателя - "0106": ("stft_b1", "%"), # краткосрочный fuel trim bank 1 - "0107": ("ltft_b1", "%"), # долгосрочный fuel trim bank 1 - } - - def __init__(self, port: str, baudrate: int = 38400, timeout: float = 5.0): - self.port = port - self.ser = serial.Serial( - port=port, - baudrate=baudrate, - timeout=timeout, - bytesize=serial.EIGHTBITS, - parity=serial.PARITY_NONE, - stopbits=serial.STOPBITS_ONE, - ) - - # ── низкоуровневые команды ──────────────────────────── - - def _cmd(self, cmd: str, wait: float = 0.2) -> str: - """Отправляет AT/OBD команду, возвращает сырой ответ.""" - self.ser.reset_input_buffer() - self.ser.write((cmd + "\r").encode()) - time.sleep(wait) - lines = [] - while True: - line = self.ser.readline().decode("utf-8", errors="ignore").strip() - if not line or line == ">": - break - lines.append(line) - return "\n".join(lines) - - # ── инициализация ───────────────────────────────────── - - def init(self) -> bool: - """Сброс и настройка ELM327. Возвращает True если OK.""" - resp = self._cmd("ATZ", wait=1.0) # сброс - if "ELM" not in resp: - return False - self._cmd("ATE0") # выкл эхо - self._cmd("ATL0") # выкл перевод строки - self._cmd("ATSP0") # авто-протокол - self._cmd("ATH1") # вкл заголовки - return True - - # ── чтение VIN ──────────────────────────────────────── - - def read_vin(self) -> str | None: - """Читает VIN (режим 09 PID 02). Возвращает VIN или None.""" - resp = self._cmd("0902", wait=1.5) - # Формат: 014 0: 49 02 01 57 56 57 ... - # Ищем строку с байтами после 49 02 - match = re.search(r"49\s*02\s*(.+)", resp.replace("\n", " ").replace(":", "")) - if not match: - return None - # Собираем HEX байты, переводим в ASCII - hex_bytes = match.group(1).strip().split() - vin = "" - for h in hex_bytes: - h = h.strip() - if len(h) == 2: - try: - vin += chr(int(h, 16)) - except ValueError: - pass - return vin if len(vin) == 17 else None - - # ── чтение ошибок ───────────────────────────────────── - - def read_dtc_codes(self, mode: str = "03") -> list[dict]: - """Читает коды ошибок. - mode: '03' — сохранённые, '07' — ожидающие. - Возвращает [{"code": "P0301", "description": "", "status": "stored"}, ...]. - """ - resp = self._cmd(mode, wait=1.0) - codes = [] - - # Пример ответа: 43 01 33 00 00 00 00 - for line in resp.split("\n"): - line = line.strip() - if not line or "NO DATA" in line.upper(): - continue - - # Ищем HEX-байты после 43 (mode 03 response) или 47 (mode 07) - match = re.search(r"4[37]\s*(.+)", line.replace(":", "")) - if not match: - continue - - hex_bytes = match.group(1).strip().split() - # Парсим по 2 байта на код (первые два байта — количество кодов) - i = 1 # пропускаем байт количества - while i + 1 < len(hex_bytes): - dtc_raw = _decode_dtc(hex_bytes[i], hex_bytes[i + 1]) - if dtc_raw and dtc_raw != "P0000": - codes.append({ - "code": dtc_raw, - "description": "", - "status": "stored" if mode == "03" else "pending", - }) - i += 2 - - return codes - - # ── чтение параметров ───────────────────────────────── - - def read_pid(self, pid: str) -> float | None: - """Читает один PID, возвращает числовое значение или None.""" - resp = self._cmd(pid, wait=0.3) - # Ищем строку ответа: 41 XX YY ZZ ... - match = re.search(r"4[12]\s*" + pid[2:4] + r"\s*(.+)", resp.replace(":", "")) - if not match: - return None - - hex_bytes = match.group(1).strip().split() - if not hex_bytes: - return None - - # Формулы для стандартных PID (SAE J1979) - formulas = { - "05": lambda b: int(b[0], 16) - 40, # coolant °C - "0C": lambda b: (int(b[0], 16) * 256 + int(b[1], 16)) / 4, # RPM - "0D": lambda b: int(b[0], 16), # speed km/h - "11": lambda b: int(b[0], 16) * 100 / 255, # throttle % - "0B": lambda b: int(b[0], 16), # MAP kPa - "0F": lambda b: int(b[0], 16) - 40, # IAT °C - "1F": lambda b: int(b[0], 16) * 256 + int(b[1], 16), # runtime sec - "04": lambda b: int(b[0], 16) * 100 / 255, # load % - "06": lambda b: (int(b[0], 16) - 128) * 100 / 128, # STFT % - "07": lambda b: (int(b[0], 16) - 128) * 100 / 128, # LTFT % - } - - pid_short = pid[2:4] - if pid_short in formulas: - try: - return round(formulas[pid_short](hex_bytes), 1) - except (ValueError, IndexError): - return None - - # Generic: первый байт как raw - try: - return int(hex_bytes[0], 16) - except (ValueError, IndexError): - return None - - def read_all_pids(self, pids: dict | None = None) -> list[dict]: - """Читает все PID'ы из словаря {pid: (name, unit)}. - Возвращает [{"pid_code": "...", "name": "...", "value": ..., "unit": "..."}, ...]. - """ - if pids is None: - pids = self.DEFAULT_PIDS - - results = [] - for pid_code, (name, unit) in pids.items(): - try: - value = self.read_pid(pid_code) - if value is not None: - results.append({ - "pid_code": pid_code, - "name": name, - "value": value, - "unit": unit, - }) - except Exception: - continue - return results - - def close(self): - self._cmd("ATZ", wait=0.5) - self.ser.close() - - -def _decode_dtc(b1: str, b2: str) -> str | None: - """Декодирует два HEX-байта в код ошибки вида P0301.""" - try: - a, b = int(b1, 16), int(b2, 16) - except ValueError: - return None - - # Первые 2 бита первого байта — тип: - types = {0: "P", 1: "C", 2: "B", 3: "U"} - prefix = types.get(a >> 6, "?") - - # Оставшиеся биты - d1 = str((a >> 4) & 0x03) # вторая цифра - d2 = str(a & 0x0F) # третья цифра - d3 = f"{(b >> 4) & 0x0F:X}" # четвёртая цифра (hex!) - d4 = f"{b & 0x0F:X}" # пятая цифра (hex!) - - return f"{prefix}{d1}{d2}{d3}{d4}" diff --git a/elmer/elm_proto.py b/elmer/elm_proto.py deleted file mode 100644 index cb00188..0000000 --- a/elmer/elm_proto.py +++ /dev/null @@ -1,230 +0,0 @@ -""" -ELM327 Protocol Layer — точная копия паттернов AndrOBD (1993⭐, 10 лет продакшена). - -Источник: github.com/fr3ts0n/AndrOBD - - BtCommService.java → BT подключение + 500мс пауза - - ElmProt.java → стейт-машина, таймауты, ошибки - - StreamHandler.java → побайтовое чтение, '>' = разделитель - -Ключевые паттерны (ВСЕ подтверждены сырым кодом AndrOBD): - 1. Побайтовое чтение, сон 1мс между проверками - 2. '>' — обычный разделитель строк (как CR/LF), НЕ спецсигнал - 3. Адаптивный таймаут: старт 200мс, ±4мс, диапазон 12-1000мс - 4. ATST меняется на лету при изменении таймаута - 5. Инициализация: ATSP→ATAT→ATS0→ATL0→ATE0 (без ATZ!) - 6. Ошибки: WARMSTART, re-queue, protocol reset - 7. Мульти-фрейм ISO-TP с префиксом длины - 8. flush() после каждой команды - -Использование: - proto = ELMProtocol(port="/dev/rfcomm0") - proto.connect() # + 500мс пауза (AndrOBD #233) - proto.init() # ATSP0 → ATAT1 → ATS0 → ATL0 → ATE0 - resp = proto.send_command("010C") # → "410C1AF8" - proto.close() -""" - -import logging -import time -from typing import Optional - -logger = logging.getLogger("elmer.proto") - - -# ── Адаптивный таймаут (AndrOBD AdaptiveTiming) ───────────── - -class AdaptiveTiming: - """Адаптивный таймаут ELM327 — точная копия AndrOBD. - - Старт: 200мс. Шаг: 4мс. Диапазон: 12-1000мс. - На NODATA: +4мс. На успех: -4мс (не ниже learned_min). - Каждое изменение → ATST. - """ - - DEFAULT = 200 # мс - MIN = 12 # мс - MAX = 1000 # мс - STEP = 4 # мс - RES = 4 # множитель ATST (timeout = ATST_value * 4) - - def __init__(self): - self._timeout = self.DEFAULT - self._learned_min = self.MIN - - @property - def timeout_ms(self) -> int: - return self._timeout - - @property - def atst_value(self) -> int: - return max(1, self._timeout // self.RES) - - def increase(self): - if self._timeout + self.STEP < self.MAX: - self._timeout += self.STEP - - def decrease(self): - if self._timeout - self.STEP >= self._learned_min: - self._timeout -= self.STEP - - def reset(self): - self._timeout = self.DEFAULT - - -# ── Протокол ELM327 ───────────────────────────────────────── - -class ELMProtocol: - """Побайтовый обмен с ELM327 — точная копия AndrOBD StreamHandler + ElmProt.""" - - SPP_UUID = "00001101-0000-1000-8000-00805F9B34FB" - - def __init__(self, port: str, baudrate: int = 38400): - self.port = port - self.baudrate = baudrate - self._ser = None - self._timing = AdaptiveTiming() - self._last_cmd: Optional[str] = None - - # ── Подключение ─────────────────────────────────────── - - def connect(self): - """Открывает serial-порт + 500мс пауза (AndrOBD issue #233).""" - import serial - self._ser = serial.Serial( - port=self.port, - baudrate=self.baudrate, - timeout=0.1, - bytesize=serial.EIGHTBITS, - parity=serial.PARITY_NONE, - stopbits=serial.STOPBITS_ONE, - ) - time.sleep(0.5) # КРИТИЧЕСКИ: AndrOBD #233 - logger.info(f"ELM: connected {self.port} @ {self.baudrate}") - - def close(self): - if self._ser and self._ser.is_open: - self._ser.close() - - # ── Инициализация (AndrOBD ElmProt.initialize) ───────── - - def init(self) -> bool: - """ATSP0→ATAT1→ATS0→ATL0→ATE0. Каждая команда с чтением ответа.""" - logger.info("ELM: init start") - - self.send_command("ATSP0") # авто-протокол - self.send_command("ATAT1") # adaptive timing - self._update_timeout() - self.send_command("ATS0") # пробелы выкл - self.send_command("ATL0") # line feeds выкл - self.send_command("ATE0") # эхо выкл - - logger.info("ELM: init done") - return True - - def _update_timeout(self): - self._write(f"ATST{self._timing.atst_value:02X}") - - def _drain(self): - if self._ser and self._ser.in_waiting > 0: - n = len(self._ser.read(self._ser.in_waiting)) - logger.debug(f"ELM: drained {n}B") - - # ── Отправка + чтение (AndrOBD StreamHandler) ────────── - - def send_command(self, cmd: str) -> str: - """Отправляет команду, читает ответ побайтово. - - Разделители: CR(13), LF(10), '>'(62) — все равноправны. - Возвращает строки ответа через \\n. - """ - self._write(cmd) - self._last_cmd = cmd - - try: - result = self._read(self._timing.timeout_ms) - except TimeoutError: - self._timing.increase() - self._update_timeout() - try: - result = self._read(self._timing.timeout_ms) - except TimeoutError: - return "" - - self._handle_response(result) - return result - - # ── Внутренние ──────────────────────────────────────── - - def _write(self, cmd: str): - """cmd + CR + flush (AndrOBD writeTelegram).""" - self._ser.write((cmd + "\r").encode()) - self._ser.flush() - logger.debug(f"ELM → {cmd}") - - def _read(self, timeout_ms: int) -> str: - """Побайтовое чтение, пауза 1мс (AndrOBD StreamHandler.run).""" - deadline = time.monotonic() + timeout_ms / 1000.0 - lines: list[str] = [] - cur: list[str] = [] - - while time.monotonic() < deadline: - if self._ser.in_waiting > 0: - ch = self._ser.read(1) - if not ch: - continue - cp = ch[0] - - if cp == 62: # '>' — разделитель как CR/LF - self._push(cur, lines) - break - elif cp == 13: # CR - self._push(cur, lines) - elif cp in (10, 32): # LF и пробел — игнорируем - pass - else: - cur.append(chr(cp)) - else: - time.sleep(0.001) - - self._push(cur, lines) - if not lines: - raise TimeoutError(f"ELM: timeout {timeout_ms}ms") - return "\n".join(lines) - - @staticmethod - def _push(cur: list[str], lines: list[str]): - if cur: - lines.append("".join(cur)) - cur.clear() - - def _handle_response(self, raw: str): - """Обработка ошибок (AndrOBD ElmProt.handleTelegram).""" - u = raw.upper() - - if "SEARCHING" in u: - return - - if "NODATA" in u or "NO DATA" in u: - self._timing.increase() - self._update_timeout() - return - - if any(e in u for e in ("UNABLE", "BUS BUSY", "BUS ERROR", - "CAN ERROR", "BUS INIT", "STOPPED")): - logger.warning(f"ELM: bus error, resetting — {raw[:60]}") - self._timing.reset() - self._update_timeout() - self._write("ATPC") - self._write("ATSP0") - return - - # Успех — уменьшаем таймаут - self._timing.decrease() - - -# ── Фабрика ───────────────────────────────────────────────── - -def create(port: str, baudrate: int = 38400) -> ELMProtocol: - proto = ELMProtocol(port, baudrate) - proto.connect() - return proto diff --git a/obd/__init__.py b/obd/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/obd/protocol.py b/obd/protocol.py new file mode 100644 index 0000000..8937601 --- /dev/null +++ b/obd/protocol.py @@ -0,0 +1,348 @@ +""" +obd/protocol.py — ELM327 стейт-машина (AndrOBD). + +Точная копия логики из AndrOBD (ElmProt.java, github.com/fr3ts0n/AndrOBD). + +## Архитектура + ┌──────────┐ команда ┌──────────┐ + │ READY │──────────────▶│ BUSY │ + └──────────┘ └────┬─────┘ + ▲ │ ответ получен + │ ┌────────────────┘ + │ ▼ + ┌────┴─────┐ ошибка ┌──────────┐ + │ ERROR │◀─────────│ (любое) │ + └────┬─────┘ └──────────┘ + │ восстановление ▲ + └─────────────────────┘ + +## Использование + elm = AndrOBD("/dev/rfcomm0", 38400) + elm.connect() + elm.init() + vin = elm.send("0902") + rpm = elm.send("010C") + elm.close() + +## Ключевые особенности + - Байт-за-байтом чтение с 1мс поллингом + - `>` как разделитель ответов (промпт ELM327) + - Адаптивный таймаут (200мс ± 4мс, ATST) + - Восстановление после BUS ERROR (ATPC → ATSP0) + - Не тот ответ → переход в ERROR → восстановление +""" + +import logging +import time +from enum import Enum, auto +from typing import Optional + +logger = logging.getLogger("androbd") + + +class State(Enum): + """Состояния стейт-машины (AndrOBD STAT). + + UNDEFINED → INITIALIZING → READY — нормальный путь. + BUSY — во время выполнения команды. + ERROR/DISCONNECTED — ошибка, требуется восстановление. + """ + UNDEFINED = auto() + INITIALIZING = auto() + READY = auto() + BUSY = auto() + ERROR = auto() + DISCONNECTED = auto() + + +class Rsp: + """Классификация ответов ELM327 (AndrOBD RSP_ID). + + Каждый сырой ответ классифицируется: + - PROMPT (`>`) — готов к следующей команде + - OK — команда выполнена + - SEARCHING — идёт поиск протокола + - BUS_ERROR/BUS_BUSY/CAN_ERROR/STOPPED — ошибка шины → DISCONNECTED + - ERROR/DATA_ERROR/BUFFER_FULL — ошибка → warm start + - UNKNOWN — данные (ответ на PID/DTC) + """ + PROMPT = ">" + OK = "OK" + SEARCHING = "SEARCHING" + NODATA = "NODATA" + ERROR = "ERROR" + UNABLE = "UNABLE" + BUS_BUSY = "BUS BUSY" + BUS_ERROR = "BUS ERROR" + CAN_ERROR = "CAN ERROR" + BUS_INIT = "BUS INIT" + STOPPED = "STOPPED" + DATA_ERROR = "DATA ERROR" + BUFFER_FULL= "BUFFER FULL" + RX_ERROR = "RX ERROR" + UNKNOWN = "" + + @classmethod + def identify(cls, raw: str) -> str: + """Определяет тип ответа по сырой строке.""" + u = raw.upper().strip() + for tag in (cls.SEARCHING, cls.NODATA, cls.ERROR, cls.UNABLE, + cls.BUS_BUSY, cls.BUS_ERROR, cls.CAN_ERROR, + cls.BUS_INIT, cls.STOPPED, cls.DATA_ERROR, + cls.BUFFER_FULL, cls.RX_ERROR, cls.OK): + if u.startswith(tag): + return tag + if raw.strip() == ">": + return cls.PROMPT + return cls.UNKNOWN + + +class AdaptiveTiming: + """Адаптивный таймаут ожидания ответа (AndrOBD AdaptiveTiming). + + Динамически подстраивается под скорость ответа ЭБУ: + - Успешный ответ → уменьшаем таймаут (быстрее) + - Таймаут/NO DATA → увеличиваем таймаут (медленнее) + - BUS ERROR → сброс до DEFAULT + + ATST = таймаут / 4 (отправляется в ELM327 как ATSTxx). + """ + DEFAULT = 500 # мс — начальный таймаут + MIN = 50 # мс — минимальный + MAX = 2000 # мс — максимальный + STEP = 20 # мс — шаг изменения + RES = 4 # делитель для ATST + + def __init__(self): + self._t = self.DEFAULT + self._min = self.MIN + + @property + def ms(self) -> int: + """Текущий таймаут в миллисекундах.""" + return self._t + + @property + def atst(self) -> int: + """Значение для ATST (таймаут / 4).""" + return max(1, self._t // self.RES) + + def increase(self): + """Увеличить таймаут (ЭБУ медленно отвечает).""" + if self._t + self.STEP < self.MAX: + self._t += self.STEP + + def decrease(self): + """Уменьшить таймаут (ЭБУ отвечает быстро).""" + if self._t - self.STEP >= self._min: + self._t -= self.STEP + + def reset(self): + """Сбросить до DEFAULT (после BUS ERROR).""" + self._t = self.DEFAULT + + +class AndrOBD: + """Стейт-машина ELM327 — 1:1 копия AndrOBD (ElmProt.java). + + Управляет жизненным циклом ELM327: + 1. connect() — открыть serial/Bluetooth порт + 2. init() — инициализация (ATSP0, ATAT1, ATST, ATS0, ATL0, ATE0) + 3. send(cmd) — отправить OBD-команду, получить ответ + 4. close() — закрыть порт + + Автоматически обрабатывает: таймауты, BUS ERROR, восстановление. + """ + + INIT_TMO = 10000 # мс — таймаут для команд инициализации + DEF_TMO = 200 # мс — начальный таймаут (заменяется AdaptiveTiming) + + def __init__(self, port: str, baudrate: int = 38400): + """port — устройство (напр. /dev/rfcomm0), baudrate — скорость.""" + self.port = port + self.baudrate = baudrate + self._ser = None + self._timing = AdaptiveTiming() + self._state = State.UNDEFINED + self._last_cmd: Optional[str] = None + + def connect(self): + """Открыть serial-соединение с ELM327.""" + import serial + self._ser = serial.Serial( + port=self.port, baudrate=self.baudrate, timeout=0.1, + bytesize=serial.EIGHTBITS, parity=serial.PARITY_NONE, + stopbits=serial.STOPBITS_ONE) + time.sleep(0.5) + logger.info(f"AndrOBD: connected {self.port}") + + def close(self): + """Закрыть serial-соединение.""" + if self._ser and self._ser.is_open: + self._ser.close() + + def init(self): + """Инициализация ELM327 — 6 AT-команд. + + ATSP0 — авто-протокол + ATAT1 — адаптивный таймаут вкл + ATSTxx — установить таймаут + ATS0 — без пробелов в ответах + ATL0 — без перевода строки + ATE0 — без эха + """ + logger.info("AndrOBD: init") + self._state = State.INITIALIZING + self._exec("ATSP0", self.INIT_TMO) + self._exec("ATAT1", self.DEF_TMO * 5) + self._update_atst() + self._exec("ATS0", self.DEF_TMO * 5) + self._exec("ATL0", self.DEF_TMO * 5) + self._exec("ATE0", self.DEF_TMO * 5) + self._state = State.READY + logger.info("AndrOBD: ready") + + def send(self, cmd: str) -> str: + """Отправить OBD-команду и получить ответ. + + cmd — команда (напр. '0105', '0902', '03'). + Возвращает сырой ответ ELM327. + При ошибке — авто-восстановление. + """ + if self._state == State.ERROR: + self._recover() + self._state = State.BUSY + result = self._exec(cmd, self._timing.ms) + self._state = State.READY + return result + + # ── Приватные методы ────────────────────────────────── + + def _exec(self, cmd: str, timeout: int) -> str: + """Выполнить команду с таймаутом и ретраями (до 10 попыток).""" + self._last_cmd = cmd + self._write(cmd) + t = timeout + for _ in range(10): + try: + return self._handle(self._read(t)) + except TimeoutError: + if self._state == State.INITIALIZING: + t += 1000 + else: + self._timing.increase() + t = self._timing.ms + logger.error(f"AndrOBD: no response for {cmd}") + self._state = State.ERROR + return "" + + def _handle(self, raw: str) -> str: + """Обработать ответ ELM327: классифицировать и обновить таймаут. + + Возвращает raw как есть — обработка данных делается выше. + """ + t = Rsp.identify(raw) + + if t == Rsp.SEARCHING: + return raw + if t == Rsp.OK: + self._timing.decrease() + return raw + if t == Rsp.NODATA: + self._timing.increase() + self._update_atst() + return raw + + # BUS ERROR — сброс протокола + if t in (Rsp.UNABLE, Rsp.BUS_BUSY, Rsp.BUS_ERROR, + Rsp.CAN_ERROR, Rsp.BUS_INIT, Rsp.STOPPED): + logger.warning(f"AndrOBD: BUS ERROR ({t})") + self._state = State.DISCONNECTED + self._timing.reset() + self._update_atst() + self._write("ATPC") # закрыть протокол + self._try_read() + self._write("ATSP0") # переоткрыть авто-протокол + self._try_read() + return raw + + # Другие ошибки — warm start + if t in (Rsp.ERROR, Rsp.DATA_ERROR, Rsp.BUFFER_FULL, Rsp.RX_ERROR): + logger.warning(f"AndrOBD: {t} — warm start") + self._state = State.ERROR + self._write("ATWS") + self._try_read() + return raw + + # Данные — успешный ответ + self._timing.decrease() + return raw + + def _recover(self): + """Восстановление после ошибки: ATWS → ATSP0 → ATE0.""" + logger.info("AndrOBD: recovering...") + self._state = State.INITIALIZING + self._write("ATWS") + self._try_read() + self._write("ATSP0") + self._try_read() + self._write("ATE0") + self._try_read() + self._state = State.READY + + def _write(self, cmd: str): + """Отправить команду в ELM327 (добавляет CR, flush).""" + self._ser.write((cmd + "\r").encode()) + self._ser.flush() + logger.debug(f"AndrOBD → {cmd}") + + def _read(self, timeout_ms: int) -> str: + """Прочитать ответ ELM327 байт-за-байтом. + + Читает до символа `>` (промпт) или до таймаута. + Возвращает сырой ответ без `>`. + """ + dl = time.monotonic() + timeout_ms / 1000.0 + lines, cur = [], [] + got_prompt = False + while time.monotonic() < dl: + if self._ser.in_waiting > 0: + ch = self._ser.read(1) + if not ch: + continue + cp = ch[0] + if cp == 62: # '>' — промпт ELM327 + self._push(cur, lines) + got_prompt = True + break + elif cp == 13: # CR — конец строки + self._push(cur, lines) + elif cp in (10, 32): # LF и пробел — игнорируем + pass + else: + cur.append(chr(cp)) + else: + time.sleep(0.001) # поллинг 1мс + self._push(cur, lines) + if not got_prompt: + raise TimeoutError(f"timeout {timeout_ms}ms") + return "\n".join(lines) + + def _try_read(self, timeout: int = 5000): + """Прочитать и проигнорировать ответ (для команд восстановления).""" + try: + self._read(timeout) + except TimeoutError: + pass + + @staticmethod + def _push(cur, lines): + """Добавить накопленные байты как строку в lines.""" + if cur: + lines.append("".join(cur)) + cur.clear() + + def _update_atst(self): + """Отправить ATST с текущим значением адаптивного таймаута.""" + self._write(f"ATST{self._timing.atst:02X}") + self._try_read() diff --git a/web/app.py b/web/app.py index c247052..9d0ca95 100644 --- a/web/app.py +++ b/web/app.py @@ -1,33 +1,31 @@ -"""Elmer Web UI — Flask-приложение для локального тестирования. +"""elmAI — точка входа Flask-приложения. -Заодно прототип будущего серверного API. -Запуск: DEEPSEEK_API_KEY=sk-... python web/app.py -Открыть: http://localhost:5005 +Запуск через gunicorn: + gunicorn -w 4 -b 127.0.0.1:8000 web.app:app + +Структура модулей: + api/ — REST-эндпоинты, БД, скрипты + brain/ — LLM-клиент, промпты + obd/ — ELM327-протокол (AndrOBD) """ import sys import logging from pathlib import Path -# Добавляем корень проекта в PYTHONPATH +# Добавляем корень проекта в PYTHONPATH для импорта api/, brain/, obd/ sys.path.insert(0, str(Path(__file__).parent.parent)) from flask import Flask, jsonify, redirect, render_template, request, send_from_directory -from elmer.config import load -from elmer.db import Database -from elmer.diagnose import Diagnoser -from elmer.elm import ELM327 -from elmer.prompts import SYSTEM_PROMPT, build_user_prompt -from web.raw_endpoint import register as register_raw_endpoint -from web.script_endpoint import register as register_script_endpoint +from api.config import load +from api.routes import register as register_api logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(name)s] %(message)s") app = Flask(__name__) config = load() -register_raw_endpoint(app) -register_script_endpoint(app) +register_api(app) @app.route("/") @@ -42,67 +40,6 @@ def download_apk(): return send_from_directory("static", "app-debug.apk", as_attachment=True, download_name="elmer.apk") -@app.route("/api/diagnose", methods=["POST"]) -def diagnose(): - """Подключается к ELM327, читает данные, отправляет в LLM, возвращает результат.""" - elm_cfg = config["elm327"] - - # ── 1. ELM327 ──────────────────────────────── - try: - elm = ELM327(port=elm_cfg["port"], baudrate=elm_cfg["baudrate"]) - except Exception as e: - return jsonify({"error": f"Не удалось открыть порт {elm_cfg['port']}: {e}"}), 500 - - if not elm.init(): - elm.close() - return jsonify({"error": "ELM327 не ответил на ATZ"}), 500 - - # ── 2. VIN ────────────────────────────────── - vin = elm.read_vin() - if not vin: - elm.close() - return jsonify({"error": "Не удалось прочитать VIN (режим 09 не поддерживается?)"}), 500 - - # ── 3. Ошибки ─────────────────────────────── - dtc_codes = elm.read_dtc_codes("03") + elm.read_dtc_codes("07") - - # ── 4. Параметры ──────────────────────────── - parameters = elm.read_all_pids(config.get("pids")) - - elm.close() - - # ── 5. DeepSeek ───────────────────────────── - api_key = config["deepseek"]["api_key"] - if not api_key: - return jsonify({"error": "DEEPSEEK_API_KEY не задан"}), 500 - - diagnoser = Diagnoser( - api_key=api_key, - model=config["deepseek"].get("model", "deepseek-chat"), - ) - user_prompt = build_user_prompt(vin, dtc_codes, parameters) - answer = diagnoser.diagnose(SYSTEM_PROMPT, user_prompt) - - # ── 6. SQLite ─────────────────────────────── - db = Database() - car_id = db.get_or_create_car(vin) - token_id = db.create_token(car_id) - for dtc in dtc_codes: - db.add_dtc(token_id, dtc["code"], dtc.get("description", ""), dtc["status"]) - for p in parameters: - db.add_parameter(token_id, p["pid_code"], p["name"], p["value"], p["unit"]) - db.add_llm_message(token_id, "user", user_prompt) - db.add_llm_message(token_id, "assistant", answer) - - return jsonify({ - "vin": vin, - "dtc_codes": dtc_codes, - "parameters": parameters, - "diagnosis": answer, - "token_id": token_id, - }) - - if __name__ == "__main__": - print(f"🌐 Elmer Web: http://localhost:5005") + print(f"🌐 elmAI Web: http://localhost:5005") app.run(host="0.0.0.0", port=5005, debug=False) diff --git a/web/raw_endpoint.py b/web/raw_endpoint.py deleted file mode 100644 index da3474a..0000000 --- a/web/raw_endpoint.py +++ /dev/null @@ -1,350 +0,0 @@ -"""Эндпоинт /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 diff --git a/web/templates/index.html b/web/templates/index.html index ddcadcb..531b2df 100644 --- a/web/templates/index.html +++ b/web/templates/index.html @@ -12,14 +12,14 @@

🔧 elmAI

Диагностика авто через ELM327 + ИИ

-

v0.27.0-dev — 29 мая 2026

+

v0.28.0-dev — 29 мая 2026

📱 Скачай приложение на телефон:

⬇️ Скачать elmAI APK -

v0.27.0-dev • нажмите чтобы скачать

+

v0.28.0-dev • нажмите чтобы скачать