diff --git a/api/config.py b/api/config.py index fa27765..1038218 100644 --- a/api/config.py +++ b/api/config.py @@ -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(): diff --git a/api/db.py b/api/db.py index de1e4cd..5de5648 100644 --- a/api/db.py +++ b/api/db.py @@ -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() diff --git a/api/routes.py b/api/routes.py index b856d4c..0e8222f 100644 --- a/api/routes.py +++ b/api/routes.py @@ -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: diff --git a/brain/client.py b/brain/client.py index 3c0e6e1..7f6ef5e 100644 --- a/brain/client.py +++ b/brain/client.py @@ -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, diff --git a/obd/protocol.py b/obd/protocol.py index 8937601..72a0a51 100644 --- a/obd/protocol.py +++ b/obd/protocol.py @@ -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}")