fix(server): идемпотентность upload, WAL, буфер ELM, таймаут LLM

- api/db.py: WAL mode, busy_timeout, request_id UNIQUE, close(), контекстный менеджер
- api/routes.py: проверка request_id при upload, /ping-llm кэш 60с, /chat через roles
- api/config.py: lru_cache на load()
- obd/protocol.py: reset_input_buffer перед _write(), условный READY в send()
- brain/client.py: модель 120b, таймаут из параметра, LLMError класс, обработка 429/5xx
This commit is contained in:
Repinoid
2026-05-31 16:56:21 +03:00
parent eab73cd76a
commit 6893ae33f5
5 changed files with 177 additions and 100 deletions
+2
View File
@@ -1,6 +1,7 @@
"""Загрузка конфигурации из config.yaml."""
import os
from functools import lru_cache
from pathlib import Path
import yaml
@@ -8,6 +9,7 @@ import yaml
CONFIG_PATH = Path(os.environ.get("ELMER_CONFIG", Path(__file__).parent.parent / "config.yaml"))
@lru_cache(maxsize=1)
def load() -> dict:
"""Читает config.yaml, подставляет переменные окружения в значения."""
if not CONFIG_PATH.exists():
+47 -11
View File
@@ -18,10 +18,24 @@ 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 = sqlite3.connect(str(self.path), timeout=30, check_same_thread=False)
self.conn.row_factory = sqlite3.Row
self.conn.execute("PRAGMA journal_mode=WAL")
self.conn.execute("PRAGMA busy_timeout=30000")
self._init_schema()
def __enter__(self):
return self
def __exit__(self, exc_type, exc_val, exc_tb):
self.close()
return False
def close(self):
if self.conn:
self.conn.close()
self.conn = None
def _init_schema(self):
self.conn.executescript("""
CREATE TABLE IF NOT EXISTS cars (
@@ -113,22 +127,40 @@ class Database:
llm_success INTEGER DEFAULT 0,
-- Сырые данные (JSON)
raw_responses TEXT
raw_responses TEXT,
-- Идемпотентность
request_id TEXT UNIQUE,
response_json 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);
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);
CREATE INDEX IF NOT EXISTS idx_sessions_request_id ON sessions(request_id);
""")
self.conn.commit()
# ── sessions ──────────────────────────────────────────
def get_cached_response(self, request_id: str) -> dict | None:
"""Возвращает сохранённый ответ сессии по request_id, или None."""
row = self.conn.execute(
"SELECT response_json FROM sessions WHERE request_id = ?", (request_id,)
).fetchone()
if row and row["response_json"]:
return json.loads(row["response_json"])
return None
def save_session(self, client_info: dict, responses: list[dict],
diagnosis: str = "", llm_model: str = "",
llm_duration_ms: int = 0, llm_success: bool = False):
"""Сохраняет сводную запись о сессии."""
llm_duration_ms: int = 0, llm_success: bool = False,
request_id: str = "", response_json: dict | None = None):
"""Сохраняет сводную запись о сессии.
Если request_id передан и уже существует — silently return (идемпотентность).
"""
ci = client_info
# Подсчёт DTC/PID из ответов
@@ -151,8 +183,10 @@ class Database:
vin = None
break
resp_json_str = json.dumps(response_json, ensure_ascii=False) if response_json else None
self.conn.execute("""
INSERT INTO sessions (
INSERT OR IGNORE INTO sessions (
client_ip, real_ip, user_agent, content_length,
phone_model, phone_maker, android_version, android_sdk,
app_version, android_id,
@@ -163,9 +197,9 @@ class Database:
transport, mock_mode,
diagnosis_text, diagnosis_len, llm_model,
llm_duration_ms, llm_success,
raw_responses
raw_responses, request_id, response_json
) VALUES (?,?,?,?, ?,?,?,?, ?,?, ?,?,?, ?,?,?, ?,?,?, ?,?,?, ?,?,
?,?,?, ?,?, ?)
?,?,?, ?,?, ?,?,?)
""", (
ci.get("client_ip"), ci.get("real_ip"), ci.get("user_agent"),
ci.get("content_length"),
@@ -179,6 +213,8 @@ class Database:
diagnosis, len(diagnosis), llm_model,
llm_duration_ms, 1 if llm_success else 0,
json.dumps(responses, ensure_ascii=False) if responses else None,
request_id if request_id else None,
resp_json_str,
))
self.conn.commit()
+76 -67
View File
@@ -6,11 +6,21 @@ POST /api/v1/session/upload — приём батча, LLM-анализ, воз
import logging
import time
from flask import jsonify, request
from api.config import load
from api.db import Database
from api.parser import format_no_llm, parse_batch
from api.scripts import build_default_script, build_full_script
from api.parser import parse_batch, format_no_llm
from brain.client import Diagnoser, LLMError
from brain.prompts import SYSTEM_PROMPT
logger = logging.getLogger("elmer.script")
# Кэш для /ping-llm (60 секунд)
_ping_llm_cache: dict = {}
def _build_diagnosis_prompt(data: dict) -> str:
"""Строит промпт для LLM из распарсенных данных."""
@@ -55,19 +65,12 @@ def register(app):
@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
@@ -75,6 +78,15 @@ def register(app):
responses = data["responses"]
logger.info(f"Upload: {len(responses)} responses")
# ── Идемпотентность: проверяем request_id ─────
request_id = (data.get("request_id") or "").strip()
if request_id:
with Database() as db:
cached = db.get_cached_response(request_id)
if cached is not None:
logger.info(f"Upload: cached response for {request_id}")
return jsonify(cached), 200
# ── Информация о клиенте ──────────────────────
client_info = data.get("client_info", {})
client_info["client_ip"] = request.remote_addr
@@ -104,40 +116,40 @@ def register(app):
try:
diagnosis = diagnoser.diagnose(SYSTEM_PROMPT, _build_diagnosis_prompt(parsed))
llm_success = True
except Exception as e:
except LLMError as e:
logger.warning(f"LLM failed: {e}")
diagnosis = format_no_llm(parsed) + f"\n\n(LLM недоступен: {e})"
diagnosis = format_no_llm(parsed) + f"\n\n({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({
response = {
"diagnosis": diagnosis,
"parsed": _summary(parsed),
"llm_available": llm_available,
"llm_success": llm_success,
})
}
# ── Сохранение в БД ───────────────────────────
try:
with Database() as db:
db.save_session(
client_info=client_info,
responses=responses,
diagnosis=diagnosis,
llm_model=model,
llm_duration_ms=llm_duration_ms,
llm_success=llm_success,
request_id=request_id,
response_json=response if request_id else None,
)
except Exception as e:
logger.error(f"DB save failed: {e}")
return jsonify(response)
@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
@@ -146,27 +158,18 @@ def register(app):
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}"
)
# История диалога — передаём как массив messages с ролями
history_raw = data.get("history", [])
history_msgs = [
{"role": m["role"], "content": m["content"]}
for m in history_raw[-10:]
if isinstance(m, dict) and "role" in m and "content" in m
]
try:
diagnoser = Diagnoser(
@@ -176,10 +179,12 @@ def register(app):
)
answer = diagnoser.diagnose(
"Ты — лаконичный автоэксперт. Помни контекст диалога. Отвечай кратко, максимум 20 строк.",
prompt,
question,
history=history_msgs if history_msgs else None,
)
except Exception as e:
answer = f"LLM недоступен: {e}"
except LLMError as e:
logger.warning(f"Chat LLM failed: {e}")
answer = str(e)
return jsonify({"answer": answer})
@@ -190,29 +195,33 @@ def register(app):
@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
"""Быстрая проверка доступности LLM (с кэшем 60с)."""
nonlocal _ping_llm_cache
now = time.time()
if _ping_llm_cache and (now - _ping_llm_cache.get("ts", 0)) < 60:
return jsonify(_ping_llm_cache["data"])
cfg = load()
api_key = cfg["llm"]["api_key"]
if not api_key:
return jsonify({"ok": False, "error": "no API key"})
result = {"ok": False, "error": "no API key"}
else:
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)
result = {"ok": True, "ms": ms}
except Exception as e:
ms = int((time.time() - t0) * 1000)
result = {"ok": False, "ms": ms, "error": "LLM unavailable"}
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]})
_ping_llm_cache = {"ts": now, "data": result}
return jsonify(result)
def _summary(p: dict) -> dict:
+49 -21
View File
@@ -14,39 +14,67 @@ brain/client.py — LLM-клиент для диагностики авто.
qwen3-6-27b-fp8 — быстрая (но CoT leak bug)
"""
import logging
import requests
logger = logging.getLogger("brain.client")
DEFAULT_BASE = "https://api.aillm.ru/v1"
DEFAULT_MODEL = "gpt-oss-20b"
DEFAULT_MODEL = "gpt-oss-120b"
DEFAULT_TIMEOUT = 180
class LLMError(Exception):
"""Ошибка LLM API с безопасным для клиента сообщением."""
pass
class Diagnoser:
"""Отправляет данные в DeepSeek и возвращает диагноз."""
"""LLM-клиент для OpenAI-совместимого API."""
def __init__(self, api_key: str, model: str = DEFAULT_MODEL, base_url: str = DEFAULT_BASE):
def __init__(self, api_key: str, model: str = DEFAULT_MODEL,
base_url: str = DEFAULT_BASE, timeout: int = DEFAULT_TIMEOUT):
self.api_key = api_key
self.model = model
self.base_url = base_url.rstrip("/")
self.timeout = timeout
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"]
"""Отправляет сообщения в LLM API, возвращает текст ответа."""
try:
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=self.timeout,
)
resp.raise_for_status()
data = resp.json()
return data["choices"][0]["message"]["content"]
except requests.Timeout:
logger.warning(f"LLM timeout after {self.timeout}s")
raise LLMError("LLM не ответил вовремя. Попробуйте позже.")
except requests.HTTPError as e:
status = e.response.status_code if e.response is not None else 0
logger.warning(f"LLM HTTP {status}: {e}")
if status == 429:
raise LLMError("Слишком много запросов. Подождите минуту.")
elif 500 <= status < 600:
raise LLMError("LLM временно недоступен. Попробуйте позже.")
else:
raise LLMError("Ошибка LLM. Попробуйте позже.")
except Exception as e:
logger.error(f"LLM unexpected: {e}")
raise LLMError("LLM временно недоступен.")
def diagnose(
self,
+3 -1
View File
@@ -213,7 +213,8 @@ class AndrOBD:
self._recover()
self._state = State.BUSY
result = self._exec(cmd, self._timing.ms)
self._state = State.READY
if self._state == State.BUSY:
self._state = State.READY
return result
# ── Приватные методы ──────────────────────────────────
@@ -292,6 +293,7 @@ class AndrOBD:
def _write(self, cmd: str):
"""Отправить команду в ELM327 (добавляет CR, flush)."""
self._ser.reset_input_buffer() # сбросить хвосты предыдущего ответа
self._ser.write((cmd + "\r").encode())
self._ser.flush()
logger.debug(f"AndrOBD → {cmd}")