db.py: 修复迁移非幂等——gunicorn 双 worker 并发 init_db 时后到者遇 Duplicate column 会让整个服务起不来,改为吞掉 duplicate column 错误。garmin.py/scheduler.py: 退避期写入 sync_status.rate_limited_until,自动同步与手动同步在退避期内跳过,不再反复撞击 Garmin 把限流窗口越撞越深。
1080 lines
41 KiB
Python
1080 lines
41 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 json
|
||
import os
|
||
import threading
|
||
|
||
from db import execute, query_one, query_all
|
||
from config import DB_TYPE
|
||
from services import health
|
||
from services import garmin_extras as extras
|
||
|
||
# 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"),
|
||
"stage": row.get("stage"),
|
||
}
|
||
|
||
|
||
def reset_stale_syncs():
|
||
"""Clear a "syncing" status left behind by a process that went away.
|
||
|
||
The status lives in the database but the work lives in a thread. A restart
|
||
(or a crash) takes the thread and leaves the row, so the UI shows a
|
||
progress bar that will never move and refuses to start a new sync.
|
||
"""
|
||
execute(
|
||
"UPDATE sync_status SET status = 'idle', stage = NULL "
|
||
"WHERE status = 'syncing'"
|
||
)
|
||
|
||
|
||
class RateLimited(RuntimeError):
|
||
"""Garmin answered 429.
|
||
|
||
It reaches us disguised: the body is the plain text "Rate limited", and
|
||
garth feeds that to json.loads, so the exception surfaced as
|
||
`JSONDecodeError: Expecting value: line 1 column 1` — which reads like a
|
||
parsing bug rather than "stop asking". Naming it means the sync status
|
||
says what is actually wrong.
|
||
"""
|
||
|
||
|
||
# Retrying while rate limited is what deepens the limit, so once Garmin says
|
||
# 429 the whole process stands down until this passes. Bumped from 30 to 60
|
||
# minutes: the account was stuck for days because every hourly tick re-hit it,
|
||
# so a longer cooldown gives Garmin's window room to actually close.
|
||
RATE_LIMIT_BACKOFF = datetime.timedelta(minutes=60)
|
||
# In-process cache of the cooldown, kept in sync with the DB copy below and
|
||
# still the lever the tests reach for via _rate_limited_until.clear().
|
||
_rate_limited_until = {}
|
||
|
||
|
||
def _parse_dt(v):
|
||
"""Coerce a stored rate-limit time into a naive UTC datetime, or None."""
|
||
if v is None:
|
||
return None
|
||
if isinstance(v, datetime.datetime):
|
||
return v
|
||
if isinstance(v, str):
|
||
for fmt in ("%Y-%m-%d %H:%M:%S", "%Y-%m-%dT%H:%M:%S"):
|
||
try:
|
||
return datetime.datetime.strptime(v, fmt)
|
||
except ValueError:
|
||
continue
|
||
try:
|
||
return datetime.datetime.fromisoformat(v)
|
||
except ValueError:
|
||
return None
|
||
return None
|
||
|
||
|
||
def rate_limited_until(user_id):
|
||
"""When the account is still cooling down, as a UTC datetime (or None).
|
||
|
||
The cooldown is persisted to the database so every gunicorn worker and a
|
||
process restart see the same deadline — a process-local dict alone let each
|
||
worker re-hit Garmin and keep the limit alive forever.
|
||
"""
|
||
mem = _rate_limited_until.get(user_id)
|
||
try:
|
||
row = query_one(
|
||
"SELECT rate_limited_until FROM sync_status WHERE user_id = ?", [user_id]
|
||
)
|
||
except Exception:
|
||
row = None
|
||
candidates = [c for c in (mem, _parse_dt(row["rate_limited_until"] if row else None)) if c]
|
||
return max(candidates) if candidates else None
|
||
|
||
|
||
def _note_rate_limit(user_id):
|
||
"""Record a rate-limit cooldown, extending it if one is already running."""
|
||
now = datetime.datetime.utcnow()
|
||
existing = rate_limited_until(user_id)
|
||
until = (existing + datetime.timedelta(minutes=30)) if existing and existing > now \
|
||
else (now + RATE_LIMIT_BACKOFF)
|
||
_rate_limited_until[user_id] = until
|
||
try:
|
||
# Persist alongside the current status so the cooldown survives across
|
||
# workers and restarts.
|
||
cur_status = (get_sync_status(user_id) or {}).get("status") or "idle"
|
||
_set_sync_status(
|
||
user_id, cur_status, now.isoformat(timespec="seconds"),
|
||
rate_limited_until=until.isoformat(timespec="seconds"),
|
||
)
|
||
except Exception:
|
||
pass
|
||
|
||
|
||
def _rate_limit_block(user_id):
|
||
"""If the account is cooling down, return (until, message); else (None, None).
|
||
|
||
Central guard used by both the manual and scheduled sync entry points, so a
|
||
blocked account issues zero Garmin requests until the window closes.
|
||
"""
|
||
until = rate_limited_until(user_id)
|
||
if not until or datetime.datetime.utcnow() >= until:
|
||
return None, None
|
||
mins = max(1, int((until - datetime.datetime.utcnow()).total_seconds() // 60))
|
||
return until, (
|
||
f"Garmin 仍在限制请求频率,预计约 {mins} 分钟后自动恢复。"
|
||
"已自动退避,请耐心等待——反复点击正是把限流撞得更深的原因,令牌本身没有失效。"
|
||
)
|
||
|
||
|
||
def _is_rate_limited(e):
|
||
"""429 from Garmin, however it happens to be dressed.
|
||
|
||
Usually it arrives as a JSONDecodeError, because garth calls .json() on a
|
||
429 whose body is the plain text "Rate limited" — the status code is gone
|
||
by the time the exception reaches us, and the message is the useless
|
||
"Expecting value: line 1 column 1". What survives is the body itself, on
|
||
the exception's `doc` attribute, so that is where to look.
|
||
"""
|
||
response = getattr(e, "response", None)
|
||
if response is not None and getattr(response, "status_code", None) == 429:
|
||
return True
|
||
|
||
haystacks = [str(e), str(getattr(e, "doc", "") or "")]
|
||
return any("rate limit" in h.lower() for h in haystacks)
|
||
|
||
|
||
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")])
|
||
# A re-bind means the old session is stale; the next call must build a
|
||
# fresh one rather than keep using the session the old token minted.
|
||
forget_client(user_id)
|
||
|
||
|
||
def has_token(user_id):
|
||
return load_token(user_id) is not None
|
||
|
||
|
||
# An authenticated client, reused across requests in this process.
|
||
#
|
||
# Building one costs ~11s against Garmin — loading the token, refreshing the
|
||
# OAuth2 grant and fetching the profile — which dwarfed the ~4s of actual data
|
||
# fetching behind an activity-detail request. The session is a requests.Session
|
||
# underneath, so it is reusable; it is dropped after CLIENT_TTL so a refreshed
|
||
# or revoked token is picked up rather than being cached indefinitely.
|
||
CLIENT_TTL_SECONDS = 900
|
||
_clients = {}
|
||
_clients_lock = threading.Lock()
|
||
|
||
|
||
def _cached_client(user_id):
|
||
entry = _clients.get(user_id)
|
||
if entry and (datetime.datetime.utcnow() - entry[1]).total_seconds() < CLIENT_TTL_SECONDS:
|
||
return entry[0]
|
||
return None
|
||
|
||
|
||
def _cache_client(user_id, client):
|
||
if user_id:
|
||
_clients[user_id] = (client, datetime.datetime.utcnow())
|
||
|
||
|
||
def forget_client(user_id):
|
||
"""Drop the cached session — call after re-binding an account."""
|
||
_clients.pop(user_id, 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.
|
||
"""
|
||
if user_id:
|
||
with _clients_lock:
|
||
cached = _cached_client(user_id)
|
||
if cached is not None:
|
||
return cached
|
||
|
||
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)
|
||
|
||
blocked = rate_limited_until(user_id)
|
||
if blocked and datetime.datetime.utcnow() < blocked:
|
||
raise RateLimited(
|
||
"Garmin 暂时限制了请求频率,稍后会自动恢复(约 "
|
||
f"{max(1, int((blocked - datetime.datetime.utcnow()).total_seconds() // 60))} 分钟)。"
|
||
)
|
||
|
||
# Only when it has actually expired. Refreshing on every connect spends
|
||
# quota for nothing, and that is what walked the account into a 429.
|
||
oauth2 = getattr(client.garth, "oauth2_token", None)
|
||
if oauth2 is None or getattr(oauth2, "expired", True):
|
||
try:
|
||
client.garth.refresh_oauth2()
|
||
except Exception as e: # noqa: BLE001 - re-raised, just named better
|
||
if _is_rate_limited(e):
|
||
_note_rate_limit(user_id)
|
||
raise RateLimited(
|
||
"Garmin 暂时限制了请求频率。这通常是短时间内连接过于频繁,"
|
||
"等待约半小时后会自动恢复,令牌本身没有失效。"
|
||
) from e
|
||
raise
|
||
# 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"]
|
||
with _clients_lock:
|
||
_cache_client(user_id, client)
|
||
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)
|
||
with _clients_lock:
|
||
_cache_client(user_id, 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
|
||
|
||
|
||
# --- one activity, in full ---------------------------------------------------
|
||
|
||
# Garmin will return thousands of samples per activity. A phone chart cannot
|
||
# draw more than a few hundred usefully, and the payload is stored as a row, so
|
||
# the series are thinned on the way in rather than on every read.
|
||
DETAIL_MAX_POINTS = 300
|
||
|
||
# Descriptor key -> the name the UI charts by. Anything not listed is dropped:
|
||
# the full descriptor set runs to dozens of fields, most of them empty.
|
||
SERIES_KEYS = {
|
||
"directTimestamp": "timestamp",
|
||
"sumElapsedDuration": "elapsed",
|
||
"sumDuration": "duration",
|
||
"sumDistance": "distance",
|
||
"directHeartRate": "heartRate",
|
||
"directSpeed": "speed",
|
||
"directElevation": "elevation",
|
||
"directRunCadence": "cadence",
|
||
"directBikeCadence": "cadence",
|
||
"directDoubleCadence": "cadence",
|
||
"directPower": "power",
|
||
"directAirTemperature": "temperature",
|
||
}
|
||
|
||
|
||
def _thin(values, limit=DETAIL_MAX_POINTS):
|
||
"""Evenly sample a list down to `limit` points, keeping first and last."""
|
||
if len(values) <= limit:
|
||
return values
|
||
step = (len(values) - 1) / (limit - 1)
|
||
return [values[int(round(i * step))] for i in range(limit)]
|
||
|
||
|
||
def _series_from_details(details):
|
||
"""Turn Garmin's column-store detail payload into per-metric arrays.
|
||
|
||
The response is a descriptor list plus rows of parallel values, so every
|
||
metric has to be read out by the index its descriptor names.
|
||
"""
|
||
descriptors = details.get("metricDescriptors") or []
|
||
rows = details.get("activityDetailMetrics") or []
|
||
if not descriptors or not rows:
|
||
return {}
|
||
|
||
index = {}
|
||
for d in descriptors:
|
||
name = SERIES_KEYS.get(d.get("key"))
|
||
if name and name not in index:
|
||
index[name] = d.get("metricsIndex")
|
||
|
||
rows = _thin(rows)
|
||
out = {}
|
||
for name, position in index.items():
|
||
if position is None:
|
||
continue
|
||
column = []
|
||
for row in rows:
|
||
metrics = row.get("metrics") or []
|
||
column.append(metrics[position] if position < len(metrics) else None)
|
||
# A column of nothing but nulls is a sensor the watch does not have.
|
||
if any(v is not None for v in column):
|
||
out[name] = column
|
||
return out
|
||
|
||
|
||
def _lap_rows(splits):
|
||
laps = []
|
||
for i, lap in enumerate((splits or {}).get("lapDTOs") or [], start=1):
|
||
laps.append({
|
||
"index": lap.get("lapIndex") or i,
|
||
"duration": _num(lap.get("duration")),
|
||
"movingDuration": _num(lap.get("movingDuration")),
|
||
"distance": _num(lap.get("distance")),
|
||
"averageSpeed": _num(lap.get("averageSpeed")),
|
||
"maxSpeed": _num(lap.get("maxSpeed")),
|
||
"calories": _num(lap.get("calories")),
|
||
"averageHR": _num(lap.get("averageHR")),
|
||
"maxHR": _num(lap.get("maxHR")),
|
||
"elevationGain": _num(lap.get("elevationGain")),
|
||
"elevationLoss": _num(lap.get("elevationLoss")),
|
||
})
|
||
return laps
|
||
|
||
|
||
def _hr_zones(zones):
|
||
out = []
|
||
for z in zones or []:
|
||
out.append({
|
||
"zone": z.get("zoneNumber"),
|
||
"seconds": _num(z.get("secsInZone")) or 0,
|
||
"lowBoundary": _num(z.get("zoneLowBoundary")),
|
||
})
|
||
return sorted(out, key=lambda z: z.get("zone") or 0)
|
||
|
||
|
||
def _build_detail(client, activity_id):
|
||
"""Assemble everything Garmin knows about one activity.
|
||
|
||
Each call is wrapped: a watch without a barometer has no weather, a
|
||
treadmill run has no gear, and a missing optional endpoint must leave the
|
||
rest of the page intact rather than fail the request.
|
||
"""
|
||
summary = _safe(lambda: client.get_activity_evaluation(activity_id), {}) or {}
|
||
# 500 is already more samples than the 300 we keep, and asking for 2000
|
||
# triples the payload for points that get thinned away anyway.
|
||
details = _safe(
|
||
lambda: client.get_activity_details(activity_id, maxchart=500, maxpoly=0), {}
|
||
) or {}
|
||
|
||
return {
|
||
"activityId": str(activity_id),
|
||
"summary": summary.get("summaryDTO") or {},
|
||
"activityName": summary.get("activityName"),
|
||
"activityType": (summary.get("activityTypeDTO") or {}).get("typeKey"),
|
||
"eventType": (summary.get("eventTypeDTO") or {}).get("typeKey"),
|
||
"laps": _lap_rows(_safe(lambda: client.get_activity_splits(activity_id), {})),
|
||
"hrZones": _hr_zones(
|
||
_safe(lambda: client.get_activity_hr_in_timezones(activity_id), [])
|
||
),
|
||
"weather": _safe(lambda: client.get_activity_weather(activity_id), {}) or {},
|
||
"gear": _safe(lambda: client.get_activity_gear(activity_id), []) or [],
|
||
"exerciseSets": (
|
||
_safe(lambda: client.get_activity_exercise_sets(activity_id), {}) or {}
|
||
).get("exerciseSets") or [],
|
||
"series": _series_from_details(details),
|
||
}
|
||
|
||
|
||
def _store_detail(user_id, activity_id, detail):
|
||
cols = ["activity_id", "user_id", "payload", "fetched_at"]
|
||
values = [activity_id, user_id, json.dumps(detail, default=str),
|
||
datetime.datetime.utcnow().isoformat(timespec="seconds")]
|
||
placeholders = ", ".join(["?"] * len(cols))
|
||
if DB_TYPE == "mariadb":
|
||
updates = ", ".join(f"{c}=VALUES({c})" for c in cols if c != "activity_id")
|
||
sql = (f"INSERT INTO activity_details ({', '.join(cols)}) "
|
||
f"VALUES ({placeholders}) ON DUPLICATE KEY UPDATE {updates}")
|
||
else:
|
||
updates = ", ".join(f"{c}=excluded.{c}" for c in cols if c != "activity_id")
|
||
sql = (f"INSERT INTO activity_details ({', '.join(cols)}) VALUES "
|
||
f"({placeholders}) ON CONFLICT(activity_id) DO UPDATE SET {updates}")
|
||
execute(sql, values)
|
||
|
||
|
||
def read_activity_detail(user_id, activity_id):
|
||
"""The stored detail for one activity, or None if it was never synced."""
|
||
row = query_one(
|
||
"SELECT payload FROM activity_details WHERE user_id = ? AND activity_id = ?",
|
||
[user_id, str(activity_id)],
|
||
)
|
||
if not row or not row.get("payload"):
|
||
return None
|
||
try:
|
||
return json.loads(row["payload"])
|
||
except ValueError:
|
||
# A truncated row is worth re-syncing, not worth crashing on.
|
||
return None
|
||
|
||
|
||
def sync_activity_details(client, user_id, limit=None, on_progress=None):
|
||
"""Fetch and store the full detail for activities that lack one.
|
||
|
||
Detail used to be fetched when the user opened an activity, which meant
|
||
seven Garmin calls on the critical path of a tap: slow at best, and a
|
||
timeout whenever the link was poor. It belongs in the sync, so the screen
|
||
only ever reads the local database.
|
||
"""
|
||
rows = query_all(
|
||
"SELECT a.id FROM activities a "
|
||
"LEFT JOIN activity_details d ON d.activity_id = a.id "
|
||
"WHERE a.user_id = ? AND d.activity_id IS NULL "
|
||
"ORDER BY a.start_time DESC",
|
||
[user_id],
|
||
)
|
||
if limit:
|
||
rows = rows[:limit]
|
||
|
||
stored = 0
|
||
for i, row in enumerate(rows):
|
||
try:
|
||
_store_detail(user_id, row["id"], _build_detail(client, row["id"]))
|
||
stored += 1
|
||
except Exception: # noqa: BLE001 - one bad activity must not stop the rest
|
||
continue
|
||
if on_progress:
|
||
on_progress(i + 1, len(rows))
|
||
return stored
|
||
|
||
|
||
_backfill_progress = {}
|
||
|
||
|
||
def backfill_status(user_id):
|
||
"""Progress of the historical backfill for this account."""
|
||
return _backfill_progress.get(user_id) or {
|
||
"running": False, "stage": None, "done": 0, "total": 0, "error": None,
|
||
}
|
||
|
||
|
||
def _set_backfill(user_id, **fields):
|
||
state = dict(_backfill_progress.get(user_id) or {})
|
||
state.update(fields)
|
||
_backfill_progress[user_id] = state
|
||
|
||
|
||
def days_missing_series(user_id, limit=None):
|
||
"""Days that have a health row but no within-day curves stored."""
|
||
rows = query_all(
|
||
"SELECT h.date FROM health_data h "
|
||
"LEFT JOIN daily_series s ON s.user_id = h.user_id AND s.date = h.date "
|
||
"WHERE h.user_id = ? AND s.date IS NULL "
|
||
"GROUP BY h.date ORDER BY h.date DESC",
|
||
[user_id],
|
||
)
|
||
dates = [str(r["date"])[:10] for r in rows]
|
||
return dates[:limit] if limit else dates
|
||
|
||
|
||
def start_backfill(user_id, limit=None):
|
||
"""Fill in everything the per-day sync leaves out, in the background.
|
||
|
||
Two long jobs share one runner because they share a cause — an account
|
||
whose history predates these features — and because the user should press
|
||
one button, not two. Each activity costs several Garmin calls and each day
|
||
of curves costs five, so this runs for minutes; the UI polls.
|
||
"""
|
||
state = _backfill_progress.get(user_id)
|
||
if state and state.get("running"):
|
||
return state
|
||
|
||
_set_backfill(user_id, running=True, stage="启动中", done=0, total=0, error=None)
|
||
|
||
def run():
|
||
try:
|
||
client = _connect({}, user_id=user_id)
|
||
|
||
_set_backfill(user_id, stage="运动详情", done=0, total=0)
|
||
sync_activity_details(
|
||
client, user_id, limit,
|
||
on_progress=lambda d, n: _set_backfill(
|
||
user_id, stage="运动详情", done=d, total=n),
|
||
)
|
||
|
||
dates = days_missing_series(user_id, limit)
|
||
_set_backfill(user_id, stage="每日曲线", done=0, total=len(dates))
|
||
for i, date in enumerate(dates):
|
||
try:
|
||
extras.sync_daily_series(client, user_id, date)
|
||
except Exception: # noqa: BLE001 - one day must not stop the rest
|
||
pass
|
||
_set_backfill(user_id, stage="每日曲线", done=i + 1,
|
||
total=len(dates))
|
||
|
||
_set_backfill(user_id, running=False, stage="完成", error=None)
|
||
except Exception as e: # noqa: BLE001 - reported through the status endpoint
|
||
_set_backfill(user_id, running=False, stage=None, error=describe(e))
|
||
|
||
threading.Thread(target=run, daemon=True, name=f"backfill-{user_id}").start()
|
||
return backfill_status(user_id)
|
||
|
||
|
||
# Up to this many days, a sync also pulls each day's within-day curves inline.
|
||
# Beyond it the curves are left to the background backfill: five extra calls
|
||
# per day would turn a year's sync into an hour.
|
||
SERIES_INLINE_DAYS = 14
|
||
|
||
|
||
# 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")
|
||
until, msg = _rate_limit_block(user_id)
|
||
if until:
|
||
_set_sync_status(
|
||
user_id, "rate_limited", now,
|
||
records_synced=0, progress_current=0, progress_total=days,
|
||
started_at=now, last_error=msg,
|
||
)
|
||
return {
|
||
"status": "rate_limited",
|
||
"message": msg,
|
||
"retryAfterSeconds": int((until - datetime.datetime.utcnow()).total_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")
|
||
until, msg = _rate_limit_block(user_id)
|
||
if until:
|
||
# A blocked account must issue zero Garmin requests — that is the whole
|
||
# point. Return immediately without touching the network.
|
||
_set_sync_status(
|
||
user_id, "rate_limited", now,
|
||
records_synced=0, last_error=msg,
|
||
)
|
||
return {
|
||
"status": "rate_limited",
|
||
"recordsSynced": 0,
|
||
"message": msg,
|
||
"retryAfterSeconds": int((until - datetime.datetime.utcnow()).total_seconds()),
|
||
}
|
||
_set_sync_status(
|
||
user_id, "syncing", now,
|
||
records_synced=0, progress_current=0, progress_total=days,
|
||
stage="连接 Garmin",
|
||
)
|
||
|
||
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)
|
||
record.update(extras.daily_extras(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
|
||
# Within-day curves for short syncs only. A year-long backfill
|
||
# would add five calls per day on top of everything else; those
|
||
# days are filled by start_backfill instead.
|
||
if days <= SERIES_INLINE_DAYS:
|
||
try:
|
||
extras.sync_daily_series(client, user_id, date_str)
|
||
except Exception as e: # noqa: BLE001
|
||
day_errors.append(f"{date_str} series: {describe(e)}")
|
||
|
||
# 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, stage=f"每日数据 {date_str}",
|
||
)
|
||
|
||
activities_synced = 0
|
||
_set_sync_status(user_id, "syncing", now, records_synced=days_synced,
|
||
progress_current=days, progress_total=days,
|
||
stage="运动记录")
|
||
try:
|
||
activities_synced = _sync_activities(
|
||
client, user_id, start_date, today.isoformat()
|
||
)
|
||
except Exception as e:
|
||
day_errors.append(f"activities: {describe(e)}")
|
||
|
||
# Full detail for any activity that does not have it stored yet, so that
|
||
# opening one later is a local read.
|
||
details_synced = 0
|
||
try:
|
||
details_synced = sync_activity_details(
|
||
client, user_id,
|
||
on_progress=lambda d, n: _set_sync_status(
|
||
user_id, "syncing", now, records_synced=days_synced,
|
||
progress_current=days, progress_total=days,
|
||
stage=f"运动详情 {d}/{n}"),
|
||
)
|
||
except Exception as e:
|
||
day_errors.append(f"activity_details: {describe(e)}")
|
||
|
||
# Everything else Garmin holds: body composition, blood pressure, race
|
||
# predictions, challenges and devices. Account-wide, so once per sync.
|
||
extra_counts = {}
|
||
stage_names = {
|
||
"bodyComposition": "身体成分", "bloodPressure": "血压",
|
||
"racePredictions": "成绩预测", "challenges": "挑战赛", "devices": "设备",
|
||
}
|
||
for name, call in (
|
||
("bodyComposition",
|
||
lambda: extras.sync_body_composition(client, user_id, start_date,
|
||
today.isoformat())),
|
||
("bloodPressure",
|
||
lambda: extras.sync_blood_pressure(client, user_id, start_date,
|
||
today.isoformat())),
|
||
("racePredictions",
|
||
lambda: extras.sync_race_predictions(client, user_id)),
|
||
("challenges", lambda: extras.sync_challenges(client, user_id)),
|
||
("devices", lambda: extras.sync_devices(client, user_id)),
|
||
):
|
||
_set_sync_status(user_id, "syncing", now, records_synced=days_synced,
|
||
progress_current=days, progress_total=days,
|
||
stage=stage_names.get(name, name))
|
||
try:
|
||
extra_counts[name] = call()
|
||
except Exception as e: # noqa: BLE001 - one section must not fail the sync
|
||
day_errors.append(f"{name}: {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, stage=None,
|
||
last_error="; ".join(day_errors[:3]) if day_errors else None,
|
||
)
|
||
message = (
|
||
f"同步完成,更新 {days_synced} 天数据、{activities_synced} 条运动记录"
|
||
f"(含 {details_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,
|
||
}
|