From 682936b0b69e48a500166431372c2ba7d2d55fa1 Mon Sep 17 00:00:00 2001 From: ericwyuan Date: Thu, 3 Sep 2026 23:24:48 +0800 Subject: [PATCH] =?UTF-8?q?fix(garmin):=20=E5=88=B7=E6=96=B0=E5=87=BA?= =?UTF-8?q?=E6=9D=A5=E7=9A=84=E4=BB=A4=E7=89=8C=E4=BB=8E=E6=9D=A5=E6=B2=A1?= =?UTF-8?q?=E5=86=99=E5=9B=9E=E5=BA=93=EF=BC=8C=E4=BA=8E=E6=98=AF=E6=AF=8F?= =?UTF-8?q?=E6=AC=A1=E8=BF=9E=E6=8E=A5=E9=83=BD=E9=87=8D=E6=8D=A2=E4=B8=80?= =?UTF-8?q?=E6=AC=A1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 「又被限流了」的根因找到了,不是请求量,是令牌。 `_connect` 里 `refresh_oauth2()` 换来的新 OAuth2 令牌只活在进程内存里—— `save_token` 只在绑定账号时调用过一次。于是每次客户端缓存过期(15 分钟)、 每个 gunicorn worker、每次部署重启,都从库里读回**同一个已过期的令牌**, 然后再做一次真实 SSO 换令牌。而 SSO 端点是按账号限流最狠的那个,社区报告能 封 48 小时(garth #217、python-garminconnect #337)。我今天为了部署重启了 八次服务,每次都清掉缓存。 - `refresh_oauth2()` 成功后 `_persist_token()` 写回。拆出这个函数是因为它和 `save_token` 想要的正好相反:重新绑定要作废现有会话,持久化刷新结果必须 保住刚刚产出它的那个会话 - 写回时不带 garmin_email,否则 upsert 会把绑定邮箱刷成 NULL,数据同步页会 忘记绑的是哪个账号 - 刷新加进程内锁,并在拿到锁后重读一次库:另一个线程刚换过就直接用它的, 不再自己去换一次 - 五条测试盯住这个不变量,包括「冷缓存不该再换一次」(这条如果回归,就是同一 个 bug 再来一遍) 顺带把数据端点也节流了——那是另外一半问题,不是这次的病因,但一天历史要 9 次 调用,730 天全历史 6600 个请求全速打出去,不该指望佳明一直容忍: - services/garmin_throttle.py:代理包住 client,所有调用(含以后新加的)都经 同一个收口,按间隔排队并计数 - 0.5s 是查过的:garmin-data-export 默认 0.15s、garmin-connect-scraper 默认 3s、官方合作方 API 100 次/分钟(0.6s)。依据写在文件顶部 - 单次同步 1200 个请求预算,跑满就干净收尾、下次接着跑(已存的天数本来就跳过) - 运动详情每次最多补 40 条——新账号几百条,不限量就是一次性打光预算 Co-Authored-By: Claude Opus 5 --- CLAUDE.md | 13 +- backend/.env.example | 15 ++ backend/services/garmin.py | 146 ++++++++++++++++--- backend/services/garmin_throttle.py | 202 ++++++++++++++++++++++++++ backend/tests/conftest.py | 5 + backend/tests/test_garmin_sync.py | 129 ++++++++++++++++ backend/tests/test_garmin_throttle.py | 122 ++++++++++++++++ 7 files changed, 606 insertions(+), 26 deletions(-) create mode 100644 backend/services/garmin_throttle.py create mode 100644 backend/tests/test_garmin_throttle.py diff --git a/CLAUDE.md b/CLAUDE.md index 85925db..aa7a3f4 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -111,9 +111,16 @@ cd backend && .venv/bin/python tests/smoke.py ——自定义组件里的按钮要显式 `width: auto`(Copilot 浮窗 ✕/发送按钮曾因此 错位,见 `Copilot.css` 注释)。F7 `.navbar .left/.right` 的 frosted pill 也 需在 `f7theme.css` 覆盖。 -4. **同步限流**:Garmin 账号 429 退避 24h,`rate_limited_until` 以 **DB 为准** - (多 worker 内存不一致会卡死守卫);`_is_rate_limited` 沿异常 `__cause__` - 链识别 RetryError 里的 429。 +4. **同步限流**:真正把账号打进 429 的是 **SSO 换令牌端点**,不是数据端点。 + `refresh_oauth2()` 刷出来的新令牌**必须写回 DB**(`_persist_token`)——不写回 + 的话,每次客户端缓存过期(15 分钟)、每个 gunicorn worker、每次部署重启,都会 + 从库里读回同一个过期令牌再换一次,而这个端点按账号限流、能封几小时。 + 2026-09-03 排查一整天的根因就是这个。 + 数据端点另有节流:所有调用经 `services/garmin_throttle.py` 的代理,默认间隔 + 0.5s + 单次预算 1200 次(依据见该文件顶部)。 + 429 退避 24h,`rate_limited_until` 以 **DB 为准**(多 worker 内存不一致会卡死 + 守卫),且只作**信息**不作闸门;`_is_rate_limited` 沿异常 `__cause__` 链识别 + RetryError 里的 429。 5. **AI 生成慢**:网关上游推理模型单次 2~5 分钟,晨报走后台生成 + 轮询, 接口超时/流式不可靠见 `services/ai.py` 的非流式回退;生产队列并发 `AI_JOB_CONCURRENCY=2`(共享网关,太高会把别家打成 502)。 diff --git a/backend/.env.example b/backend/.env.example index f81cb8c..ac53126 100644 --- a/backend/.env.example +++ b/backend/.env.example @@ -101,3 +101,18 @@ AI_JOB_GAP_SECONDS=5 # Set AI_JOBS=false to stop consuming entirely (screens then show the computed # figures with no model reading). AI_JOBS=true + +# --- Garmin 请求节流 --- +# 一天的历史要 9 次 API 调用(_extract_daily 6 + daily_extras 3),短同步再加 +# 5 次曲线,每条没有详情的运动 1 次。730 天全历史 ≈ 6600 个请求。以前是能发多 +# 快发多快。 +# +# 0.5 秒的依据(2026-09-03 调研,见 services/garmin_throttle.py 顶部注释): +# sirredbeard/garmin-data-export 默认 0.15s、evg656e/garmin-connect-scraper +# 默认 3s、佳明官方合作方 API 是 100 次/分钟(合 0.6s)。 +GARMIN_MIN_INTERVAL_SECONDS=0.5 +# 单次同步的请求预算。跑满就干净收尾,下次接着跑(已存的天数会跳过)。 +# 1200 × 0.5s ≈ 10 分钟,够覆盖四个月历史。 +GARMIN_REQUEST_BUDGET=1200 +# 每次同步补多少条运动详情。新账号有几百条,不限量就是一次性打光预算。 +GARMIN_DETAILS_PER_SYNC=40 diff --git a/backend/services/garmin.py b/backend/services/garmin.py index b2f0de7..a0534e4 100644 --- a/backend/services/garmin.py +++ b/backend/services/garmin.py @@ -30,6 +30,7 @@ from db import execute, query_one, query_all from config import DB_TYPE from services import health from services import garmin_extras as extras +from services import garmin_throttle as throttle # How many days back a sync reaches. DEFAULT_SYNC_DAYS = int(os.environ.get("GARMIN_SYNC_DAYS") or 7) @@ -372,7 +373,13 @@ def load_token(user_id): return row["token"] if row else None -def save_token(user_id, token, garmin_email=None): +def _persist_token(user_id, token, garmin_email=None): + """Write the token row. Does **not** drop the cached session. + + Separate from `save_token` because the two callers want opposite things: a + re-bind must invalidate the live session, while persisting a refreshed + OAuth2 token must keep the very session that just produced it. + """ cols = ["user_id", "token", "garmin_email", "updated_at"] placeholders = ", ".join(["?"] * len(cols)) if DB_TYPE == "mariadb": @@ -383,8 +390,33 @@ def save_token(user_id, token, garmin_email=None): 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")]) + params = [user_id, token, garmin_email, + datetime.datetime.utcnow().isoformat(timespec="seconds")] + if garmin_email is None: + # A refresh must not blank the remembered email: the column is part of + # the upsert, so passing None would overwrite it with NULL and the + # 数据同步 page would forget which account is bound. + cols_kept = [c for c in cols if c != "garmin_email"] + placeholders_kept = ", ".join(["?"] * len(cols_kept)) + if DB_TYPE == "mariadb": + updates_kept = ", ".join( + f"{c}=VALUES({c})" for c in cols_kept if c != "user_id") + sql = (f"INSERT INTO garmin_tokens ({', '.join(cols_kept)}) " + f"VALUES ({placeholders_kept}) " + f"ON DUPLICATE KEY UPDATE {updates_kept}") + else: + updates_kept = ", ".join( + f"{c}=excluded.{c}" for c in cols_kept if c != "user_id") + sql = (f"INSERT INTO garmin_tokens ({', '.join(cols_kept)}) " + f"VALUES ({placeholders_kept}) " + f"ON CONFLICT(user_id) DO UPDATE SET {updates_kept}") + params = [user_id, token, params[-1]] + execute(sql, params) + + +def save_token(user_id, token, garmin_email=None): + """Store a token from a fresh login, and stand the old session down.""" + _persist_token(user_id, token, garmin_email) # 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) @@ -433,6 +465,8 @@ def delete_token(user_id): CLIENT_TTL_SECONDS = 900 _clients = {} _clients_lock = threading.Lock() +# Held only around an OAuth2 refresh — see _connect. +_refresh_lock = threading.Lock() def _cached_client(user_id): @@ -483,16 +517,45 @@ def _connect(creds, user_id=None): # 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 + # One refresh at a time per process. Two threads reaching a cold + # client cache together would otherwise each mint a token against + # the endpoint that rate-limits hardest. + with _refresh_lock: + oauth2 = getattr(client.garth, "oauth2_token", None) + stored = load_token(user_id) if user_id else None + if stored and stored != token: + # Another thread refreshed while we waited: adopt its + # token instead of spending an SSO call of our own. + client.garth.loads(stored) + 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, named better + if _is_rate_limited(e): + _note_rate_limit(user_id) + raise RateLimited( + "Garmin 暂时限制了请求频率。这通常是短时间内连接过于频繁," + "等待约半小时后会自动恢复,令牌本身没有失效。" + ) from e + raise + # Keep it. This line is the whole reason the account kept + # getting throttled: the refreshed OAuth2 token used to + # live only in this process's memory, so every cold client + # cache — a 15-minute timeout, a second gunicorn worker, a + # deploy restart — re-read the same expired token from the + # database and minted another one against Garmin's SSO + # endpoint. That endpoint limits per account and blocks for + # hours. Persisting it turns "an SSO call per connect" into + # "one per token lifetime" (about an hour). + if user_id: + try: + _persist_token(user_id, client.garth.dumps()) + except Exception as e: # noqa: BLE001 - never fail a sync + logging.getLogger(__name__).warning( + "refreshed Garmin token failed to persist " + "for %s: %s", user_id, e + ) # 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"] @@ -1051,6 +1114,12 @@ ALWAYS_REFETCH_DAYS = 3 # per day would turn a year's sync into an hour. SERIES_INLINE_DAYS = 14 +# Activity details fetched per sync run. One request each, and a fresh account +# has hundreds — fetching them all in one run is the single biggest burst a +# sync can produce. The rest are picked up by the next run and by +# start_backfill, so nothing is lost by spreading them out. +DETAILS_PER_SYNC = int(os.environ.get("GARMIN_DETAILS_PER_SYNC") or 40) + # 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. @@ -1136,6 +1205,14 @@ def sync_data(user_id, creds, days=None, client=None, trigger="manual"): "lastSyncTime": now, }) + # Every Garmin call from here on is spaced and counted. Applied here + # rather than in `_connect` so an injected client is paced identically — + # the pacing is a property of a sync run, not of how the session was + # obtained. See services/garmin_throttle.py for why this is necessary at + # all: a day of history costs nine requests, and a full backfill used to + # fire six thousand of them flat out. + client = throttle.pace(client, user_id) + today = datetime.date.today() # -1 means "incremental sync": pick up from the latest date already in the @@ -1177,6 +1254,7 @@ def sync_data(user_id, creds, days=None, client=None, trigger="manual"): reached_the_start = False day_errors = [] rate_limited_mid_run = None + budget_spent = False for i in range(days): date_str = (today - datetime.timedelta(days=i)).isoformat() if i >= ALWAYS_REFETCH_DAYS and date_str in stored: @@ -1185,6 +1263,12 @@ def sync_data(user_id, creds, days=None, client=None, trigger="manual"): try: record = _extract_daily(client, date_str) record.update(extras.daily_extras(client, date_str)) + except throttle.BudgetExhausted: + # Not a failure: this run has issued as many requests as it is + # allowed. Keep what is stored and stop — the next run skips the + # days already written and carries on from here. + budget_spent = True + break except Exception as e: # A 429 inside the day loop used to be filed as one more skipped # day, so a 730-day backfill kept hammering Garmin for another @@ -1252,9 +1336,12 @@ def sync_data(user_id, creds, days=None, client=None, trigger="manual"): progress_current=days, progress_total=days, stage="运动记录") try: - activities_synced = _sync_activities( - client, user_id, start_date, today.isoformat() - ) + if not budget_spent: + activities_synced = _sync_activities( + client, user_id, start_date, today.isoformat() + ) + except throttle.BudgetExhausted: + budget_spent = True except Exception as e: day_errors.append(f"activities: {describe(e)}") @@ -1262,13 +1349,19 @@ def sync_data(user_id, creds, days=None, client=None, trigger="manual"): # 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}"), - ) + if not budget_spent: + # Capped per run. On a fresh account this is one request for each + # of a couple of hundred activities; unbounded, it spends the whole + # budget (and used to spend the account's goodwill) in one burst. + details_synced = sync_activity_details( + client, user_id, limit=DETAILS_PER_SYNC, + 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 throttle.BudgetExhausted: + budget_spent = True except Exception as e: day_errors.append(f"activity_details: {describe(e)}") @@ -1335,12 +1428,19 @@ def sync_data(user_id, creds, days=None, client=None, trigger="manual"): ) if day_errors: message += f"({len(day_errors)} 项跳过)" + if budget_spent: + # Said plainly: an unexplained "only 120 of 730 days" reads as a bug, + # and the next run picking up where this one stopped is the design. + message += ";本次已达请求预算上限,剩余部分下次同步继续" + requests_used = getattr(getattr(client, "limiter", None), "used", None) return finish({ "status": "success", "recordsSynced": days_synced, "activitiesSynced": activities_synced, "badgesSynced": badges_synced, "personalRecordsSynced": records_synced_pr, + "requestsUsed": requests_used, + "partial": budget_spent, "message": message, "lastSyncTime": now, }) diff --git a/backend/services/garmin_throttle.py b/backend/services/garmin_throttle.py new file mode 100644 index 0000000..fd4b311 --- /dev/null +++ b/backend/services/garmin_throttle.py @@ -0,0 +1,202 @@ +""" +Client-side pacing for Garmin Connect. + +Why this exists +--------------- +A sync is not a handful of requests. Each day costs **nine** API calls +(`_extract_daily` makes six, `daily_extras` three), short syncs add five more +for the within-day curves, every activity without a stored detail costs one, +and five account-wide calls finish the run. So: + + 7 days ~ 100 requests + 30 days ~ 300 requests + 730 days ~ 6600 requests + +and until now they were fired as fast as the network allowed. That is what +walks the account into a 429 — repeatedly, because the only defence was a 24h +cooldown *after* Garmin had already punished us. Backing off after the fact +does not stop it happening again; spacing the requests does. + +Garmin publishes no limit for these endpoints, so the default comes from what +other consumers actually get away with (checked 2026-09-03): + +* `sirredbeard/garmin-data-export` paces at **0.15s** by default, backs off on + a 429, ramps back up, and pauses every 250 calls. Its README suggests + `--delay 1.0` "if you want to play it safe". +* `evg656e/garmin-connect-scraper` defaults to **3s**, warning that smaller + values "may result in timeout penalties from Garmin". +* Garmin's *official* (partner) API allows **100 requests/minute** — 0.6s + each — which is the closest thing to a number Garmin itself stands behind. + +0.5s sits between them: slower than the official allowance, three times more +cautious than the aggressive scraper, six times faster than the timid one. + +Worth being clear about what this does and does not fix. The brutal +multi-hour blocks reported across those projects are on the **SSO/login** +endpoint and are keyed per account; this app's own lockout came from there, +not from data volume (it was re-minting a token on every connect — see +`_connect` in garmin.py). Pacing the data calls is the other half: a full +backfill firing 6600 requests flat out is not something to rely on Garmin +tolerating. + +What it does +------------ +* **Spacing.** Every call waits until at least `min_interval` has passed since + the previous one, per account. +* **A per-run budget.** A run may issue at most N requests; the sync then stops + cleanly and resumes next time (it already skips days it has stored). Without + this, "全部历史" is a single unbounded burst no matter how well spaced. + +Both apply at one choke point — a proxy around the client object — so every +call site is covered, including ones added later. Anything that reaches Garmin +without going through the client is not paced, which today means only the +login/token calls on `client.garth`: a couple per sync, not thousands. +""" +import functools +import os +import threading +import time + + +class BudgetExhausted(Exception): + """This run has issued as many requests as it is allowed. + + Not an error: the caller keeps what it stored and continues next run. + """ + + +def min_interval(): + """Seconds between requests. Read per call so it is tunable without a + restart, and so tests can set it to 0.""" + return float(os.environ.get("GARMIN_MIN_INTERVAL_SECONDS") or 0.5) + + +def default_budget(): + """Requests one sync run may issue. + + 1200 at 0.5s is about 10 minutes of wall clock — long enough to cover four + months of history in one go (9 requests per day), short enough that a full + backfill is a handful of background runs rather than one long hammering. + """ + return int(os.environ.get("GARMIN_REQUEST_BUDGET") or 1200) + + +class Limiter: + """Per-account spacing and per-run request accounting.""" + + def __init__(self, interval=None): + self._interval = interval + self._lock = threading.Lock() + self._next_at = 0.0 + self.used = 0 + self.budget = None + + @property + def interval(self): + return min_interval() if self._interval is None else self._interval + + def start_run(self, budget=None): + """Begin a run. `budget` None means unlimited (a short, known-size job).""" + with self._lock: + self.used = 0 + self.budget = budget + + @property + def remaining(self): + if self.budget is None: + return None + return max(0, self.budget - self.used) + + def acquire(self): + """Claim a slot, waiting for it if necessary.""" + with self._lock: + if self.budget is not None and self.used >= self.budget: + raise BudgetExhausted( + f"本次同步已达请求预算上限({self.budget} 次)" + ) + self.used += 1 + now = time.monotonic() + wait = self._next_at - now + # Reserve this call's slot before releasing the lock, so two + # threads get consecutive slots rather than the same one. + self._next_at = max(now, self._next_at) + self.interval + # Slept outside the lock: two callers should wait out their own + # (non-overlapping) slots concurrently, not queue behind each other's. + if wait > 0: + time.sleep(wait) + + +_limiters = {} +_limiters_lock = threading.Lock() + + +def limiter_for(user_id): + """The limiter for one account, created on first use. + + Per account and per process. Two Gunicorn workers therefore pace + independently — acceptable because a given account syncs from one place at + a time (the scheduler claims its tick through the database, and a manual + sync sets the status to `syncing`), and the budget bounds the damage if + that ever stops holding. + """ + key = user_id or "_anon" + with _limiters_lock: + found = _limiters.get(key) + if found is None: + found = Limiter() + _limiters[key] = found + return found + + +class PacedClient: + """Attribute proxy that paces every public method call. + + Non-callables (`display_name`, `garth`) pass through untouched, as do + private names, and attribute writes reach the wrapped client — `_connect` + sets `client.display_name` and the library reads it back on every request. + """ + + def __init__(self, inner, limiter): + object.__setattr__(self, "_inner", inner) + object.__setattr__(self, "_limiter", limiter) + + @property + def limiter(self): + return object.__getattribute__(self, "_limiter") + + @property + def inner(self): + return object.__getattribute__(self, "_inner") + + def __getattr__(self, name): + attr = getattr(object.__getattribute__(self, "_inner"), name) + if name.startswith("_") or not callable(attr): + return attr + limiter = object.__getattribute__(self, "_limiter") + + @functools.wraps(attr) + def paced(*args, **kwargs): + limiter.acquire() + return attr(*args, **kwargs) + + return paced + + def __setattr__(self, name, value): + setattr(object.__getattribute__(self, "_inner"), name, value) + + def __repr__(self): + inner = object.__getattribute__(self, "_inner") + return f"" + + +def pace(client, user_id, budget=None): + """Wrap `client` so its calls are spaced, and start a run. + + Idempotent: wrapping an already-wrapped client returns it (with the run + restarted) rather than stacking two layers of waiting. + """ + limiter = limiter_for(user_id) + limiter.start_run(default_budget() if budget is None else budget) + if isinstance(client, PacedClient): + return client + return PacedClient(client, limiter) diff --git a/backend/tests/conftest.py b/backend/tests/conftest.py index 784ea75..81322f1 100644 --- a/backend/tests/conftest.py +++ b/backend/tests/conftest.py @@ -31,6 +31,11 @@ from auth import sign_token # noqa: E402 # process — a developer's real AI_MODEL_CHAIN or API keys would silently change # what the suite exercises (and could bill real API calls). Clear them here; # individual tests opt back in through the `keys` / `gateway` fixtures. +# Pacing is real time: at the default 0.5s a single 30-day sync test would +# sleep for over two minutes. Tests exercise the *accounting* (budgets, counts) +# with the wait set to zero. +os.environ.setdefault("GARMIN_MIN_INTERVAL_SECONDS", "0") + _AI_ENV_VARS = ( "AI_MODEL_CHAIN", "AI_DAY_BUDGET", diff --git a/backend/tests/test_garmin_sync.py b/backend/tests/test_garmin_sync.py index ddab01e..1ba784a 100644 --- a/backend/tests/test_garmin_sync.py +++ b/backend/tests/test_garmin_sync.py @@ -999,3 +999,132 @@ class TestSyncHistory: r = client.get("/api/garmin/sync-history", headers=auth) assert r.status_code == 200 assert r.get_json() == {"items": []} + + +class TestRefreshedTokenIsKept: + """The account's repeated lockouts came from here. + + `refresh_oauth2()` mints a token against Garmin's SSO endpoint — the one + that limits per account and blocks for hours. The refreshed token used to + live only in the process's memory, so every cold client cache (a 15-minute + timeout, a second gunicorn worker, a deploy restart) re-read the same + expired token from the database and minted another one. + """ + + @staticmethod + def _garmin(refreshes, expired_after_load=True): + class StubOAuth2: + def __init__(self, expired): + self.expired = expired + + class StubGarth: + def __init__(self): + self.oauth2_token = None + self.profile = {"displayName": "Tester"} + self.dumped = "refreshed-token" + + def configure(self, **kwargs): + pass + + def loads(self, token): + self.oauth2_token = StubOAuth2( + expired_after_load if token == "stored-token" else False + ) + + def refresh_oauth2(self): + refreshes.append(1) + self.oauth2_token = StubOAuth2(False) + + def dumps(self): + return self.dumped + + class StubGarmin: + def __init__(self, is_cn=False): + self.garth = StubGarth() + self.display_name = None + + return StubGarmin + + def test_the_refreshed_token_is_written_back(self, db, user, monkeypatch): + refreshes = [] + monkeypatch.setattr(garmin_svc, "_import_garmin", + lambda: self._garmin(refreshes)) + garmin_svc.save_token(user["id"], "stored-token", "a@example.com") + + garmin_svc._connect({}, user["id"]) + + assert refreshes == [1] + assert garmin_svc.load_token(user["id"]) == "refreshed-token", ( + "a refresh that is not persisted makes the next connect mint " + "another token against the endpoint that blocks accounts" + ) + + def test_persisting_a_refresh_keeps_the_bound_email(self, db, user, monkeypatch): + refreshes = [] + monkeypatch.setattr(garmin_svc, "_import_garmin", + lambda: self._garmin(refreshes)) + garmin_svc.save_token(user["id"], "stored-token", "a@example.com") + garmin_svc.forget_client(user["id"]) + + garmin_svc._connect({}, user["id"]) + + assert garmin_svc.get_remembered_email(user["id"]) == "a@example.com", ( + "the 数据同步 page shows this; a refresh must not blank it" + ) + + def test_a_cold_cache_does_not_refresh_again(self, db, user, monkeypatch): + """The whole point: the second connect reads a *valid* stored token.""" + refreshes = [] + monkeypatch.setattr(garmin_svc, "_import_garmin", + lambda: self._garmin(refreshes)) + garmin_svc.save_token(user["id"], "stored-token", "a@example.com") + + garmin_svc._connect({}, user["id"]) + garmin_svc.forget_client(user["id"]) # as a restart or TTL would + garmin_svc._connect({}, user["id"]) + + assert refreshes == [1], f"minted {len(refreshes)} SSO tokens, expected 1" + + def test_persisting_the_refresh_does_not_drop_the_live_session( + self, db, user, monkeypatch + ): + """`save_token` invalidates the cached client; the refresh path must + not, or every connect would throw its own session away.""" + refreshes = [] + monkeypatch.setattr(garmin_svc, "_import_garmin", + lambda: self._garmin(refreshes)) + garmin_svc.save_token(user["id"], "stored-token", "a@example.com") + + first = garmin_svc._connect({}, user["id"]) + second = garmin_svc._connect({}, user["id"]) + assert second is first, "the session was cached, not rebuilt" + + def test_a_rate_limited_refresh_still_surfaces_as_rate_limited( + self, db, user, monkeypatch + ): + class StubGarth: + oauth2_token = None + profile = {"displayName": "Tester"} + + def configure(self, **kwargs): + pass + + def loads(self, token): + pass + + def refresh_oauth2(self): + raise Exception("429 Client Error: Too Many Requests") + + def dumps(self): + return "unused" + + class StubGarmin: + def __init__(self, is_cn=False): + self.garth = StubGarth() + + monkeypatch.setattr(garmin_svc, "_import_garmin", lambda: StubGarmin) + garmin_svc.save_token(user["id"], "stored-token", "a@example.com") + + with pytest.raises(garmin_svc.RateLimited): + garmin_svc._connect({}, user["id"]) + assert garmin_svc.rate_limited_until(user["id"]) is not None diff --git a/backend/tests/test_garmin_throttle.py b/backend/tests/test_garmin_throttle.py new file mode 100644 index 0000000..34258e4 --- /dev/null +++ b/backend/tests/test_garmin_throttle.py @@ -0,0 +1,122 @@ +""" +Unit tests for the Garmin request pacer. + +Nothing here touches the network, and the interval is zero (see conftest) so +the accounting is tested without the waiting. +""" +import time + +import pytest + +from services import garmin_throttle as throttle + + +class Recorder: + """Stands in for a garminconnect client.""" + + def __init__(self): + self.calls = [] + self.display_name = "Tester" + self.garth = object() + + def get_user_summary(self, date): + self.calls.append(date) + return {"date": date} + + def get_sleep_data(self, date): + self.calls.append(date) + return {} + + +class TestPacedClient: + def test_calls_reach_the_wrapped_client(self): + inner = Recorder() + client = throttle.pace(inner, "u") + assert client.get_user_summary("2026-09-01") == {"date": "2026-09-01"} + assert inner.calls == ["2026-09-01"] + + def test_every_call_is_counted(self): + client = throttle.pace(Recorder(), "u") + client.get_user_summary("d") + client.get_sleep_data("d") + assert client.limiter.used == 2 + + def test_non_callables_pass_through(self): + """`display_name` is read by the library on every request, and `garth` + carries the login flow — neither is an API call.""" + inner = Recorder() + client = throttle.pace(inner, "u") + assert client.display_name == "Tester" + assert client.garth is inner.garth + assert client.limiter.used == 0 + + def test_attribute_writes_reach_the_client(self): + inner = Recorder() + client = throttle.pace(inner, "u") + client.display_name = "Someone" + assert inner.display_name == "Someone" + + def test_wrapping_twice_does_not_stack_two_waits(self): + once = throttle.pace(Recorder(), "u") + twice = throttle.pace(once, "u") + assert twice is once + + def test_a_missing_method_still_raises_attribute_error(self): + client = throttle.pace(Recorder(), "u") + with pytest.raises(AttributeError): + client.get_something_that_does_not_exist + + +class TestBudget: + def test_the_budget_stops_the_run(self): + client = throttle.pace(Recorder(), "u", budget=3) + for _ in range(3): + client.get_sleep_data("d") + with pytest.raises(throttle.BudgetExhausted): + client.get_sleep_data("d") + + def test_remaining_counts_down(self): + client = throttle.pace(Recorder(), "u", budget=5) + client.get_sleep_data("d") + assert client.limiter.remaining == 4 + + def test_a_new_run_resets_the_budget(self): + client = throttle.pace(Recorder(), "u", budget=2) + client.get_sleep_data("d") + client.get_sleep_data("d") + throttle.pace(client, "u", budget=2) + client.get_sleep_data("d") # must not raise + assert client.limiter.used == 1 + + def test_no_budget_means_unlimited(self): + client = throttle.pace(Recorder(), "u", budget=0) + # 0 is falsy but explicit; only None means unlimited. + with pytest.raises(throttle.BudgetExhausted): + client.get_sleep_data("d") + limiter = throttle.limiter_for("u") + limiter.start_run(None) + client.get_sleep_data("d") + assert limiter.remaining is None + + +class TestSpacing: + def test_calls_are_spaced_by_the_interval(self, monkeypatch): + monkeypatch.setenv("GARMIN_MIN_INTERVAL_SECONDS", "0.05") + client = throttle.pace(Recorder(), f"spacing-{time.monotonic()}") + started = time.monotonic() + for _ in range(4): + client.get_sleep_data("d") + # Three gaps between four calls; the first claims a free slot. + assert time.monotonic() - started >= 0.05 * 3 + + def test_the_interval_is_read_per_call(self, monkeypatch): + monkeypatch.setenv("GARMIN_MIN_INTERVAL_SECONDS", "2.5") + assert throttle.min_interval() == 2.5 + monkeypatch.delenv("GARMIN_MIN_INTERVAL_SECONDS") + assert throttle.min_interval() == 0.5, "the researched default" + + def test_accounts_are_paced_independently(self): + a = throttle.limiter_for("account-a") + b = throttle.limiter_for("account-b") + assert a is not b + assert throttle.limiter_for("account-a") is a