refactor: архитектура сервера — obd/ brain/ api/

- obd/protocol.py — ELM327 стейт-машина AndrOBD (подробные докстринги)
- brain/client.py — LLM-клиент Diagnoser
- brain/prompts.py — системные промпты
- api/db.py — SQLite (sessions + старые таблицы)
- api/scripts.py — сборка диагностических скриптов
- api/parser.py — парсинг ответов ELM327
- api/routes.py — все REST-эндпоинты
- api/config.py — загрузка конфига
- web/app.py — точка входа Flask (упрощена)

Удалено: elmer/elm.py, elmer/elm_proto.py, web/raw_endpoint.py
This commit is contained in:
Repinoid
2026-05-29 17:24:33 +03:00
parent 2028c7ff0e
commit 2db6cb68a8
16 changed files with 1209 additions and 866 deletions
View File
+39
View File
@@ -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)
+275
View File
@@ -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]
+130
View File
@@ -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)
+224
View File
@@ -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"],
}
+29
View File
@@ -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"},
],
}