Files
elmer/api/db.py
T

342 lines
15 KiB
Python

"""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), 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 (
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,
device_uuid TEXT,
phone_lang TEXT,
phone_tz TEXT,
phone_display 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,
-- Идемпотентность
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_request_id ON sessions(request_id);
CREATE INDEX IF NOT EXISTS idx_sessions_uuid ON sessions(device_uuid);
""")
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,
request_id: str = "", response_json: dict | None = None):
"""Сохраняет сводную запись о сессии.
Если request_id передан и уже существует — silently return (идемпотентность).
"""
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
resp_json_str = json.dumps(response_json, ensure_ascii=False) if response_json else None
self.conn.execute("""
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, device_uuid, phone_lang, phone_tz, phone_display,
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, request_id, response_json
) 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("device_uuid"),
ci.get("phone_lang"), ci.get("phone_tz"), ci.get("phone_display"),
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,
request_id if request_id else None,
resp_json_str,
))
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]
def save_dtc_scan(self, client_info: dict, dtc_codes: list[str]):
"""Сохраняет быстрый скан кодов ошибок."""
self.conn.execute("""
INSERT INTO sessions (
client_ip, real_ip, user_agent,
phone_model, phone_maker, android_version, android_sdk,
app_version, android_id, device_uuid,
elm_mac, elm_bt_name,
dtc_count, response_count,
script_mode, transport,
raw_responses
) VALUES (?,?,?, ?,?,?,?, ?,?,?, ?,?, ?,?, ?,?,?)
""", (
client_info.get("client_ip"), client_info.get("real_ip"), client_info.get("user_agent"),
client_info.get("phone_model"), client_info.get("phone_maker"), client_info.get("android_version"),
client_info.get("android_sdk"), client_info.get("app_version"), client_info.get("android_id"),
client_info.get("device_uuid"),
client_info.get("elm_mac"), client_info.get("elm_bt_name"),
len(dtc_codes), 0,
"dtc_scan", client_info.get("transport", "bt"),
json.dumps([{"decoded": f"DTC stored: {c}"} for c in dtc_codes], ensure_ascii=False)
))
self.conn.commit()
# ── 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]