Files
GarminHealthLab/backend/db.py
ericwyuan acc6a2474b [阶段4.4] AI 建议结果缓存 - 页面不再阻塞等待 160 秒
网关首选的推理模型一次生成约 160 秒,每次打开建议页都重跑不可用。
结果落库缓存,页面读缓存,用户想要新的再手动触发。

db.py:
- 新增 ai_recommendations 表,每用户一行(重新生成是替换不是累积)
- fingerprint 列记录这条建议是基于哪份数据算出来的

services/analysis.py:
- _fingerprint() 对全部每日指标 + 运动条数取 sha256,任何一次同步
  新增或修正了数值都会让摘要变化,从而使缓存失效
- TTL 默认 24 小时(AI_CACHE_TTL_HOURS 可调)
- 指定 model 参数时绕过缓存:点名某个模型意味着想要那个模型的答案
- 降级到规则引擎的结果不写缓存,避免把兜底答案当成 AI 结果存下来
- 缓存写入失败只打日志,不影响本次请求返回

routes: ?refresh=1 强制重新生成

前端:
- "重新生成" 按钮走 refresh,并提示需要 1-3 分钟、可以离开本页
- meta 栏显示是否为缓存结果及生成时间,以及网关的上游厂商
- axios 该请求超时放宽到 240s(冷生成远超默认超时)

tests/test_ai_cache.py (20 通过):
- 第二次调用不再打模型
- 新增一天数据 / 修正某天数值 / 新增一条运动记录,三种情况都失效
- TTL 边界两侧各一条(刚过期重算、未过期沿用)
- 缓存按用户隔离,A 的结果不会答给 B
- payload 损坏时重新生成而不是抛异常
- 规则兜底结果和无数据用户都不落缓存

Co-Authored-By: Claude Haiku 4.5 <noreply@anthropic.com>
2026-08-23 17:55:12 +08:00

244 lines
6.3 KiB
Python

