Files
GarminHealthLab/backend/services/garmin.py
ericwyuan 6a5cfa7806 [阶段7.1] 同步改为后台任务 + 进度上报,支持回补历史
趋势页提供了"一年"档,但库里只有 7 天数据,那一档形同虚设。
实测每天约 2.84 秒(一天要打 5 个端点),回补一年需要约 28 分钟,
远超任何 HTTP 超时能等的时间。

- sync_status 新增 progress_current / progress_total / started_at
- start_sync() 起后台线程并立即返回,sync_data 每 5 天写一次进度
  (写库便宜但不免费,而前端本来就是 2 秒一轮询)
- POST /api/garmin/sync 改为 202 立即返回,接受 days 参数并
  夹在 1..730;进度经 GET /status 轮询
- 一次新同步会清掉上一次的错误,避免旧错误一直挂在界面上

前端:
- 同步页给出 7 / 30 / 90 / 365 天四个选项,日常与首次回补分开
- 进度条显示"第 N / 共 M 天"与预计耗时,并说明可以离开本页
- 页面挂载时若发现正在同步会接着轮询 —— 回补比页面存活时间长,
  刷新后必须能接上进度

tests (+7, 共 299):
- start_sync 在工作完成前就返回,且返回前已把 total 写好
- 进度随同步推进,结束时等于总天数
- days 超范围被夹到 730
- 新同步清除上一次的错误

Co-Authored-By: Claude Haiku 4.5 <noreply@anthropic.com>
2026-08-23 21:31:50 +08:00

544 lines
20 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
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
import threading
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"],
"progressCurrent": row.get("progress_current"),
"progressTotal": row.get("progress_total"),
"startedAt": row.get("started_at"),
}
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
# Above this many days a sync is long enough that the caller must not block
# on it — a year takes roughly 20 minutes at ~3s per day.
BACKGROUND_THRESHOLD_DAYS = 14
def start_sync(user_id, creds, days=None):
"""Run a sync in the background and return immediately.
Progress lands in sync_status, which the UI polls; a full backfill runs
far longer than any sensible HTTP timeout.
"""
days = days or DEFAULT_SYNC_DAYS
now = datetime.datetime.utcnow().isoformat(timespec="seconds")
_set_sync_status(
user_id, "syncing", now,
records_synced=0, progress_current=0, progress_total=days,
started_at=now, last_error=None,
)
thread = threading.Thread(
target=sync_data, args=(user_id, creds, days), daemon=True
)
thread.start()
return {"status": "syncing", "days": days}
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, progress_current=0, progress_total=days,
)
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
# Reported every few days rather than every day: the write is cheap
# but not free, and the UI polls on a 2s cadence anyway.
if (i + 1) % 5 == 0 or i + 1 == days:
_set_sync_status(
user_id, "syncing", now,
records_synced=days_synced, progress_current=i + 1,
progress_total=days,
)
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,
progress_current=days, progress_total=days,
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,
}