fix: threading lock in Database for concurrent writes
This commit is contained in:
@@ -10,6 +10,7 @@
|
|||||||
|
|
||||||
import json
|
import json
|
||||||
import sqlite3
|
import sqlite3
|
||||||
|
import threading
|
||||||
from datetime import datetime, timezone
|
from datetime import datetime, timezone
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|
||||||
@@ -21,6 +22,7 @@ class Database:
|
|||||||
self.conn.row_factory = sqlite3.Row
|
self.conn.row_factory = sqlite3.Row
|
||||||
self.conn.execute("PRAGMA journal_mode=WAL")
|
self.conn.execute("PRAGMA journal_mode=WAL")
|
||||||
self.conn.execute("PRAGMA busy_timeout=30000")
|
self.conn.execute("PRAGMA busy_timeout=30000")
|
||||||
|
self._lock = threading.Lock()
|
||||||
self._init_schema()
|
self._init_schema()
|
||||||
|
|
||||||
def __enter__(self):
|
def __enter__(self):
|
||||||
@@ -175,65 +177,65 @@ class Database:
|
|||||||
|
|
||||||
Если request_id передан и уже существует — silently return (идемпотентность).
|
Если request_id передан и уже существует — silently return (идемпотентность).
|
||||||
"""
|
"""
|
||||||
ci = client_info
|
with self._lock:
|
||||||
|
ci = client_info
|
||||||
|
|
||||||
# Подсчёт DTC/PID из ответов
|
# Подсчёт DTC/PID из ответов
|
||||||
dtc_count = 0
|
dtc_count = 0
|
||||||
pid_count = 0
|
pid_count = 0
|
||||||
for r in responses:
|
for r in responses:
|
||||||
dec = (r.get("decoded") or "").lower()
|
dec = (r.get("decoded") or "").lower()
|
||||||
if dec.startswith("dtc"):
|
if dec.startswith("dtc"):
|
||||||
dtc_count += 1
|
dtc_count += 1
|
||||||
elif ":" in dec and not dec.startswith(("vin", "dtc", "elm", "protocol")):
|
elif ":" in dec and not dec.startswith(("vin", "dtc", "elm", "protocol")):
|
||||||
pid_count += 1
|
pid_count += 1
|
||||||
|
|
||||||
# VIN из ответов
|
# VIN из ответов
|
||||||
vin = None
|
vin = None
|
||||||
for r in responses:
|
for r in responses:
|
||||||
dec = (r.get("decoded") or "")
|
dec = (r.get("decoded") or "")
|
||||||
if dec.startswith("VIN:"):
|
if dec.startswith("VIN:"):
|
||||||
vin = dec[4:].strip()
|
vin = dec[4:].strip()
|
||||||
if len(vin) != 17:
|
if len(vin) != 17:
|
||||||
vin = None
|
vin = None
|
||||||
break
|
break
|
||||||
|
|
||||||
resp_json_str = json.dumps(response_json, ensure_ascii=False) if response_json else None
|
resp_json_str = json.dumps(response_json, ensure_ascii=False) if response_json else None
|
||||||
|
|
||||||
self.conn.execute("""
|
self.conn.execute("""
|
||||||
INSERT OR IGNORE INTO sessions (
|
INSERT OR IGNORE INTO sessions (
|
||||||
client_ip, real_ip, user_agent, content_length,
|
client_ip, real_ip, user_agent, content_length,
|
||||||
phone_model, phone_maker, android_version, android_sdk,
|
phone_model, phone_maker, android_version, android_sdk,
|
||||||
app_version, android_id, device_uuid, phone_lang, phone_tz, phone_display,
|
app_version, android_id, device_uuid, phone_lang, phone_tz, phone_display,
|
||||||
elm_mac, elm_bt_name, obd_protocol,
|
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, car_info,
|
||||||
|
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,
|
vin, dtc_count, pid_count,
|
||||||
duration_ms, response_count, error_count,
|
ci.get("duration_ms"), len(responses), ci.get("error_count", 0),
|
||||||
retry_count, timeout_count, script_mode,
|
ci.get("retry_count", 0), ci.get("timeout_count", 0),
|
||||||
transport, mock_mode, car_info,
|
ci.get("script_mode"), ci.get("transport"), ci.get("mock_mode", 0),
|
||||||
diagnosis_text, diagnosis_len, llm_model,
|
ci.get("car_info", ""),
|
||||||
llm_duration_ms, llm_success,
|
diagnosis, len(diagnosis), llm_model,
|
||||||
raw_responses, request_id, response_json
|
llm_duration_ms, 1 if llm_success else 0,
|
||||||
) VALUES (?,?,?,?, ?,?,?,?,?, ?,?,?,?,?, ?,?,?, ?,?,?, ?,?,?, ?,?,?, ?,?, ?,?,?,?, ?,?,?,?,?)
|
json.dumps(responses, ensure_ascii=False) if responses else None,
|
||||||
|
request_id if request_id else None,
|
||||||
""", (
|
resp_json_str,
|
||||||
ci.get("client_ip"), ci.get("real_ip"), ci.get("user_agent"),
|
))
|
||||||
ci.get("content_length"),
|
self.conn.commit()
|
||||||
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),
|
|
||||||
ci.get("car_info", ""),
|
|
||||||
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]:
|
def get_recent_sessions(self, limit: int = 50) -> list[dict]:
|
||||||
"""Последние N сессий."""
|
"""Последние N сессий."""
|
||||||
@@ -244,27 +246,28 @@ class Database:
|
|||||||
|
|
||||||
def save_dtc_scan(self, client_info: dict, dtc_codes: list[str]):
|
def save_dtc_scan(self, client_info: dict, dtc_codes: list[str]):
|
||||||
"""Сохраняет быстрый скан кодов ошибок."""
|
"""Сохраняет быстрый скан кодов ошибок."""
|
||||||
self.conn.execute("""
|
with self._lock:
|
||||||
INSERT INTO sessions (
|
self.conn.execute("""
|
||||||
client_ip, real_ip, user_agent,
|
INSERT INTO sessions (
|
||||||
phone_model, phone_maker, android_version, android_sdk,
|
client_ip, real_ip, user_agent,
|
||||||
app_version, android_id, device_uuid,
|
phone_model, phone_maker, android_version, android_sdk,
|
||||||
elm_mac, elm_bt_name,
|
app_version, android_id, device_uuid,
|
||||||
dtc_count, response_count,
|
elm_mac, elm_bt_name,
|
||||||
script_mode, transport,
|
dtc_count, response_count,
|
||||||
raw_responses
|
script_mode, transport,
|
||||||
) VALUES (?,?,?, ?,?,?,?, ?,?,?, ?,?, ?,?, ?,?,?)
|
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("client_ip"), client_info.get("real_ip"), client_info.get("user_agent"),
|
||||||
client_info.get("android_sdk"), client_info.get("app_version"), client_info.get("android_id"),
|
client_info.get("phone_model"), client_info.get("phone_maker"), client_info.get("android_version"),
|
||||||
client_info.get("device_uuid"),
|
client_info.get("android_sdk"), client_info.get("app_version"), client_info.get("android_id"),
|
||||||
client_info.get("elm_mac"), client_info.get("elm_bt_name"),
|
client_info.get("device_uuid"),
|
||||||
len(dtc_codes), 0,
|
client_info.get("elm_mac"), client_info.get("elm_bt_name"),
|
||||||
"dtc_scan", client_info.get("transport", "bt"),
|
len(dtc_codes), 0,
|
||||||
json.dumps([{"decoded": f"DTC stored: {c}"} for c in dtc_codes], ensure_ascii=False)
|
"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()
|
))
|
||||||
|
self.conn.commit()
|
||||||
|
|
||||||
# ── device_profiles ──────────────────────────────────
|
# ── device_profiles ──────────────────────────────────
|
||||||
|
|
||||||
@@ -285,33 +288,34 @@ class Database:
|
|||||||
|
|
||||||
profile — результат obd.probe.probe().
|
profile — результат obd.probe.probe().
|
||||||
"""
|
"""
|
||||||
now = datetime.now(timezone.utc).isoformat()
|
with self._lock:
|
||||||
self.conn.execute("""
|
now = datetime.now(timezone.utc).isoformat()
|
||||||
INSERT INTO device_profiles
|
self.conn.execute("""
|
||||||
(mac, level, elm_version, elm_desc, protocol, voltage,
|
INSERT INTO device_profiles
|
||||||
supported, unsupported, errors, first_seen, last_seen)
|
(mac, level, elm_version, elm_desc, protocol, voltage,
|
||||||
VALUES (?,?,?,?,?,?, ?,?,?, ?,?)
|
supported, unsupported, errors, first_seen, last_seen)
|
||||||
ON CONFLICT(mac) DO UPDATE SET
|
VALUES (?,?,?,?,?,?, ?,?,?, ?,?)
|
||||||
level = excluded.level,
|
ON CONFLICT(mac) DO UPDATE SET
|
||||||
elm_version = excluded.elm_version,
|
level = excluded.level,
|
||||||
elm_desc = excluded.elm_desc,
|
elm_version = excluded.elm_version,
|
||||||
protocol = excluded.protocol,
|
elm_desc = excluded.elm_desc,
|
||||||
voltage = excluded.voltage,
|
protocol = excluded.protocol,
|
||||||
supported = excluded.supported,
|
voltage = excluded.voltage,
|
||||||
unsupported = excluded.unsupported,
|
supported = excluded.supported,
|
||||||
errors = excluded.errors,
|
unsupported = excluded.unsupported,
|
||||||
last_seen = excluded.last_seen
|
errors = excluded.errors,
|
||||||
""", (
|
last_seen = excluded.last_seen
|
||||||
mac,
|
""", (
|
||||||
profile.get("level", -1),
|
mac,
|
||||||
profile.get("elm_version"),
|
profile.get("level", -1),
|
||||||
profile.get("elm_desc"),
|
profile.get("elm_version"),
|
||||||
profile.get("protocol"),
|
profile.get("elm_desc"),
|
||||||
profile.get("voltage"),
|
profile.get("protocol"),
|
||||||
json.dumps(profile.get("supported", []), ensure_ascii=False),
|
profile.get("voltage"),
|
||||||
json.dumps(profile.get("unsupported", []), ensure_ascii=False),
|
json.dumps(profile.get("supported", []), ensure_ascii=False),
|
||||||
json.dumps(profile.get("errors", []), ensure_ascii=False),
|
json.dumps(profile.get("unsupported", []), ensure_ascii=False),
|
||||||
now, now,
|
json.dumps(profile.get("errors", []), ensure_ascii=False),
|
||||||
))
|
now, now,
|
||||||
self.conn.commit()
|
))
|
||||||
|
self.conn.commit()
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user