"""
Pluggable data layer for Garmin Health Lab.
Supports both SQLite (stdlib, local dev) and MariaDB (PyMySQL, NAS production)
through a single unified API:
init_db() -> create tables if missing
execute(sql, params) -> INSERT/UPDATE/DELETE, returns {id, changes}
query_one(sql, params) -> one row as dict or None
query_all(sql, params) -> list of row dicts
Both backends accept `?` placeholders; the SQL is translated to `%s` for
MariaDB automatically. Upserts must use backend-specific SQL (see services).
"""
import os
import sqlite3
import threading
import queue
import datetime
from config import (
DB_TYPE,
SQLITE_PATH,
MARIADB_SOCKET,
MARIADB_HOST,
MARIADB_PORT,
MARIADB_USER,
MARIADB_PASSWORD,
MARIADB_DATABASE,
)
SCHEMA = """
CREATE TABLE IF NOT EXISTS users (
id VARCHAR(64) PRIMARY KEY,
email VARCHAR(255) NOT NULL UNIQUE,
garmin_email VARCHAR(255) NOT NULL,
garmin_password_hash TEXT NOT NULL,
jwt_token TEXT,
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP
);
CREATE TABLE IF NOT EXISTS health_data (
id VARCHAR(64) PRIMARY KEY,
user_id VARCHAR(64) NOT NULL,
date DATE NOT NULL,
steps INT,
heart_rate INT,
heart_rate_variability DOUBLE,
blood_pressure_systolic INT,
blood_pressure_diastolic INT,
sleep_duration INT,
sleep_quality DOUBLE,
stress INT,
calories_burned DOUBLE,
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP,
UNIQUE(user_id, date),
FOREIGN KEY (user_id) REFERENCES users(id)
);
CREATE TABLE IF NOT EXISTS activities (
id VARCHAR(64) PRIMARY KEY,
user_id VARCHAR(64) NOT NULL,
activity_type VARCHAR(255) NOT NULL,
start_time DATETIME NOT NULL,
end_time DATETIME NOT NULL,
duration INT,
distance DOUBLE,
calories DOUBLE,
heart_rate_average INT,
heart_rate_max INT,
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
FOREIGN KEY (user_id) REFERENCES users(id)
);
CREATE TABLE IF NOT EXISTS sync_status (
user_id VARCHAR(64) PRIMARY KEY,
last_sync_time DATETIME,
status VARCHAR(32) DEFAULT 'idle',
last_error TEXT,
records_synced INT DEFAULT 0,
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP,
FOREIGN KEY (user_id) REFERENCES users(id)
);
-- One cached LLM answer per user. Generating one takes minutes against a
-- large reasoning model, which is far too slow to sit in a page load, so the
-- result is stored and reused until the underlying data changes.
-- `fingerprint` identifies the health data the advice was derived from.
CREATE TABLE IF NOT EXISTS ai_recommendations (
user_id VARCHAR(64) PRIMARY KEY,
fingerprint VARCHAR(64) NOT NULL,
model VARCHAR(64),
upstream VARCHAR(64),
days INT,
payload TEXT NOT NULL,
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
FOREIGN KEY (user_id) REFERENCES users(id)
);
"""
# --- MariaDB pool (lazy) ----------------------------------------------------
_mariadb_pool = None
_pool_lock = threading.Lock()
def _new_mariadb_conn():
import pymysql
from pymysql.cursors import DictCursor
kwargs = dict(
user=MARIADB_USER,
password=MARIADB_PASSWORD,
database=MARIADB_DATABASE,
charset="utf8mb4",
autocommit=True,
cursorclass=DictCursor,
connect_timeout=10,
)
if MARIADB_SOCKET:
kwargs["unix_socket"] = MARIADB_SOCKET
else:
kwargs["host"] = MARIADB_HOST
kwargs["port"] = MARIADB_PORT
return pymysql.connect(**kwargs)
def _mariadb_acquire():
global _mariadb_pool
if _mariadb_pool is None:
with _pool_lock:
if _mariadb_pool is None:
_mariadb_pool = queue.Queue(maxsize=10)
for _ in range(10):
_mariadb_pool.put(_new_mariadb_conn())
try:
return _mariadb_pool.get(block=False)
except queue.Empty:
return _new_mariadb_conn()
def _mariadb_release(conn):
try:
conn.ping(reconnect=False)
_mariadb_pool.put(conn)
except Exception:
try:
conn.close()
except Exception:
pass
# --- SQLite connection ------------------------------------------------------
def _sqlite_connect():
data_dir = os.path.dirname(SQLITE_PATH)
if data_dir and not os.path.exists(data_dir):
os.makedirs(data_dir, exist_ok=True)
conn = sqlite3.connect(SQLITE_PATH, isolation_level=None)
conn.row_factory = sqlite3.Row
conn.execute("PRAGMA foreign_keys = ON")
return conn
def _connect():
if DB_TYPE == "mariadb":
return _mariadb_acquire()
return _sqlite_connect()
def _disconnect(conn):
if DB_TYPE == "mariadb":
_mariadb_release(conn)
else:
conn.close()
def _adapt_sql(sql):
# pymysql uses %s placeholders; sqlite3 uses ?. Business code writes ?.
return sql.replace("?", "%s") if DB_TYPE == "mariadb" else sql
def _serialize(value):
if value is None:
return None
if isinstance(value, (datetime.datetime, datetime.date)):
return value.isoformat()
return value
def _row_to_dict(row):
if row is None:
return None
if isinstance(row, dict):
return {k: _serialize(v) for k, v in row.items()}
return {k: _serialize(row[k]) for k in row.keys()}
# --- Public API -------------------------------------------------------------
def init_db():
conn = _connect()
try:
cur = conn.cursor()
for stmt in SCHEMA.split(";"):
stmt = stmt.strip()
if not stmt:
continue
cur.execute(_adapt_sql(stmt))
finally:
_disconnect(conn)
def execute(sql, params=None):
params = params or []
conn = _connect()
try:
cur = conn.cursor()
cur.execute(_adapt_sql(sql), params)
return {"id": cur.lastrowid, "changes": cur.rowcount}
finally:
_disconnect(conn)
def query_one(sql, params=None):
params = params or []
conn = _connect()
try:
cur = conn.cursor()
cur.execute(_adapt_sql(sql), params)
return _row_to_dict(cur.fetchone())
finally:
_disconnect(conn)
def query_all(sql, params=None):
params = params or []
conn = _connect()
try:
cur = conn.cursor()
cur.execute(_adapt_sql(sql), params)
return [_row_to_dict(r) for r in cur.fetchall()]
finally:
_disconnect(conn)