Files
GarminHealthLab/backend/db.py
ericwyuan e7c336b9df fix(db): 数据库晚一步启动就会让服务永久挂掉
今天 11:33 站点整个不可用。不是代码的问题,是启动顺序:
NAS 上 MariaDB 在 11:34 才起来,而应用 11:33 就尝试连接,
init_db() 抛异常 → gunicorn 报 Worker failed to boot → master 退出。
一分钟后数据库好了,但已经没有进程在跑,没人会去重试——
站点就一直躺到有人手工重启为止。

init_db 改为在 DB_INIT_RETRY_SECONDS(默认 120 秒)内重试等待数据库,
超时仍然如实抛错,不会假装启动成功。NAS 重启时应用和数据库一起起来,
这个竞争是常态而不是意外。

Co-Authored-By: Claude Haiku 4.5 <noreply@anthropic.com>
2026-08-25 14:23:07 +08:00

584 lines
18 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 time
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)
);
-- Garmin OAuth tokens, obtained once through an interactive login.
-- Garmin accounts with two-factor auth cannot be logged into unattended: the
-- library asks for an MFA code on stdin, which a gunicorn worker does not
-- have. Storing the resulting tokens lets every later sync skip the login
-- entirely (they stay valid for roughly a year).
CREATE TABLE IF NOT EXISTS garmin_tokens (
user_id VARCHAR(64) PRIMARY KEY,
token TEXT NOT NULL,
garmin_email VARCHAR(255),
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP,
FOREIGN KEY (user_id) REFERENCES users(id)
);
-- Badges earned on Garmin Connect ("奖励"). Keyed by Garmin's own badge id so
-- a re-sync updates rather than duplicates.
CREATE TABLE IF NOT EXISTS badges (
id VARCHAR(64) NOT NULL,
user_id VARCHAR(64) NOT NULL,
badge_key VARCHAR(128),
name VARCHAR(255),
category_id INT,
difficulty_id INT,
earned_date DATETIME,
earned_count INT,
points INT,
PRIMARY KEY (user_id, id),
FOREIGN KEY (user_id) REFERENCES users(id)
);
-- Personal records (个人纪录), e.g. fastest 5k, longest run.
CREATE TABLE IF NOT EXISTS personal_records (
id VARCHAR(64) NOT NULL,
user_id VARCHAR(64) NOT NULL,
type_id INT,
activity_id VARCHAR(64),
activity_name VARCHAR(255),
activity_type VARCHAR(64),
value DOUBLE,
achieved_at DATETIME,
PRIMARY KEY (user_id, id),
FOREIGN KEY (user_id) REFERENCES users(id)
);
-- Coordination for work that must happen once per interval regardless of how
-- many gunicorn workers are running. A worker claims a job by writing its
-- row, and the others see a fresh claim and stand down. Without this the
-- hourly sync would fire once per worker.
CREATE TABLE IF NOT EXISTS job_locks (
name VARCHAR(64) PRIMARY KEY,
holder VARCHAR(64),
claimed_at DATETIME,
last_run_at DATETIME
);
-- Rendezvous for the interactive MFA login.
-- garth asks for the code through a *blocking* callback, so the login parks in
-- a background thread while the code arrives in a separate HTTP request that
-- may land on a different gunicorn worker. The handoff therefore goes through
-- the database rather than process memory.
-- Holds no password: that stays in the waiting thread's memory only.
CREATE TABLE IF NOT EXISTS garmin_mfa_sessions (
id VARCHAR(64) PRIMARY KEY,
user_id VARCHAR(64) NOT NULL,
status VARCHAR(32) NOT NULL,
code VARCHAR(16),
error TEXT,
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP,
FOREIGN KEY (user_id) REFERENCES users(id)
);
-- Per-user profile and preferences.
-- Height/weight/birth date/sex are here rather than on `users` because they
-- are body measurements the owner edits over time, not identity; and because
-- the rating bands and the fitness-age estimate need them, an account without
-- them still works, just with fewer personalised readings.
CREATE TABLE IF NOT EXISTS user_settings (
user_id VARCHAR(64) PRIMARY KEY,
height_cm DOUBLE,
weight_kg DOUBLE,
birth_date DATE,
sex VARCHAR(16),
units VARCHAR(16),
auto_sync INT,
auto_sync_minutes INT,
history_days INT,
updated_at DATETIME,
FOREIGN KEY (user_id) REFERENCES users(id)
);
-- Full detail for one activity, exactly as Garmin returned it.
-- The list view stores only the summary columns; opening an activity needs
-- laps, heart-rate zones and the sampled series, which are far too large to
-- carry on every list request. Fetched on demand and kept, so the second
-- visit costs nothing and works offline.
CREATE TABLE IF NOT EXISTS activity_details (
activity_id VARCHAR(64) PRIMARY KEY,
user_id VARCHAR(64) NOT NULL,
payload MEDIUMTEXT,
fetched_at DATETIME,
FOREIGN KEY (user_id) REFERENCES users(id)
);
-- Weight and body composition, one row per measurement day.
-- Separate from health_data because it arrives from the scale rather than the
-- watch, on its own irregular schedule — most days simply have no row.
CREATE TABLE IF NOT EXISTS body_composition (
id VARCHAR(96) PRIMARY KEY,
user_id VARCHAR(64) NOT NULL,
date DATE NOT NULL,
weight_kg DOUBLE,
bmi DOUBLE,
body_fat_pct DOUBLE,
body_water_pct DOUBLE,
bone_mass_kg DOUBLE,
muscle_mass_kg DOUBLE,
physique_rating DOUBLE,
visceral_fat DOUBLE,
metabolic_age DOUBLE,
source VARCHAR(32),
UNIQUE(user_id, date),
FOREIGN KEY (user_id) REFERENCES users(id)
);
-- Blood pressure readings. Manually entered in Garmin Connect, so there may
-- be none at all; the table exists so that there is somewhere to put them.
CREATE TABLE IF NOT EXISTS blood_pressure (
id VARCHAR(96) PRIMARY KEY,
user_id VARCHAR(64) NOT NULL,
measured_at DATETIME NOT NULL,
systolic INT,
diastolic INT,
pulse INT,
note TEXT,
UNIQUE(user_id, measured_at),
FOREIGN KEY (user_id) REFERENCES users(id)
);
-- Garmin's predicted race times, in seconds. One row per day it recalculates.
CREATE TABLE IF NOT EXISTS race_predictions (
id VARCHAR(96) PRIMARY KEY,
user_id VARCHAR(64) NOT NULL,
date DATE NOT NULL,
time_5k INT,
time_10k INT,
time_half INT,
time_marathon INT,
UNIQUE(user_id, date),
FOREIGN KEY (user_id) REFERENCES users(id)
);
-- Within-day sample series: heart rate, stress, body battery, respiration,
-- SpO2. One generic table rather than five near-identical ones — they differ
-- only in what the numbers mean, and the daily screen reads them the same way.
CREATE TABLE IF NOT EXISTS daily_series (
id VARCHAR(96) PRIMARY KEY,
user_id VARCHAR(64) NOT NULL,
date DATE NOT NULL,
kind VARCHAR(32) NOT NULL,
payload MEDIUMTEXT,
fetched_at DATETIME,
UNIQUE(user_id, date, kind),
FOREIGN KEY (user_id) REFERENCES users(id)
);
-- Badge challenges and ad-hoc challenges. Distinct from `badges`: a badge is
-- earned once, a challenge has a period, a target and a standing.
CREATE TABLE IF NOT EXISTS challenges (
id VARCHAR(96) PRIMARY KEY,
user_id VARCHAR(64) NOT NULL,
challenge_uuid VARCHAR(96),
kind VARCHAR(32),
name VARCHAR(255),
status VARCHAR(64),
start_date DATE,
end_date DATE,
payload MEDIUMTEXT,
FOREIGN KEY (user_id) REFERENCES users(id)
);
-- Paired devices, so the app can say which watch a number came from.
CREATE TABLE IF NOT EXISTS devices (
id VARCHAR(96) PRIMARY KEY,
user_id VARCHAR(64) NOT NULL,
device_id VARCHAR(96),
name VARCHAR(255),
model VARCHAR(255),
serial VARCHAR(96),
software_version VARCHAR(64),
last_used_at DATETIME,
payload MEDIUMTEXT,
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):
"""Return a connection to the pool, or close it if the pool is full.
`put()` blocks when the queue is at maxsize — and the queue can be full,
because `_mariadb_acquire` opens an extra connection whenever the pool is
empty rather than waiting. Under enough concurrency (a long sync plus
ordinary requests) a thread would park here forever holding its request
open. `put_nowait` plus closing the surplus keeps the pool bounded and the
thread free.
"""
try:
conn.ping(reconnect=False)
except Exception:
try:
conn.close()
except Exception:
pass
return
try:
_mariadb_pool.put_nowait(conn)
except queue.Full:
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()}
# Columns added after the first release. `CREATE TABLE IF NOT EXISTS` does
# nothing to a table that already exists, so new metrics need an explicit
# additive migration or they silently never appear in production.
MIGRATIONS = {
"sync_status": [
# A full backfill runs for many minutes, so the UI needs to show how
# far along it is rather than an indefinite spinner.
("progress_current", "INT"),
("progress_total", "INT"),
("started_at", "DATETIME"),
# Which part of the sync is running. "0 / 730 天" says nothing about
# what is actually happening for the several minutes of it.
("stage", "VARCHAR(64)"),
],
"health_data": [
# activity / energy
("distance_meters", "DOUBLE"),
("active_calories", "DOUBLE"),
("bmr_calories", "DOUBLE"),
("floors_ascended", "DOUBLE"),
("floors_descended", "DOUBLE"),
("intensity_minutes", "INT"),
("step_goal", "INT"),
("sedentary_seconds", "INT"),
("active_seconds", "INT"),
# heart / stress
("heart_rate_max", "INT"),
("heart_rate_min", "INT"),
("stress_max", "INT"),
# body battery
("body_battery_high", "INT"),
("body_battery_low", "INT"),
("body_battery_charged", "INT"),
("body_battery_drained", "INT"),
# breathing / blood oxygen
("spo2_avg", "DOUBLE"),
("spo2_min", "INT"),
("respiration_avg", "DOUBLE"),
("respiration_min", "DOUBLE"),
("respiration_max", "DOUBLE"),
# sleep detail
("sleep_deep_seconds", "INT"),
("sleep_light_seconds", "INT"),
("sleep_rem_seconds", "INT"),
("sleep_awake_seconds", "INT"),
("sleep_spo2_avg", "DOUBLE"),
("sleep_respiration_avg", "DOUBLE"),
("sleep_stress_avg", "DOUBLE"),
# training
("training_readiness", "INT"),
("vo2max", "DOUBLE"),
("endurance_score", "INT"),
# hill score, hydration and weight — daily scalars that were being
# fetched from Garmin by nothing at all until now
("hill_score", "INT"),
("hydration_ml", "INT"),
("hydration_goal_ml", "INT"),
("sweat_loss_ml", "INT"),
("weight_kg", "DOUBLE"),
("body_fat_pct", "DOUBLE"),
("bmi", "DOUBLE"),
],
}
def _existing_columns(cur, table):
if DB_TYPE == "mariadb":
cur.execute(
"SELECT COLUMN_NAME FROM information_schema.COLUMNS "
"WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = %s",
[table],
)
return {r["COLUMN_NAME"] if isinstance(r, dict) else r[0] for r in cur.fetchall()}
cur.execute(f"PRAGMA table_info({table})")
return {row[1] for row in cur.fetchall()}
def _migrate(cur):
for table, columns in MIGRATIONS.items():
present = _existing_columns(cur, table)
for name, coltype in columns:
if name in present:
continue
# SQLite has no "ADD COLUMN IF NOT EXISTS"; the membership check
# above is what keeps this idempotent on both backends.
cur.execute(f"ALTER TABLE {table} ADD COLUMN {name} {coltype}")
# --- Public API -------------------------------------------------------------
def _statements(schema):
"""Split a schema script into statements.
Comments are stripped first: splitting the raw text on ';' would cut a
comment that happens to contain one in half and hand the remainder to the
database as SQL.
"""
lines = [ln for ln in schema.splitlines() if not ln.strip().startswith("--")]
for stmt in "\n".join(lines).split(";"):
stmt = stmt.strip()
if stmt:
yield stmt
# How long to keep waiting for the database at startup.
#
# On a NAS reboot the app and MariaDB come up together and the app usually
# wins the race. Without this it raised, gunicorn reported "Worker failed to
# boot", the master shut down — and when MariaDB appeared seconds later there
# was nothing left running to notice. The site stayed down until someone
# restarted it by hand.
INIT_RETRY_SECONDS = int(os.environ.get("DB_INIT_RETRY_SECONDS") or 120)
INIT_RETRY_INTERVAL = 3
def init_db():
deadline = time.monotonic() + INIT_RETRY_SECONDS
attempt = 0
while True:
attempt += 1
try:
conn = _connect()
break
except Exception as e: # noqa: BLE001 - any connection failure is worth retrying
if time.monotonic() >= deadline:
raise
if attempt == 1:
print(f"[db] 数据库还没准备好,重试中:{e}")
time.sleep(INIT_RETRY_INTERVAL)
if attempt > 1:
print(f"[db] 第 {attempt} 次尝试后连上数据库")
try:
cur = conn.cursor()
for stmt in _statements(SCHEMA):
cur.execute(_adapt_sql(stmt))
_migrate(cur)
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)