原来每天只存 7 个指标,而 get_user_summary 一次就返回 60+ 字段, 另有睡眠分期、训练准备度、耐力分等独立端点从未被调用。 db.py: - health_data 新增 31 列(距离/活动卡路里/基础代谢/爬楼/强度分钟/ 久坐时长/最高最低心率/最大压力/身体电量四项/血氧/呼吸/ 睡眠深浅REM清醒分期/睡眠血氧/睡眠呼吸/睡眠压力/训练准备度/ VO2max/耐力分) - 新增 badges 与 personal_records 两张表,均以 (user_id, garmin_id) 为主键,重复同步更新而非累积 - 新增增量迁移: CREATE TABLE IF NOT EXISTS 对已存在的表不生效, 新列必须显式 ALTER,否则生产库上永远不会出现。按列名比对后 逐个补齐,SQLite 与 MariaDB 都幂等 services/garmin.py: - _extract_daily 改为汇总 user_summary + sleep + hrv + training_readiness + training_status + endurance_score 五个端点 - 每个可选端点用 _safe 包裹:某项设备不记录时留 NULL,不影响当天其余数据 - 新增 sync_badges / sync_personal_records(账号级,每次同步取一次) fix(garmin): 个人纪录整批写入失败 - Garmin 在同一份数据里混用 ISO 字符串和 Unix 毫秒时间戳, prStartTimeGmt 是 1570961412000,写进 DATETIME 列被 MariaDB 以 1292 拒绝,导致 11 项个人纪录一条都没存进去 - 新增 _to_datetime 统一处理 ISO / 毫秒 / 秒三种形状,并优先取 Garmin 自己提供的 *Formatted 字段 services/ai.py: - 送给模型的 CSV 从 7 列扩到 23 列,纳入身体电量、血氧、呼吸、 训练准备度、耐力分和睡眠分期 接口: GET /api/health/badges、/api/health/personal-records tests (+13, 共 292): - 徽章/纪录的往返、重复同步不累积、按用户隔离 - 两个用户可持有同一个 Garmin 徽章 id 而不冲突 - 时间戳三种形状的归一化及无效值不抛异常 NAS 实测: 7 天数据每天 31 项指标、65 个奖励、11 项个人纪录 Co-Authored-By: Claude Haiku 4.5 <noreply@anthropic.com>
502 lines
19 KiB
Python
502 lines
19 KiB
Python
"""
|
||
Garmin sync service.
|
||
|
||
Pulls daily summaries + activities through the `garminconnect` library and
|
||
upserts them. The library and real Garmin credentials are required to actually
|
||
run a sync; without them the endpoint reports a clear error instead of
|
||
crashing.
|
||
|
||
Garmin credentials: the app only stores a scrypt/PBKDF2 *hash* of the Garmin
|
||
password (so it cannot be recovered), therefore a live sync needs the
|
||
plaintext garminEmail/garminPassword supplied in the request body.
|
||
|
||
On the library's API — these were verified against garminconnect 0.2.8:
|
||
* get_user_summary(cdate) -> one day of daily totals
|
||
* get_sleep_data(cdate) -> sleep, NOT part of the summary
|
||
* get_hrv_data(cdate) -> HRV, also separate
|
||
* get_activities_by_date(start, end) -> activities in a date range
|
||
* get_activities(start, limit) -> PAGINATION, not dates
|
||
The last two are easy to confuse: `get_activities` takes an offset and a count,
|
||
so passing it a date silently asks for activity number "2026-08-23".
|
||
"""
|
||
import datetime
|
||
import os
|
||
|
||
from db import execute, query_one
|
||
from config import DB_TYPE
|
||
from services import health
|
||
|
||
# How many days back a sync reaches.
|
||
DEFAULT_SYNC_DAYS = int(os.environ.get("GARMIN_SYNC_DAYS") or 7)
|
||
|
||
|
||
def _set_sync_status(user_id, status, now, **fields):
|
||
cols = ["user_id", "status", "last_sync_time"] + list(fields.keys())
|
||
placeholders = ", ".join(["?"] * len(cols))
|
||
if DB_TYPE == "mariadb":
|
||
updates = ", ".join(f"{c}=VALUES({c})" for c in cols if c != "user_id")
|
||
sql = (
|
||
f"INSERT INTO sync_status ({', '.join(cols)}) VALUES ({placeholders}) "
|
||
f"ON DUPLICATE KEY UPDATE {updates}"
|
||
)
|
||
else:
|
||
updates = ", ".join(f"{c}=excluded.{c}" for c in cols if c != "user_id")
|
||
sql = (
|
||
f"INSERT INTO sync_status ({', '.join(cols)}) VALUES ({placeholders}) "
|
||
f"ON CONFLICT(user_id) DO UPDATE SET {updates}"
|
||
)
|
||
execute(sql, [user_id, status, now] + list(fields.values()))
|
||
|
||
|
||
def get_sync_status(user_id):
|
||
row = query_one("SELECT * FROM sync_status WHERE user_id = ?", [user_id])
|
||
if not row:
|
||
return {
|
||
"status": "idle",
|
||
"lastSyncTime": None,
|
||
"recordsSynced": 0,
|
||
"lastError": None,
|
||
}
|
||
return {
|
||
"status": row["status"],
|
||
"lastSyncTime": row["last_sync_time"],
|
||
"recordsSynced": row["records_synced"],
|
||
"lastError": row["last_error"],
|
||
}
|
||
|
||
|
||
class MFARequired(RuntimeError):
|
||
"""Raised when a password login needs a code this process cannot obtain."""
|
||
|
||
|
||
# garth sends a browser User-Agent, which its SSO flow needs. The data API
|
||
# treats that same UA as a browser hitting it directly and answers every
|
||
# request with HTTP 200 and an empty array — no error, just no data. The
|
||
# official app's UA (and in fact any non-browser one) returns real data, so
|
||
# the header is swapped after login, before any API call.
|
||
API_USER_AGENT = "com.garmin.android.apps.connectmobile"
|
||
|
||
|
||
def _use_api_user_agent(client):
|
||
try:
|
||
client.garth.sess.headers["User-Agent"] = API_USER_AGENT
|
||
except AttributeError:
|
||
pass # a stubbed client in tests has no session
|
||
|
||
|
||
def _is_cn():
|
||
# Selects Garmin's China service, a separate backend with separate
|
||
# accounts. This project tracks an international account.
|
||
return (os.environ.get("GARMIN_IS_CN") or "").lower() in ("1", "true", "yes")
|
||
|
||
|
||
def _import_garmin():
|
||
try:
|
||
from garminconnect import Garmin
|
||
except ImportError:
|
||
raise RuntimeError(
|
||
"GARMIN_LIB_MISSING: 请先运行 `pip install garminconnect` 以启用同步"
|
||
)
|
||
return Garmin
|
||
|
||
|
||
def load_token(user_id):
|
||
row = query_one("SELECT token FROM garmin_tokens WHERE user_id = ?", [user_id])
|
||
return row["token"] if row else None
|
||
|
||
|
||
def save_token(user_id, token, garmin_email=None):
|
||
cols = ["user_id", "token", "garmin_email", "updated_at"]
|
||
placeholders = ", ".join(["?"] * len(cols))
|
||
if DB_TYPE == "mariadb":
|
||
updates = ", ".join(f"{c}=VALUES({c})" for c in cols if c != "user_id")
|
||
sql = (f"INSERT INTO garmin_tokens ({', '.join(cols)}) VALUES ({placeholders}) "
|
||
f"ON DUPLICATE KEY UPDATE {updates}")
|
||
else:
|
||
updates = ", ".join(f"{c}=excluded.{c}" for c in cols if c != "user_id")
|
||
sql = (f"INSERT INTO garmin_tokens ({', '.join(cols)}) VALUES ({placeholders}) "
|
||
f"ON CONFLICT(user_id) DO UPDATE SET {updates}")
|
||
execute(sql, [user_id, token, garmin_email,
|
||
datetime.datetime.utcnow().isoformat(timespec="seconds")])
|
||
|
||
|
||
def has_token(user_id):
|
||
return load_token(user_id) is not None
|
||
|
||
|
||
def _connect(creds, user_id=None):
|
||
"""Obtain a logged-in Garmin client.
|
||
|
||
Prefers stored OAuth tokens: an account with two-factor auth cannot be
|
||
logged into from a web worker, because the library asks for the code on
|
||
stdin and there is none (the failure surfaces as
|
||
"EOFError: EOF when reading a line"). Tokens are minted once by
|
||
`garmin_login.py`, which runs in a terminal where a code can be typed.
|
||
"""
|
||
Garmin = _import_garmin()
|
||
client = Garmin(is_cn=_is_cn())
|
||
|
||
token = load_token(user_id) if user_id else None
|
||
if token:
|
||
client.garth.loads(token)
|
||
_use_api_user_agent(client)
|
||
# Proves the token still works, and refreshes it if near expiry.
|
||
client.garth.refresh_oauth2()
|
||
# garminconnect builds most of its URLs from display_name, so leaving
|
||
# it unset sends every request to ".../None".
|
||
client.display_name = client.garth.profile["displayName"]
|
||
return client
|
||
|
||
if not creds.get("garminPassword"):
|
||
raise RuntimeError("缺少 Garmin 密码,且未找到已保存的登录令牌")
|
||
|
||
client.username = creds["garminEmail"]
|
||
client.password = creds["garminPassword"]
|
||
try:
|
||
client.login()
|
||
except EOFError as e:
|
||
# garth's default MFA prompt calls input(); under gunicorn stdin is
|
||
# closed, so it raises EOFError rather than anything descriptive.
|
||
raise MFARequired(
|
||
"该 Garmin 账号开启了两步验证。请在「数据同步」页面用密码重新绑定,"
|
||
"系统会提示你输入验证码。"
|
||
) from e
|
||
_use_api_user_agent(client)
|
||
return client
|
||
|
||
|
||
def describe(e):
|
||
"""A message that is never empty.
|
||
|
||
Some exceptions carry no text at all — a bare `assert` raises
|
||
AssertionError with str(e) == "" — and storing that produced a failed
|
||
sync whose recorded reason was blank, which is undiagnosable.
|
||
"""
|
||
text = str(e).strip()
|
||
return f"{type(e).__name__}: {text}" if text else type(e).__name__
|
||
|
||
|
||
def _num(*values):
|
||
"""First value that is a usable number."""
|
||
for v in values:
|
||
if isinstance(v, (int, float)) and not isinstance(v, bool):
|
||
return v
|
||
return None
|
||
|
||
|
||
def _safe(fn, default=None):
|
||
"""Call an optional endpoint; a metric the device does not record must not
|
||
abort the whole day."""
|
||
try:
|
||
return fn()
|
||
except Exception:
|
||
return default
|
||
|
||
|
||
def _first(seq):
|
||
return seq[0] if isinstance(seq, list) and seq else {}
|
||
|
||
|
||
def _to_datetime(*values):
|
||
"""Normalise Garmin's several timestamp shapes into an ISO string.
|
||
|
||
The same payload mixes ISO strings ("2019-10-13T10:10:12.0") with epoch
|
||
milliseconds (1570961412000); handing the latter to a DATETIME column is
|
||
rejected outright, so a personal record whose only timestamp was numeric
|
||
failed the whole batch.
|
||
"""
|
||
for v in values:
|
||
if v is None or v == "":
|
||
continue
|
||
if isinstance(v, str):
|
||
return v[:26]
|
||
if isinstance(v, (int, float)) and not isinstance(v, bool):
|
||
# Values past ~1e11 are milliseconds, below that seconds.
|
||
seconds = v / 1000 if v > 1e11 else v
|
||
try:
|
||
return datetime.datetime.utcfromtimestamp(seconds).isoformat(
|
||
timespec="seconds"
|
||
)
|
||
except (ValueError, OverflowError, OSError):
|
||
continue
|
||
return None
|
||
|
||
|
||
def _extract_daily(client, date_str):
|
||
"""Everything Garmin exposes for one day.
|
||
|
||
The daily summary is the bulk of it, but sleep, HRV, training readiness
|
||
and endurance each live behind their own endpoint — none of them appear in
|
||
get_user_summary. Each is fetched defensively so a metric this device does
|
||
not record leaves a NULL instead of failing the day.
|
||
"""
|
||
s = client.get_user_summary(date_str) or {}
|
||
|
||
sleep_dto = (_safe(lambda: client.get_sleep_data(date_str)) or {}).get(
|
||
"dailySleepDTO"
|
||
) or {}
|
||
scores = sleep_dto.get("sleepScores") if isinstance(
|
||
sleep_dto.get("sleepScores"), dict
|
||
) else {}
|
||
sleep_seconds = _num(sleep_dto.get("sleepTimeSeconds"))
|
||
|
||
hrv_summary = (_safe(lambda: client.get_hrv_data(date_str)) or {}).get(
|
||
"hrvSummary"
|
||
) or {}
|
||
|
||
readiness = _first(_safe(lambda: client.get_training_readiness(date_str), []))
|
||
training = _safe(lambda: client.get_training_status(date_str), {}) or {}
|
||
vo2 = (training.get("mostRecentVO2Max") or {}).get("generic") or {}
|
||
endurance = _safe(lambda: client.get_endurance_score(date_str), {}) or {}
|
||
|
||
def secs(key):
|
||
return _num(sleep_dto.get(key))
|
||
|
||
return {
|
||
"date": date_str,
|
||
# --- activity / energy ---
|
||
"steps": _num(s.get("totalSteps")),
|
||
"stepGoal": _num(s.get("dailyStepGoal")),
|
||
"distanceMeters": _num(s.get("totalDistanceMeters")),
|
||
"caloriesBurned": _num(s.get("totalKilocalories")),
|
||
"activeCalories": _num(s.get("activeKilocalories")),
|
||
"bmrCalories": _num(s.get("bmrKilocalories")),
|
||
"floorsAscended": _num(s.get("floorsAscended")),
|
||
"floorsDescended": _num(s.get("floorsDescended")),
|
||
"intensityMinutes": (
|
||
(_num(s.get("moderateIntensityMinutes")) or 0)
|
||
+ (_num(s.get("vigorousIntensityMinutes")) or 0)
|
||
) or None,
|
||
"sedentarySeconds": _num(s.get("sedentarySeconds")),
|
||
"activeSeconds": _num(s.get("activeSeconds")),
|
||
# --- heart / stress ---
|
||
"heartRate": _num(s.get("restingHeartRate"), s.get("averageHeartRate")),
|
||
"heartRateMax": _num(s.get("maxHeartRate")),
|
||
"heartRateMin": _num(s.get("minHeartRate")),
|
||
"heartRateVariability": _num(
|
||
hrv_summary.get("lastNightAvg"), hrv_summary.get("weeklyAvg")
|
||
),
|
||
"stress": _num(s.get("averageStressLevel")),
|
||
"stressMax": _num(s.get("maxStressLevel")),
|
||
# --- body battery ---
|
||
"bodyBatteryHigh": _num(s.get("bodyBatteryHighestValue")),
|
||
"bodyBatteryLow": _num(s.get("bodyBatteryLowestValue")),
|
||
"bodyBatteryCharged": _num(s.get("bodyBatteryChargedValue")),
|
||
"bodyBatteryDrained": _num(s.get("bodyBatteryDrainedValue")),
|
||
# --- breathing / blood oxygen ---
|
||
"spo2Avg": _num(s.get("averageSpo2")),
|
||
"spo2Min": _num(s.get("lowestSpo2")),
|
||
"respirationAvg": _num(
|
||
s.get("avgWakingRespirationValue"), s.get("latestRespirationValue")
|
||
),
|
||
"respirationMin": _num(s.get("lowestRespirationValue")),
|
||
"respirationMax": _num(s.get("highestRespirationValue")),
|
||
# --- sleep ---
|
||
"sleepDuration": round(sleep_seconds / 3600, 1) if sleep_seconds else None,
|
||
"sleepQuality": _num((scores.get("overall") or {}).get("value")),
|
||
"sleepDeepSeconds": secs("deepSleepSeconds"),
|
||
"sleepLightSeconds": secs("lightSleepSeconds"),
|
||
"sleepRemSeconds": secs("remSleepSeconds"),
|
||
"sleepAwakeSeconds": secs("awakeSleepSeconds"),
|
||
"sleepSpo2Avg": secs("averageSpO2Value"),
|
||
"sleepRespirationAvg": secs("averageRespirationValue"),
|
||
"sleepStressAvg": secs("avgSleepStress"),
|
||
# --- training ---
|
||
"trainingReadiness": _num(readiness.get("score")),
|
||
"vo2max": _num(vo2.get("vo2MaxValue")),
|
||
"enduranceScore": _num(endurance.get("overallScore")),
|
||
}
|
||
|
||
|
||
def sync_badges(client, user_id):
|
||
"""Earned badges. Keyed by Garmin's badge id, so re-syncing updates."""
|
||
badges = _safe(lambda: client.get_earned_badges(), []) or []
|
||
stored = 0
|
||
for b in badges:
|
||
bid = b.get("badgeId")
|
||
if bid is None:
|
||
continue
|
||
health.upsert_badge(user_id, {
|
||
"id": str(bid),
|
||
"badgeKey": b.get("badgeKey"),
|
||
"name": b.get("badgeName"),
|
||
"categoryId": _num(b.get("badgeCategoryId")),
|
||
"difficultyId": _num(b.get("badgeDifficultyId")),
|
||
"earnedDate": _to_datetime(b.get("badgeEarnedDate")),
|
||
"earnedCount": _num(b.get("badgeEarnedNumber")),
|
||
"points": _num(b.get("badgePoints")),
|
||
})
|
||
stored += 1
|
||
return stored
|
||
|
||
|
||
def sync_personal_records(client, user_id):
|
||
records = _safe(lambda: client.get_personal_record(), []) or []
|
||
stored = 0
|
||
for r in records:
|
||
rid = r.get("id")
|
||
if rid is None:
|
||
continue
|
||
health.upsert_personal_record(user_id, {
|
||
"id": str(rid),
|
||
"typeId": _num(r.get("typeId")),
|
||
"activityId": r.get("activityId"),
|
||
"activityName": r.get("activityName"),
|
||
"activityType": r.get("activityType"),
|
||
"value": _num(r.get("value")),
|
||
# Prefer the pre-formatted strings; the bare fields are epoch ms.
|
||
"achievedAt": _to_datetime(
|
||
r.get("prStartTimeLocalFormatted"),
|
||
r.get("prStartTimeGmtFormatted"),
|
||
r.get("activityStartDateTimeLocalFormatted"),
|
||
r.get("prStartTimeLocal"),
|
||
r.get("prStartTimeGmt"),
|
||
),
|
||
})
|
||
stored += 1
|
||
return stored
|
||
|
||
|
||
def _activity_end(start, duration_seconds):
|
||
if not start or not duration_seconds:
|
||
return start
|
||
for fmt in ("%Y-%m-%dT%H:%M:%S", "%Y-%m-%d %H:%M:%S", "%Y-%m-%dT%H:%M:%S.%f"):
|
||
try:
|
||
dt = datetime.datetime.strptime(start[:26], fmt)
|
||
return (dt + datetime.timedelta(seconds=duration_seconds)).isoformat()
|
||
except ValueError:
|
||
continue
|
||
return start
|
||
|
||
|
||
def _sync_activities(client, user_id, start_date, end_date):
|
||
"""Fetch the window's activities in one call and store the new ones."""
|
||
activities = client.get_activities_by_date(start_date, end_date) or []
|
||
stored = 0
|
||
for a in activities:
|
||
start = a.get("startTimeLocal") or a.get("startTime")
|
||
activity_type = (
|
||
(a.get("activityType") or {}).get("typeKey")
|
||
if isinstance(a.get("activityType"), dict)
|
||
else a.get("activityType")
|
||
) or "unknown"
|
||
duration = _num(a.get("duration"))
|
||
|
||
# Garmin activity ids are stable, so re-syncing a window must not
|
||
# duplicate what is already stored.
|
||
garmin_id = a.get("activityId")
|
||
if garmin_id is not None:
|
||
existing = query_one(
|
||
"SELECT id FROM activities WHERE user_id = ? AND id = ?",
|
||
[user_id, str(garmin_id)],
|
||
)
|
||
if existing:
|
||
continue
|
||
|
||
health.insert_activity(
|
||
user_id,
|
||
{
|
||
"id": str(garmin_id) if garmin_id is not None else None,
|
||
"activityType": activity_type,
|
||
"startTime": start,
|
||
"endTime": _activity_end(start, duration),
|
||
"duration": duration,
|
||
"distance": _num(a.get("distance")),
|
||
"calories": _num(a.get("calories")),
|
||
"heartRateAverage": _num(a.get("averageHR")),
|
||
"heartRateMax": _num(a.get("maxHR")),
|
||
},
|
||
)
|
||
stored += 1
|
||
return stored
|
||
|
||
|
||
def sync_data(user_id, creds, days=None, client=None):
|
||
"""Pull the last `days` days from Garmin Connect into the local database.
|
||
|
||
`client` exists so tests can inject a stub instead of reaching Garmin.
|
||
"""
|
||
days = days or DEFAULT_SYNC_DAYS
|
||
now = datetime.datetime.utcnow().isoformat(timespec="seconds")
|
||
_set_sync_status(user_id, "syncing", now, records_synced=0)
|
||
|
||
try:
|
||
client = client or _connect(creds, user_id)
|
||
except Exception as e:
|
||
message = describe(e)
|
||
_set_sync_status(user_id, "error", now, records_synced=0, last_error=message)
|
||
return {
|
||
"status": "error",
|
||
"recordsSynced": 0,
|
||
"message": message,
|
||
"mfaRequired": isinstance(e, MFARequired),
|
||
"lastSyncTime": now,
|
||
}
|
||
|
||
today = datetime.date.today()
|
||
start_date = (today - datetime.timedelta(days=days - 1)).isoformat()
|
||
|
||
days_synced = 0
|
||
day_errors = []
|
||
for i in range(days):
|
||
date_str = (today - datetime.timedelta(days=i)).isoformat()
|
||
try:
|
||
record = _extract_daily(client, date_str)
|
||
except Exception as e:
|
||
day_errors.append(f"{date_str}: {describe(e)}")
|
||
continue
|
||
# A day Garmin has no data for comes back all-None; storing it would
|
||
# create an empty row that the metric endpoints then have to filter.
|
||
if any(record[k] is not None for k in record if k != "date"):
|
||
health.upsert_health_daily(user_id, record)
|
||
days_synced += 1
|
||
|
||
activities_synced = 0
|
||
try:
|
||
activities_synced = _sync_activities(
|
||
client, user_id, start_date, today.isoformat()
|
||
)
|
||
except Exception as e:
|
||
day_errors.append(f"activities: {describe(e)}")
|
||
|
||
# Badges and personal records are account-wide rather than per-day, so
|
||
# they are fetched once per sync rather than inside the day loop.
|
||
badges_synced = 0
|
||
records_synced_pr = 0
|
||
try:
|
||
badges_synced = sync_badges(client, user_id)
|
||
except Exception as e:
|
||
day_errors.append(f"badges: {describe(e)}")
|
||
try:
|
||
records_synced_pr = sync_personal_records(client, user_id)
|
||
except Exception as e:
|
||
day_errors.append(f"personal_records: {describe(e)}")
|
||
|
||
# Every single day failing means something systemic (expired session,
|
||
# API change) — reporting that as a clean success would hide it.
|
||
if days_synced == 0 and len(day_errors) >= days:
|
||
message = "; ".join(day_errors[:3])
|
||
_set_sync_status(user_id, "error", now, records_synced=0, last_error=message)
|
||
return {"status": "error", "recordsSynced": 0,
|
||
"message": f"同步失败:{message}", "lastSyncTime": now}
|
||
|
||
_set_sync_status(
|
||
user_id, "idle", now, records_synced=days_synced,
|
||
last_error="; ".join(day_errors[:3]) if day_errors else None,
|
||
)
|
||
message = (
|
||
f"同步完成,更新 {days_synced} 天数据、{activities_synced} 条运动记录、"
|
||
f"{badges_synced} 个奖励、{records_synced_pr} 项个人纪录"
|
||
)
|
||
if day_errors:
|
||
message += f"({len(day_errors)} 项跳过)"
|
||
return {
|
||
"status": "success",
|
||
"recordsSynced": days_synced,
|
||
"activitiesSynced": activities_synced,
|
||
"badgesSynced": badges_synced,
|
||
"personalRecordsSynced": records_synced_pr,
|
||
"message": message,
|
||
"lastSyncTime": now,
|
||
}
|