diff --git a/eufy_sync/__init__.py b/eufy_sync/__init__.py index b5a26ca..3600946 100644 --- a/eufy_sync/__init__.py +++ b/eufy_sync/__init__.py @@ -1,6 +1,6 @@ """Sync Eufy smart scale body composition data to Garmin Connect and Strava.""" -__version__ = "1.14.0" +__version__ = "1.15.0" # Public API for programmatic use from eufy_sync.eufy_client import EufyClient, EufyMeasurement diff --git a/eufy_sync/cli/doctor.py b/eufy_sync/cli/doctor.py index 3bbffcc..49053be 100644 --- a/eufy_sync/cli/doctor.py +++ b/eufy_sync/cli/doctor.py @@ -302,7 +302,7 @@ def _check_version(report) -> None: if latest is None: report("WARN", "version", "could not check") return - if latest == __version__: + if not updater.is_newer(latest, __version__): report("PASS", "version", f"{__version__} (up to date)") return report( diff --git a/eufy_sync/cli/status.py b/eufy_sync/cli/status.py index 77626e5..cc03c4f 100644 --- a/eufy_sync/cli/status.py +++ b/eufy_sync/cli/status.py @@ -101,6 +101,21 @@ def _print_summary( ) +def _retry_queue_line(state, user) -> str | None: + """'2 uploads waiting to retry (Garmin)', or None when nothing waits. + Targets no longer configured are left out: nothing will retry them.""" + waiting = { + target: count + for target, count in state.waiting_upload_retries(user.name).items() + if getattr(user, target, None) is not None + } + total = sum(waiting.values()) + if not total: + return None + names = ", ".join(target.capitalize() for target in sorted(waiting)) + return f"{total} upload{'' if total == 1 else 's'} waiting to retry ({names})" + + def _show_status(state, users: list) -> None: """Print detailed sync status for all users.""" for user in users: @@ -118,6 +133,10 @@ def _show_status(state, users: list) -> None: else: print("Last synced measurement: never") + retry_line = _retry_queue_line(state, user) + if retry_line: + print(retry_line) + # Eufy token health from eufy_sync.eufy_client import EufyClient eufy_status = EufyClient(user.eufy).token_status() diff --git a/eufy_sync/cli/updater.py b/eufy_sync/cli/updater.py index a977d7c..81efc06 100644 --- a/eufy_sync/cli/updater.py +++ b/eufy_sync/cli/updater.py @@ -13,6 +13,23 @@ from eufy_sync.cli import shared +def _parse(v: str) -> tuple: + # Compare numeric prefix only - tolerates suffixes like "1.7.2rc1" or "1.7.2.dev0". + if not v or len(v) > 64: + raise ValueError(f"Implausible version string: {v!r}") + match = re.match(r"^\d+(?:\.\d+)*", v) + if not match: + raise ValueError(f"Implausible version string: {v!r}") + return tuple(int(x) for x in match.group(0).split(".")) + + +def is_newer(latest: str, current: str) -> bool: + """True when latest is a strictly newer release than current. An install + ahead of PyPI (a fresh release the index hasn't caught up with, or a dev + checkout) is never offered a downgrade.""" + return _parse(latest) > _parse(current) + + def _latest_pypi_version() -> str | None: """Return the latest eufy-sync version on PyPI, or None if unreachable.""" try: @@ -44,15 +61,6 @@ def _check_for_updates() -> None: from eufy_sync import __version__ - def _parse(v: str) -> tuple: - # Compare numeric prefix only - tolerates suffixes like "1.7.2rc1" or "1.7.2.dev0". - if not v or len(v) > 64: - raise ValueError(f"Implausible version string: {v!r}") - match = re.match(r"^\d+(?:\.\d+)*", v) - if not match: - raise ValueError(f"Implausible version string: {v!r}") - return tuple(int(x) for x in match.group(0).split(".")) - latest_parsed = _parse(latest) current_parsed = _parse(__version__) @@ -91,12 +99,12 @@ def _self_update() -> None: if latest is None: print("Could not reach PyPI. Check your connection and try again.") return - if latest == __version__: - print(f"Already on the latest version (v{__version__}).") - return if not re.match(r"^\d+(?:\.\d+)*", latest): print(f"Unexpected version from PyPI ({latest!r}); update manually with pipx.") return + if not is_newer(latest, __version__): + print(f"Already on the latest version (v{__version__}).") + return # A uv or pipx reinstall replaces the whole tool venv, so the pin has to # carry the browser extra along for anyone who opted into it, or the diff --git a/eufy_sync/garmin_client.py b/eufy_sync/garmin_client.py index 160d653..0f29dc0 100644 --- a/eufy_sync/garmin_client.py +++ b/eufy_sync/garmin_client.py @@ -1,12 +1,25 @@ """Garmin Connect client. Delegates login, refresh, and upload to python-garminconnect; keeps a same-date duplicate check. + +The library logs in through curl_cffi but sends every data call through plain +requests, and on some networks (VPNs, datacenter IPs) Cloudflare refuses that +TLS fingerprint with a 403 even though the token is fine (upstream issue #444). +A call refused that way is sent once more through curl_cffi with a browser +fingerprint, reusing the session's token, before anything here treats the +session as dead. """ from __future__ import annotations import logging +import re +import time from datetime import datetime, timezone -from garminconnect import GarminConnectAuthenticationError, GarminConnectConnectionError +from garminconnect import ( + GarminConnectAuthenticationError, + GarminConnectConnectionError, + GarminConnectTooManyRequestsError, +) from eufy_sync.config import GarminConfig from eufy_sync.garmin_auth import GarminAuth @@ -21,6 +34,10 @@ # far short of a separate weigh-in. _WEIGHT_TOLERANCE_KG = 0.1 _TIMESTAMP_TOLERANCE_SECONDS = 120 +# The lookup that confirms a 409 is retried on its own, a few times with a +# short pause, rather than by sync's _retry, which would resend the upload. +_CONFIRM_LOOKUP_ATTEMPTS = 3 +_CONFIRM_LOOKUP_BACKOFF_SECONDS = 2 def _entry_instants(entry: dict) -> list[datetime]: @@ -50,38 +67,218 @@ def _match_uploaded_entry(entries: list[dict], uploaded_at: datetime) -> dict | is not unique. Weight alone is not enough: a manual weigh-in on the same day within the weight window would match too, and deleting it would throw away data eufy-sync never created. Entries that carry a timestamp must - therefore match on it. Only when Garmin returns no timestamps at all does - a single weight match stand on its own - one response carries the same - fields for every entry, so the two cases do not mix in practice.""" - timestamped = [e for e in entries if _entry_instants(e)] - if timestamped: - matches = [ - e for e in timestamped - if any( - abs((instant - uploaded_at).total_seconds()) <= _TIMESTAMP_TOLERANCE_SECONDS - for instant in _entry_instants(e) - ) - ] - else: - matches = entries + therefore match on it. + + Only delete_weight_entry uses this, and it deliberately keeps one + permissive case: when Garmin returns no usable timestamps at all, a single + weight match stands on its own. That case replaces our own weight-only + upload with the full record for the same weigh-in, two entries in the + window leave everything alone, and a wrong guess costs one entry the next + upload restores. One response carries the same fields for every entry, so + the two cases do not mix in practice. Confirming a 409 does not get this + leeway (see _timed_entries_at).""" + matches = _entries_at(entries, uploaded_at) return matches[0] if len(matches) == 1 else None +def _entries_at(entries: list[dict], uploaded_at: datetime) -> list[dict]: + """The entries whose own timestamp sits within the tolerance of + uploaded_at, or all of them when Garmin sent no timestamps. The fallback + exists for delete_weight_entry only; see _match_uploaded_entry.""" + timestamped = [e for e in entries if _entry_instants(e)] + if not timestamped: + return list(entries) + return _timed_entries_at(timestamped, uploaded_at) + + +def _timed_entries_at(entries: list[dict], uploaded_at: datetime) -> list[dict]: + """The entries with a parseable timestamp within the tolerance of + uploaded_at. An entry without one never qualifies: confirming a 409 marks + the measurement uploaded for good, so a same-weight entry from some other + time must not stand in for ours.""" + return [ + e for e in entries + if any( + abs((instant - uploaded_at).total_seconds()) <= _TIMESTAMP_TOLERANCE_SECONDS + for instant in _entry_instants(e) + ) + ] + + +# The library raises HTTP failures with the status only in the message: +# "API Error 403 - ..." from a direct call, "API call client error (403): ..." +# once its error decorator has rewrapped it. +_STATUS_RE = re.compile(r"(?:API Error|error \(|HTTP)\s*(\d{3})\b") + +# Text that only a Cloudflare challenge carries: the "Just a moment" page and +# its challenge-platform script. The API host sits behind Cloudflare, so cf-ray, +# "server: cloudflare", and the "Cloudflare Ray ID" footer appear on ordinary +# answers and on gateway error pages too, and prove nothing on their own. +_CHALLENGE_MARKERS = ( + "just a moment", + "cf-chl", + "challenge-platform", +) +# Text on Cloudflare's firewall block page ("Sorry, you have been blocked", +# error 1020). Its cf-error-details box also sits on 5xx gateway pages, so it +# counts only with a 403. +_BLOCK_PAGE_MARKERS = ( + "attention required", + "you have been blocked", +) +# Statuses a Cloudflare challenge is served with. Any other 5xx is a gateway +# or origin failure, where the request may well have reached Garmin. +_CHALLENGE_STATUSES = (403, 429, 503) + +# The browser fingerprint for the fallback. Issue #444's reporter got 200s +# with this one, the library's native headers, and the same token that plain +# requests could not use. +_IMPERSONATE = "chrome" + + +def _status_code(exc: BaseException) -> int | None: + """The HTTP status behind a library error, or None for a non-HTTP one.""" + status = getattr(getattr(exc, "response", None), "status_code", None) + if isinstance(status, int): + return status + match = _STATUS_RE.search(str(exc)) + return int(match.group(1)) if match else None + + def _is_garmin_auth_failure(exc: Exception) -> bool: """True when a Garmin call failed because the session may be dead. Covers the dedicated auth error and the 401/403 that the library reports as a generic connection error ("API Error 401 - ..."). - A Cloudflare 403 matches too: the library drops the response body before - raising, so the message is the same "API Error 403" either way. That is - safe only because a relogin keeps the stored token until the new login - succeeds, so a 403 that was only a passing block costs one login attempt, - not the session.""" + A 403 stays ambiguous even with the response in hand: the network block + in issue #444 answers with the same JSON ForbiddenException a refused + token gets. By the time a 403 reaches the relogin, the browser-fingerprint + retry has already failed, and a relogin keeps the stored token until the + new login succeeds, so a 403 that was only a passing block costs one + login attempt, not the session.""" if isinstance(exc, GarminConnectAuthenticationError): return True - return isinstance(exc, GarminConnectConnectionError) and ( - "401" in str(exc) or "403" in str(exc) - ) + return isinstance(exc, GarminConnectConnectionError) and _status_code(exc) in (401, 403) + + +class _LastResponse: + """The status, headers, and start of the body of the most recent failed + response on the library's API session or on the fallback. + + The library drops the response before raising, so this is the only place a + Cloudflare block page and a JSON 403 from the API can be told apart.""" + + _BODY_LIMIT = 4096 + + def __init__(self): + self.clear() + + def clear(self) -> None: + self.status: int | None = None + self.headers: dict[str, str] = {} + self.body = "" + + def record(self, resp) -> None: + status = getattr(resp, "status_code", None) + if not isinstance(status, int) or status < 400: + return + self.status = status + try: + self.headers = {str(k).lower(): str(v) for k, v in resp.headers.items()} + except Exception: + self.headers = {} + try: + self.body = (resp.text or "")[: self._BODY_LIMIT] + except Exception: + self.body = "" + + def hook(self, resp, *args, **kwargs) -> None: + """requests response hook. Returns None so the response is unchanged.""" + self.record(resp) + + def is_cloudflare_block(self) -> bool: + """True only on affirmative evidence that Cloudflare stopped the + request before Garmin saw it: the cf-mitigated header, or a challenge + or block page with a status Cloudflare serves those with. A branded + 502/504 gateway page is not a block; the request may have reached the + origin, so it is left to the ordinary retry policy.""" + status = self.status + if status is None or status not in _CHALLENGE_STATUSES: + return False + if self.headers.get("cf-mitigated", "").lower() == "challenge": + return True + if "html" not in self.headers.get("content-type", "").lower(): + return False + body = self.body.lower() + if any(marker in body for marker in _CHALLENGE_MARKERS): + return True + return status == 403 and any(marker in body for marker in _BLOCK_PAGE_MARKERS) + + +def _new_impersonating_session(): + """A curl_cffi session with a browser TLS fingerprint. Kept separate so + tests can swap in a fake transport.""" + from curl_cffi import requests as cffi_requests + return cffi_requests.Session(impersonate=_IMPERSONATE) + + +class _ImpersonatingSession: + """Stands in for the library's requests.Session during one fallback call. + + The library still builds the URL, the auth headers, and its own error + handling; only the transport changes. curl_cffi refuses requests' files= + argument, so a file upload is rebuilt as a CurlMime multipart body, which + curl_cffi does support.""" + + def __init__(self, recorder: _LastResponse): + self._recorder = recorder + + def request(self, method, url, headers=None, files=None, **kwargs): + mime = _to_curl_mime(files) if files else None + if mime is not None: + kwargs["multipart"] = mime + sess = _new_impersonating_session() + try: + resp = sess.request(method, url, headers=headers, **kwargs) + finally: + if mime is not None: + mime.close() + sess.close() + self._recorder.record(resp) + return resp + + +def _to_curl_mime(files: dict): + """Rebuild a requests-style files= mapping as a curl_cffi multipart body. + + Accepts the shapes the library uses: {"file": (name, bytes_or_fileobj)} + with an optional third content-type element, or a bare bytes/file value.""" + from curl_cffi import CurlMime + mime = CurlMime() + try: + for field, value in files.items(): + filename, content, content_type = None, value, "application/octet-stream" + if isinstance(value, (tuple, list)): + filename, content = value[0], value[1] + if len(value) > 2 and value[2]: + content_type = value[2] + if hasattr(content, "read"): + content = content.read() + if isinstance(content, str): + content = content.encode() + mime.addpart(name=field, filename=filename, content_type=content_type, data=bytes(content)) + except Exception: + mime.close() + raise + return mime + + +class _Conflict: + """An upload Garmin answered with 409, handed out of the upload call so + the confirming lookup runs outside it.""" + + def __init__(self, error: Exception): + self.error = error class GarminClient: @@ -99,12 +296,33 @@ def __init__(self, config: GarminConfig): # of a Garmin 429. self._reauth_attempted = False self._reauth_error: Exception | None = None + self._last_response = _LastResponse() + # Set once a browser-fingerprint retry succeeds. The block is a + # property of the network, not the session, so from then on every call + # in this run goes straight through curl_cffi instead of paying a + # refused plain request first. A new GarminClient (a new run) starts + # over on plain requests. + self._impersonate_always = False def authenticate(self, allow_interactive: bool = True) -> None: self._allow_interactive = allow_interactive self._garmin = self._auth.login(interactive=allow_interactive) + self._watch_responses(self._garmin) logger.info("Authenticated to Garmin Connect as %s", self.config.email) + def _watch_responses(self, garmin) -> None: + """Hook the library's API session so a failed call's response can be + inspected after the library has raised without it. _api_session is + private but present in every release we allow (0.3.10 onward); without + it the hook is skipped and every 403 counts as ambiguous.""" + sess = getattr(getattr(garmin, "client", None), "_api_session", None) + hooks = getattr(sess, "hooks", None) + if not isinstance(hooks, dict): + return + response_hooks = hooks.setdefault("response", []) + if self._last_response.hook not in response_hooks: + response_hooks.append(self._last_response.hook) + def _reauth(self) -> None: """Replace a dead session with a fresh login, prompting when a person is present. A scheduled run has nobody to prompt, but a password login @@ -126,6 +344,71 @@ def _reauth(self) -> None: except Exception as e: self._reauth_error = e raise + self._watch_responses(self._garmin) + + def _attempt(self, call): + """Run call once, forgetting the response of any earlier failure so a + network error is never judged by a stale block page.""" + self._last_response.clear() + return call() + + def _call_impersonating(self, call): + """Run call once more with the library's API session swapped for a + curl_cffi one. The token, URL, and headers stay the library's.""" + client = self._garmin.client + original = client._api_session + client._api_session = _ImpersonatingSession(self._last_response) + try: + return self._attempt(call) + finally: + client._api_session = original + + def _can_impersonate(self) -> bool: + client = getattr(self._garmin, "client", None) + if getattr(client, "_api_session", None) is None: + return False + try: + import curl_cffi # noqa: F401 + except ImportError: + return False + return True + + def _refused_with_403(self, exc: Exception) -> bool: + """Whether the failed attempt was answered with a 403. The recorded + response wins over the error text: the call is replayed on this + answer, and an upload POST must never be replayed after a 5xx or a + timeout, where Garmin may already have stored it.""" + recorded = self._last_response.status + if recorded is not None: + return recorded == 403 + return _status_code(exc) == 403 + + def _call_with_fallback(self, call): + """Run call; when the network refused it, run it once more through a + browser fingerprint. A Cloudflare challenge or block qualifies, and so + does any 403, because the #444 block answers with a JSON 403 that + looks exactly like a refused token. Nothing else does: a 5xx or a + network failure is never replayed here, since the request may have + reached Garmin and an upload would be sent twice. Trying the fingerprint first costs one + request; a relogin costs a login and risks a 429, and on a blocked + network its own token check fails the same way.""" + if self._impersonate_always and self._can_impersonate(): + return self._call_impersonating(call) + try: + return self._attempt(call) + except (GarminConnectAuthenticationError, GarminConnectConnectionError) as e: + blocked = self._last_response.is_cloudflare_block() + if not (blocked or self._refused_with_403(e)) or not self._can_impersonate(): + raise + logger.info( + "Garmin refused the call (%s%s); retrying once with a browser fingerprint", + e, ", Cloudflare block page" if blocked else "", + ) + result = self._call_impersonating(call) + if not self._impersonate_always: + logger.info("Browser fingerprint got through; using it for the rest of this run") + self._impersonate_always = True + return result def _call_with_reauth(self, call): """Run a Garmin call, re-logging in once when the session is dead. @@ -133,11 +416,16 @@ def _call_with_reauth(self, call): The duplicate check is the run's first Garmin call, so a token that expired between runs used to fail here on every scheduled sync before upload's own healing could kick in. Non-auth errors, and anything the - relogin or the retry raises, travel to the caller unchanged.""" + relogin or the retry raises, travel to the caller unchanged. + + Each session gets one browser-fingerprint retry before its 403 counts + against it (see _call_with_fallback). A Cloudflare block page that + survives that retry never triggers a relogin: the token was not the + problem, and a new one would meet the same block.""" try: - return call() + return self._call_with_fallback(call) except (GarminConnectAuthenticationError, GarminConnectConnectionError) as e: - if not _is_garmin_auth_failure(e): + if not _is_garmin_auth_failure(e) or self._last_response.is_cloudflare_block(): raise if self._reauth_attempted: if self._reauth_error is not None: @@ -149,7 +437,7 @@ def _call_with_reauth(self, call): # would not change that. raise self._reauth() - return call() + return self._call_with_fallback(call) def check_connection(self) -> None: """Verify an authenticated read, propagating failures to diagnostics.""" @@ -240,15 +528,144 @@ def _add_body_composition(self, body_comp: GarminBodyComposition): ) def upload_body_composition(self, body_comp: GarminBodyComposition) -> dict: - # Transient connection errors propagate to _retry; only a dead session - # (token expired or revoked, seen as a 401) heals and retries here. - result = self._call_with_reauth(lambda: self._add_body_composition(body_comp)) + """Upload one body-composition FIT and sort out what the answer means. + + Modeled on scalebridge-sync's upload outcomes: + - 2xx: uploaded. A 409 counts too, but only once a lookup finds + the weigh-in on Garmin (see _confirm_duplicate). + - 401, or 403: a dead session or a blocked network. The fingerprint + retry and the run's one relogin heal what they can; a refusal that + outlives both raises PermanentSyncError, since _retry asking again + seconds later would meet the same answer. + - 429: GarminConnectTooManyRequestsError, which sync treats as + permanent for this run, as it does a login 429. Retrying a rate + limit in a loop only extends it. + - 408 and 5xx, and network failures: raised unchanged for _retry. + - Any other 4xx: Garmin rejected this upload itself; asking again + sends the same bytes. PermanentSyncError.""" + def upload(): + try: + return self._add_body_composition(body_comp) + except GarminConnectConnectionError as e: + if _status_code(e) == 409: + # Returned, not raised: the lookup below runs as its own + # call, so its recovery never replays this POST. + return _Conflict(e) + raise + + try: + result = self._call_with_reauth(upload) + except GarminConnectTooManyRequestsError: + raise + except (GarminConnectAuthenticationError, GarminConnectConnectionError) as e: + if e is not self._reauth_error: + # A failed relogin's own error already says what went wrong + # and travels unchanged; only Garmin's answer to the upload + # is sorted here. + self._classify_upload_failure(e) + raise + if isinstance(result, _Conflict): + return self._confirm_duplicate(body_comp, result.error) logger.info( "Uploaded body comp to Garmin: %.1f kg at %s", body_comp.weight, body_comp.timestamp, ) return result if isinstance(result, dict) else {"status": "ok"} + def _confirm_duplicate(self, body_comp: GarminBodyComposition, conflict: Exception) -> dict: + """Accept a 409 only when Garmin really holds this weigh-in: an entry + within the weight window whose own parseable timestamp is within the + time window of the instant we sent. An entry without a timestamp never + confirms; see _timed_entries_at. + + The lookup gets the same recovery as any other read (fingerprint + fallback, the sticky curl_cffi path, the run's one relogin), and + only it is repeated, never the upload. A 5xx or network failure is + retried here a few times with a short pause; if the answer is still + unknown, RetryNextRunError sends the measurement to the retry queue + without sync's _retry posting it again this run. A 429 ends Garmin + for this run like an upload 429 does. Only a lookup that answers and + lacks the weigh-in is PermanentSyncError.""" + from eufy_sync.sync import PermanentSyncError, RetryNextRunError + + uploaded_at = datetime.fromisoformat(body_comp.timestamp) + if uploaded_at.tzinfo is None: + uploaded_at = uploaded_at.astimezone() + date_str = uploaded_at.astimezone().strftime("%Y-%m-%d") + + def lookup() -> bool: + data = self._garmin.get_daily_weigh_ins(date_str) + near = [ + entry for entry in (data or {}).get("dateWeightList", []) + if abs(entry.get("weight", 0) / 1000.0 - body_comp.weight) <= _WEIGHT_TOLERANCE_KG + ] # Garmin stores grams + return bool(_timed_entries_at(near, uploaded_at)) + + for attempt in range(_CONFIRM_LOOKUP_ATTEMPTS): + try: + found = self._call_with_reauth(lookup) + break + except GarminConnectTooManyRequestsError: + raise + except (GarminConnectAuthenticationError, GarminConnectConnectionError, OSError) as e: + if e is self._reauth_error: + raise + status = _status_code(e) + if status == 429: + raise GarminConnectTooManyRequestsError( + f"Garmin rate limited the lookup confirming a 409 upload: {e}" + ) from e + if isinstance(e, GarminConnectAuthenticationError) or status in (401, 403): + # Refused after the fallback and the relogin, like an upload + # would be: the same advice applies. + self._classify_upload_failure(e) + raise + if attempt == _CONFIRM_LOOKUP_ATTEMPTS - 1: + logger.warning( + "Garmin answered the upload with 409 Conflict, and the lookup to confirm it " + "already holds the weigh-in failed %d times: %s", _CONFIRM_LOOKUP_ATTEMPTS, e, + ) + raise RetryNextRunError( + f"Garmin answered the upload with 409 Conflict, but the lookup to confirm " + f"it holds the weigh-in kept failing ({e}); trying again next run" + ) from e + delay = _CONFIRM_LOOKUP_BACKOFF_SECONDS * (2 ** attempt) + logger.info("Lookup confirming a 409 upload failed (%s); retrying in %ds", e, delay) + time.sleep(delay) + if not found: + raise PermanentSyncError( + f"Garmin answered the body comp upload with 409 Conflict ({conflict}), but has no " + f"{body_comp.weight:.1f} kg weigh-in at {body_comp.timestamp}; not counting it as uploaded" + ) from conflict + logger.info("Garmin already holds this weigh-in (409, confirmed by lookup); counting it as uploaded") + return {"status": "duplicate"} + + def _classify_upload_failure(self, exc: Exception) -> None: + """Raise the right error for an upload Garmin refused, or return to + let a transient failure travel to _retry unchanged.""" + from eufy_sync.sync import PermanentSyncError + + status = _status_code(exc) + if status == 429: + retry_after = self._last_response.headers.get("retry-after") + logger.warning( + "Garmin rate-limited the upload%s; leaving it for the next run", + f" (Retry-After: {retry_after}s)" if retry_after else "", + ) + raise GarminConnectTooManyRequestsError(f"Garmin upload rate limited: {exc}") from exc + if isinstance(exc, GarminConnectAuthenticationError) or status in (401, 403): + if self._last_response.is_cloudflare_block(): + raise PermanentSyncError( + f"Cloudflare is blocking Garmin uploads from this network (HTTP {status}). " + "A VPN or datacenter connection is the usual cause; try another network." + ) from exc + raise PermanentSyncError( + f"Garmin still refused the upload ({exc}) after a fresh login. " + "If you are on a VPN, try without it; otherwise run: eufy-sync --reauth garmin" + ) from exc + if status is not None and 400 <= status < 500 and status != 408: + raise PermanentSyncError(f"Garmin rejected the body comp upload (HTTP {status}): {exc}") from exc + def close(self) -> None: # Last chance to keep whatever the library rotated mid-run: its refresh # can hand back a new refresh token that lives in memory only. sync_user diff --git a/eufy_sync/state.py b/eufy_sync/state.py index e5a13d4..8185d1f 100644 --- a/eufy_sync/state.py +++ b/eufy_sync/state.py @@ -48,8 +48,33 @@ def _init_db(self) -> None: measurement_json TEXT NOT NULL, PRIMARY KEY(user_name, measurement_id) ); + CREATE TABLE IF NOT EXISTS upload_retries ( + user_name TEXT NOT NULL, + target TEXT NOT NULL, + measurement_id TEXT NOT NULL, + measurement_timestamp TEXT NOT NULL, + weight_kg REAL, + first_failed_at TEXT NOT NULL, + last_failed_at TEXT NOT NULL, + attempts INTEGER NOT NULL DEFAULT 1, + gave_up INTEGER NOT NULL DEFAULT 0, + last_newer_success_at TEXT, + PRIMARY KEY(user_name, target, measurement_id) + ); """) self._conn.commit() + self._migrate_retry_newer_success_column() + + def _migrate_retry_newer_success_column(self) -> None: + """Add last_newer_success_at to an upload_retries table created + before it existed. NULL (no newer upload seen yet) is the safe value + for old rows: it can only delay a give-up, never cause one.""" + cursor = self._conn.execute("PRAGMA table_info(upload_retries)") + columns = {row[1] for row in cursor.fetchall()} + if not columns or "last_newer_success_at" in columns: + return + with self._conn: + self._conn.execute("ALTER TABLE upload_retries ADD COLUMN last_newer_success_at TEXT") def _migrate_if_needed(self) -> None: """Migrate v1 schema (garmin-only) to v2 (multi-target), then v2 to v3 @@ -205,6 +230,126 @@ def clear_pending_upgrade(self, user_name: str, measurement_id: str) -> None: (user_name, measurement_id), ) + def record_upload_failure( + self, + user_name: str, + target: str, + measurement_id: str, + measurement_timestamp: str, + weight_kg: float, + failed_at: str, + ) -> int: + """Note a retryable upload failure and return its attempt count. + Holds only what sync_log already holds for a delivered measurement: + no credentials, no error text.""" + with self._conn: + self._conn.execute( + """INSERT INTO upload_retries + (user_name, target, measurement_id, measurement_timestamp, weight_kg, + first_failed_at, last_failed_at, attempts) + VALUES (?, ?, ?, ?, ?, ?, ?, 1) + ON CONFLICT(user_name, target, measurement_id) DO UPDATE SET + attempts = attempts + 1, last_failed_at = excluded.last_failed_at""", + (user_name, target, measurement_id, measurement_timestamp, weight_kg, failed_at, failed_at), + ) + row = self._conn.execute( + "SELECT attempts FROM upload_retries WHERE user_name = ? AND target = ? AND measurement_id = ?", + (user_name, target, measurement_id), + ).fetchone() + return row[0] + + def get_upload_retries(self, user_name: str) -> list[dict]: + cursor = self._conn.execute( + """SELECT target, measurement_id, measurement_timestamp, weight_kg, + first_failed_at, attempts, gave_up, last_newer_success_at + FROM upload_retries WHERE user_name = ?""", + (user_name,), + ) + return [ + { + "target": target, "measurement_id": mid, "measurement_timestamp": ts, + "weight_kg": kg, "first_failed_at": first_failed_at, + "attempts": attempts, "gave_up": bool(gave_up), + "last_newer_success_at": last_newer_success_at, + } + for target, mid, ts, kg, first_failed_at, attempts, gave_up, last_newer_success_at + in cursor.fetchall() + ] + + def note_newer_upload(self, user_name: str, target: str, timestamp: datetime, uploaded_at: str) -> None: + """Record that a measurement taken at timestamp reached the target, on + every waiting entry for an older measurement. The give-up rule needs + this: a target that keeps taking newer weigh-ins is up, so an old + entry that still fails is the measurement's problem, not Garmin's. + Compared in Python because the stored strings mix UTC offsets.""" + rows = self._conn.execute( + """SELECT measurement_id, measurement_timestamp FROM upload_retries + WHERE user_name = ? AND target = ? AND gave_up = 0""", + (user_name, target), + ).fetchall() + older = [mid for mid, ts in rows if datetime.fromisoformat(ts) < timestamp] + if older: + with self._conn: + self._conn.executemany( + """UPDATE upload_retries SET last_newer_success_at = ? + WHERE user_name = ? AND target = ? AND measurement_id = ?""", + [(uploaded_at, user_name, target, mid) for mid in older], + ) + + def get_oldest_waiting_retry_timestamp(self, user_name: str, target: str) -> int | None: + """Epoch seconds of the oldest measurement still due a retry to the + target, or None. Compared in Python because the stored strings mix + UTC offsets.""" + rows = self._conn.execute( + """SELECT measurement_timestamp FROM upload_retries + WHERE user_name = ? AND target = ? AND gave_up = 0""", + (user_name, target), + ).fetchall() + if not rows: + return None + return int(min(datetime.fromisoformat(ts).timestamp() for (ts,) in rows)) + + def waiting_upload_retries(self, user_name: str) -> dict[str, int]: + """Per-target count of failed uploads still due a retry.""" + cursor = self._conn.execute( + """SELECT target, COUNT(*) FROM upload_retries + WHERE user_name = ? AND gave_up = 0 GROUP BY target""", + (user_name,), + ) + return dict(cursor.fetchall()) + + def give_up_upload_retry(self, user_name: str, target: str, measurement_id: str) -> None: + """Stop retrying. The row stays as a marker: the fetch window can still + reach this measurement, and without it the next run would retry it.""" + with self._conn: + self._conn.execute( + "UPDATE upload_retries SET gave_up = 1 WHERE user_name = ? AND target = ? AND measurement_id = ?", + (user_name, target, measurement_id), + ) + + def clear_upload_retry(self, user_name: str, target: str, measurement_id: str) -> None: + with self._conn: + self._conn.execute( + "DELETE FROM upload_retries WHERE user_name = ? AND target = ? AND measurement_id = ?", + (user_name, target, measurement_id), + ) + + def clear_upload_retries_through(self, user_name: str, target: str, timestamp: datetime) -> None: + """Drop entries for measurements taken at or before timestamp. For a + target that holds only the current weight, a newer value replaces + them. Compared in Python because the stored strings mix UTC offsets.""" + rows = self._conn.execute( + "SELECT measurement_id, measurement_timestamp FROM upload_retries WHERE user_name = ? AND target = ?", + (user_name, target), + ).fetchall() + stale = [mid for mid, ts in rows if datetime.fromisoformat(ts) <= timestamp] + if stale: + with self._conn: + self._conn.executemany( + "DELETE FROM upload_retries WHERE user_name = ? AND target = ? AND measurement_id = ?", + [(user_name, target, mid) for mid in stale], + ) + def get_oldest_weight_only_timestamp( self, user_name: str, target: str, since: int | None = None, ) -> int | None: diff --git a/eufy_sync/sync.py b/eufy_sync/sync.py index 99026cd..4b45c0f 100644 --- a/eufy_sync/sync.py +++ b/eufy_sync/sync.py @@ -4,7 +4,7 @@ import logging import time from dataclasses import asdict -from datetime import datetime, timezone +from datetime import datetime, timedelta, timezone import httpx from garminconnect import GarminConnectTooManyRequestsError @@ -28,12 +28,41 @@ # their full record. One that has not matched within two weeks never will, # and an unbounded reach-back would refetch everything since it on every run. UPGRADE_LOOKBACK_DAYS = 14 +# A failed upload is retried because the fetch cursor stays behind it, so a +# measurement that never uploads would block every newer one to that target. +# Two weeks after its first failure (the same reach as upgrades), or after two +# weeks' worth of scheduled runs at one every 4 hours, a failing measurement is +# capped: it is still tried once each run, in order, but a retryable failure +# moves past it and newer ones keep uploading. The age counts from the first +# failure, not from when the weigh-in was taken, so an old measurement reached +# by a backfill gets the same two weeks. +RETRY_MAX_AGE_DAYS = UPGRADE_LOOKBACK_DAYS +MAX_RETRY_ATTEMPTS = 84 +# A capped measurement is given up only when it has failed for 30 days AND a +# newer weigh-in reached the same target after its first failure. A target +# that keeps taking newer weigh-ins is up, so one that still refuses this one +# after a month never will. A transient failure never gives anything up on its +# own, and during an outage nothing newer lands, so nothing is lost. +RETRY_ABANDON_DAYS = 30 +# How far past the abandon window the fetch reaches for a queued failure. An +# entry older than that whose target never took a newer weigh-in cannot be +# given up, and without a bound it would refetch all history since it on +# every run. It stays queued; --backfill-days still reaches it. +RETRY_REACH_BACK_MARGIN_DAYS = 2 class PermanentSyncError(RuntimeError): """Raised for failures that retries can't fix (bad password, revoked token).""" +class RetryNextRunError(RuntimeError): + """Retryable, but not within this run: _retry raises it straight through + and the measurement goes to the retry queue. Used when Garmin answered an + upload with 409 and the lookup confirming it kept failing. Replaying the + POST now would cost another upload for the same unanswered question; the + next run's POST 409s again and repeats the lookup.""" + + class UnsupportedMeasurementError(PermanentSyncError): """A target cannot accept this one measurement (e.g. outside its weight range). The measurement is skipped; the target itself stays healthy.""" @@ -59,7 +88,7 @@ def _retry(fn, description: str): try: return fn() except Exception as e: - if _is_permanent(e): + if _is_permanent(e) or isinstance(e, RetryNextRunError): raise if attempt == MAX_RETRIES - 1: raise @@ -69,6 +98,66 @@ def _retry(fn, description: str): time.sleep(delay) +def _triage_retries( + user_name: str, state: SyncState, target_names: list[str], dry_run: bool, +) -> tuple[set[tuple[str, str]], dict[tuple[str, str], bool]]: + """Tidy the retry queue for this run's targets. Returns (given_up, + capped): pairs of (target, measurement_id) never to upload again, and + pairs past a cap whose next failure must not stop their target. Each + capped pair maps to whether the target has taken a newer weigh-in since + it first failed; if so, the target is known to be up and another failure + of this one is not worth reporting as a target error. + + The queue does not replay anything itself: the per-target cursor already + re-fetches a failed measurement. The queue counts the attempts, so a + measurement that keeps failing stops holding back newer ones, and gives + one up under the rule at RETRY_ABANDON_DAYS. + """ + given_up: set[tuple[str, str]] = set() + capped: dict[tuple[str, str], bool] = {} + now = time.time() + oldest_allowed = now - RETRY_MAX_AGE_DAYS * 86400 + abandon_before = now - RETRY_ABANDON_DAYS * 86400 + for row in state.get_upload_retries(user_name): + target, mid = row["target"], row["measurement_id"] + if target not in target_names: + continue + if state.is_synced(user_name, mid, target): + if not dry_run: + state.clear_upload_retry(user_name, target, mid) + continue + taken_at = datetime.fromisoformat(row["measurement_timestamp"]).timestamp() + first_failed = datetime.fromisoformat(row["first_failed_at"]).timestamp() + if target in ("strava", "zwift"): + latest = state.get_latest_sync_timestamp(user_name, target) + if latest is not None and taken_at <= latest: + # A newer weight already reached it; this one is obsolete. + if not dry_run: + state.clear_upload_retry(user_name, target, mid) + continue + newer_success = row["last_newer_success_at"] + target_healthy = ( + newer_success is not None + and datetime.fromisoformat(newer_success).timestamp() > first_failed + ) + if row["gave_up"]: + given_up.add((target, mid)) + elif first_failed < abandon_before and target_healthy: + given_up.add((target, mid)) + if not dry_run: + state.give_up_upload_retry(user_name, target, mid) + logger.warning( + "%s the %s upload of %.2f kg from %s: it has failed for %d days (%d attempts) " + "while newer weigh-ins uploaded, most recently at %s", + "Would give up on" if dry_run else "Giving up on", + target.capitalize(), row["weight_kg"], row["measurement_timestamp"], + int((now - first_failed) // 86400), row["attempts"], newer_success, + ) + elif row["attempts"] >= MAX_RETRY_ATTEMPTS or first_failed < oldest_allowed: + capped[(target, mid)] = target_healthy + return given_up, capped + + def sync_user(user: UserConfig, state: SyncState, backfill_days: int | None = None, headless: bool = False, dry_run: bool = False, repair_days: int | None = None, target: str | None = None, report: SyncReport | None = None) -> tuple[dict[str, int], dict[str, str]]: """Sync one user's Eufy data to configured targets. @@ -134,6 +223,9 @@ def sync_user(user: UserConfig, state: SyncState, backfill_days: int | None = No # GarminConnectTooManyRequestsError) still work. raise first_exception + # Before the cursor: an entry given up here must not pull it back. + given_up, capped = _triage_retries(user.name, state, [name for name, _ in targets], dry_run) + # Determine how far back to fetch: from the OLDEST per-target cursor, # not a shared one. If one target was down (auth failing) while # another kept syncing, a shared cursor would advance past the outage @@ -161,9 +253,27 @@ def sync_user(user: UserConfig, state: SyncState, backfill_days: int | None = No # Processed data can arrive after newer weigh-ins have # advanced the cursor. Include small timestamp shifts. ts = min(ts, pending_ts - UPGRADE_MAX_SECONDS) if ts is not None else pending_ts - UPGRADE_MAX_SECONDS + retry_ts = state.get_oldest_waiting_retry_timestamp(user.name, name) + if retry_ts is not None: + # A capped failure sits behind the cursor once newer + # ones upload past it; reach back for it, but no + # further than the abandon window and a margin. + floor = int(time.time()) - (RETRY_ABANDON_DAYS + RETRY_REACH_BACK_MARGIN_DAYS) * 86400 + if retry_ts < floor: + logger.info( + "Some queued %s uploads are older than %d days and still waiting; " + "the regular fetch no longer reaches them (--backfill-days does)", + name.capitalize(), RETRY_ABANDON_DAYS + RETRY_REACH_BACK_MARGIN_DAYS, + ) + retry_ts = floor + ts = min(ts if ts is not None else default_cursor, retry_ts) cursors.append(ts if ts is not None else default_cursor) after_timestamp = min(cursors) + # The error of a capped measurement that failed this run, per target. + # Reported only if no newer one reaches the target after it. + capped_errors: dict[str, str] = {} + pending = {} pending_previous_ids = {} if any(name == "garmin" for name, _ in targets): @@ -247,6 +357,16 @@ def sync_user(user: UserConfig, state: SyncState, backfill_days: int | None = No continue if target_name in current_weight_timestamps and m.measurement_id != target_measurement.measurement_id: continue + # Repair is an explicit request to resend, so it overrides. A + # saved replacement for a deleted weight-only entry always + # goes through: skipping it would leave that day empty. + if ( + not repair + and (target_name, target_measurement.measurement_id) in given_up + and target_measurement.measurement_id not in pending + ): + logger.debug("Skipping %s upload that was given up: %s", target_name.capitalize(), target_measurement.measurement_id) + continue # Still consulted in repair mode: it decides whether the sync # is recorded below, since re-uploading a known id must not @@ -296,6 +416,13 @@ def sync_user(user: UserConfig, state: SyncState, backfill_days: int | None = No continue if dry_run: + if (target_name, target_measurement.measurement_id) in capped: + print( + f"[DRY RUN] Would retry {target_name}: {target_measurement.weight_kg:.1f} kg at {target_measurement.timestamp}" + " (past its retry cap; newer ones still upload if it fails)" + ) + counts[target_name] += 1 + continue print(f"[DRY RUN] Would sync to {target_name}: {target_measurement.weight_kg:.1f} kg at {target_measurement.timestamp}") counts[target_name] += 1 continue @@ -334,6 +461,11 @@ def sync_user(user: UserConfig, state: SyncState, backfill_days: int | None = No # used to abort sync_user and take Strava down with it, even # though Strava was fine. Same containment the auth loop above # already has: record the message, drop the target, keep going. + # A capped measurement has failed for days already; one try a + # run is enough, and backoff on each would stall an outage run. + capped_key = (target_name, target_measurement.measurement_id) + is_capped = capped_key in capped + attempt = (lambda fn, _description: fn()) if is_capped else _retry try: if target_name == "garmin": if upgrade_row is not None: @@ -348,14 +480,14 @@ def sync_user(user: UserConfig, state: SyncState, backfill_days: int | None = No datetime.fromisoformat(upgrade_row["measurement_timestamp"]), upgrade_row["weight_kg"], ) - # _retry calls the lambda before this iteration ends, + # attempt calls the lambda before this iteration ends, # so the loop variables it closes over are the right ones. - result = _retry( + result = attempt( lambda: client.upload_body_composition(body_comp), # noqa: B023 f"Garmin upload ({m.measurement_id})", ) else: - result = _retry( + result = attempt( lambda: client.update_weight(target_measurement.weight_kg), # noqa: B023 f"{target_name.capitalize()} weight update ({target_measurement.measurement_id})", ) @@ -377,16 +509,56 @@ def sync_user(user: UserConfig, state: SyncState, backfill_days: int | None = No logger.info("Upgraded weight-only entry to full body comp for %s", m.timestamp.astimezone().date()) if target_name == "garmin" and (upgrade_row is not None or m.measurement_id in pending): state.clear_pending_upgrade(user.name, m.measurement_id) + if target_name in current_weight_timestamps: + state.clear_upload_retries_through(user.name, target_name, target_measurement.timestamp) + else: + state.clear_upload_retry(user.name, target_name, m.measurement_id) + # Evidence for the give-up rule: the target took a + # weigh-in newer than every older queued failure. + state.note_newer_upload( + user.name, target_name, m.timestamp, datetime.now(timezone.utc).isoformat(), + ) + capped_errors.pop(target_name, None) except UnsupportedMeasurementError as e: logger.warning("Skipping %s for %s: %s", target_name.capitalize(), user.name, e) continue except Exception as e: logger.error("Upload to %s failed for %s: %s", target_name, user.name, e) - # str(e) carries the actionable text the CLI keys its - # notification off (e.g. the "--reauth" hint), so it must - # reach the caller unwrapped. - errors[target_name] = str(e) - targets = [t for t in targets if t[0] != target_name] + retryable = not _is_permanent(e) and upgrade_row is None and m.measurement_id not in pending + if retryable and is_capped: + # Past its cap: try the newer measurements instead of + # stopping here. Reported only if none of them land + # and the target has not taken a newer one since + # this first failed; status still counts it. + if not capped[capped_key]: + capped_errors[target_name] = str(e) + else: + # str(e) carries the actionable text the CLI keys its + # notification off (e.g. the "--reauth" hint), so it + # must reach the caller unwrapped. + errors[target_name] = str(e) + targets = [t for t in targets if t[0] != target_name] + # Permanent and auth failures need the user, not a retry. + # A replacement for a weight-only entry already has its + # own store (pending_upgrades) and must never be given up. + if retryable: + attempts = state.record_upload_failure( + user_name=user.name, + target=target_name, + measurement_id=target_measurement.measurement_id, + measurement_timestamp=target_measurement.timestamp.isoformat(), + weight_kg=target_measurement.weight_kg, + failed_at=datetime.now(timezone.utc).isoformat(), + ) + if target_name in current_weight_timestamps: + # Only the newest weight is worth sending later. + state.clear_upload_retries_through( + user.name, target_name, target_measurement.timestamp - timedelta(microseconds=1), + ) + logger.info( + "%s will retry %s on the next run (failed %d time%s)", + target_name.capitalize(), target_measurement.measurement_id, attempts, "" if attempts == 1 else "s", + ) continue counts[target_name] += 1 @@ -402,6 +574,10 @@ def sync_user(user: UserConfig, state: SyncState, backfill_days: int | None = No # nowhere to go. break + # No newer measurement landed after a capped failure: either nothing + # newer exists or the target is down. Keep the entry; report it. + for target_name, message in capped_errors.items(): + errors.setdefault(target_name, message) return counts, errors finally: diff --git a/pyproject.toml b/pyproject.toml index 9b017b9..f62842e 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "eufy-sync" -version = "1.14.0" +version = "1.15.0" description = "Sync Eufy smart scale data to Garmin Connect, Strava, and Zwift" readme = "README.md" license = "MIT" diff --git a/tests/conftest.py b/tests/conftest.py index ba91c0e..18f5d55 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -54,6 +54,14 @@ def delete_password(service, account): monkeypatch.setattr("keyring.get_password", get_password) monkeypatch.setattr("keyring.delete_password", delete_password) + # A test that reaches an unpatched relogin would otherwise run a real + # Garmin SSO login with its fake credentials. Tests that need a login + # patch GarminAuth or the Garmin class and never get here. + def refuse_login(self, *args, **kwargs): + raise RuntimeError("tests must not log in to Garmin; patch the relogin") + + monkeypatch.setattr("garminconnect.Garmin.login", refuse_login) + @pytest.fixture(autouse=True) def _mute_notifications(monkeypatch): diff --git a/tests/test_garmin_client.py b/tests/test_garmin_client.py index 8804b22..4cbd373 100644 --- a/tests/test_garmin_client.py +++ b/tests/test_garmin_client.py @@ -1,10 +1,14 @@ from __future__ import annotations +import json from datetime import datetime, timedelta, timezone from unittest.mock import MagicMock, patch import pytest +import requests +from garminconnect import Garmin +from eufy_sync import garmin_client from eufy_sync.config import GarminConfig from eufy_sync.garmin_client import GarminClient from eufy_sync.transform import GarminBodyComposition @@ -301,9 +305,13 @@ def test_duplicate_check_fails_open_when_the_relogin_fails(): def test_a_successful_relogin_is_not_repeated_when_the_api_keeps_refusing(): # Cloudflare 403s on the API while SSO logins succeed: the first call # relogs in, and every later call must fail with its own error instead - # of running a full login each time. + # of running a full login each time. The upload, having outlived both the + # relogin and the fingerprint retry, is classified as permanent so _retry + # does not ask again seconds later. from garminconnect import GarminConnectConnectionError + from eufy_sync.sync import PermanentSyncError + blocked = GarminConnectConnectionError("API Error 403 - ") stale = MagicMock() stale.get_body_composition.side_effect = blocked @@ -319,7 +327,7 @@ def test_a_successful_relogin_is_not_repeated_when_the_api_keeps_refusing(): assert client.has_weight_on_date(datetime(2026, 6, 10, tzinfo=timezone.utc)) is False with pytest.raises(GarminConnectConnectionError, match="403"): client.check_connection() - with pytest.raises(GarminConnectConnectionError, match="403"): + with pytest.raises(PermanentSyncError, match="403"): client.upload_body_composition(bc) reauth.assert_called_once() @@ -588,3 +596,721 @@ def test_old_session_still_serves_calls_after_a_failed_relogin(): assert client.has_weight_on_date(datetime(2026, 6, 10, tzinfo=timezone.utc)) is False assert client.upload_body_composition(bc) == {"ok": True} assert client._garmin is session + + +# --------------------------------------------------------------------------- +# Cloudflare blocks and the browser-fingerprint fallback (upstream issue #444) +# +# These run the real library against fake transports: a requests adapter for +# the library's own API session and a fake curl_cffi session for the fallback. +# Nothing leaves the machine. +# --------------------------------------------------------------------------- + +TOKEN = "token-abc" +CF_PAGE = ( + 403, + {"Content-Type": "text/html; charset=UTF-8", "Server": "cloudflare", "CF-RAY": "8f1-EWR"}, + b"Just a moment..." + b"", +) +CF_MITIGATED = (403, {"Content-Type": "text/html", "cf-mitigated": "challenge"}, b"") +# Issue #444's block: the same JSON a refused token gets, behind the same +# Cloudflare headers every API answer carries. +JSON_403 = ( + 403, + {"Content-Type": "application/json", "Server": "cloudflare", "CF-RAY": "8f1-EWR"}, + b'{"message":"HTTP 403 Forbidden","error":"ForbiddenException"}', +) + + +# A gateway error page in Cloudflare's livery. It says nothing about whether +# Garmin received the request. +CF_502 = ( + 502, + {"Content-Type": "text/html; charset=UTF-8", "Server": "cloudflare", "CF-RAY": "8f1-EWR"}, + b"502 Bad gateway
" + b"Cloudflare Ray ID: 8f1
", +) + + +def _ok(body: dict, status: int = 200): + return (status, {"Content-Type": "application/json"}, json.dumps(body).encode()) + + +class _FakeAdapter(requests.adapters.BaseAdapter): + """Answers the library's requests session from a script of responses.""" + + def __init__(self, responses): + super().__init__() + self.responses = list(responses) + self.sent = [] + + def send(self, request, **kwargs): + self.sent.append(request) + status, headers, body = self.responses.pop(0) + resp = requests.Response() + resp.status_code = status + resp.headers.update(headers) + resp._content = body + resp.url = request.url + resp.request = request + return resp + + def close(self): + pass + + +class _FakeCffiResponse: + def __init__(self, status, headers, body): + self.status_code = status + self.headers = headers + self.content = body + self.text = body.decode() + + def json(self): + return json.loads(self.content) + + +class _FakeCffi: + """Stands in for curl_cffi sessions; records what each fallback sent.""" + + def __init__(self, responses): + self.responses = list(responses) + self.calls = [] + + def session(self): + fake = self + + class _Session: + def request(self, method, url, headers=None, **kwargs): + fake.calls.append({"method": method, "url": url, "headers": headers, **kwargs}) + return _FakeCffiResponse(*fake.responses.pop(0)) + + def close(self): + pass + + return _Session() + + +def _real_garmin(responses): + garmin = Garmin("g@example.com", "pw", retry_attempts=0) + garmin.client.di_token = TOKEN + adapter = _FakeAdapter(responses) + garmin.client._api_session.mount("https://", adapter) + return garmin, adapter + + +def _client_on(garmin, interactive=False): + client = _client_with_fake_garmin(garmin) + client._allow_interactive = interactive + client._watch_responses(garmin) + return client + + +BC = GarminBodyComposition(timestamp="2026-06-10T08:00:00+00:00", weight=86.2, percent_fat=18.5) + + +def test_cloudflare_page_on_upload_retries_through_curl_cffi_without_relogin(): + garmin, adapter = _real_garmin([CF_PAGE]) + original_session = garmin.client._api_session + client = _client_on(garmin) + cffi = _FakeCffi([_ok({"detailedImportResult": {"successes": [], "failures": []}}, status=202)]) + + with patch.object(garmin_client, "_new_impersonating_session", cffi.session), \ + patch.object(client._auth, "silent_reauth") as reauth: + result = client.upload_body_composition(BC) + + reauth.assert_not_called() # the token was never the problem + assert result == {"detailedImportResult": {"successes": [], "failures": []}} + assert len(adapter.sent) == 1 and len(cffi.calls) == 1 + call = cffi.calls[0] + # Same endpoint, same token, and the FIT goes as a curl_cffi multipart + # body because curl_cffi refuses files=. + assert call["method"] == "POST" + assert call["url"] == "https://connectapi.garmin.com/upload-service/upload" + assert call["headers"]["Authorization"] == f"Bearer {TOKEN}" + assert "files" not in call and call["multipart"] is not None + assert garmin.client._api_session is original_session # swapped back + + +def test_json_403_on_a_read_tries_the_fingerprint_before_any_relogin(): + # Issue #444: the block is a JSON 403. The fingerprint retry is one + # request; a relogin is a login, a 429 risk, and on that network its own + # token check fails the same way. + garmin, _ = _real_garmin([JSON_403]) + client = _client_on(garmin) + cffi = _FakeCffi([_ok({"dateWeightList": [{"weight": 86000}]})]) + + with patch.object(garmin_client, "_new_impersonating_session", cffi.session), \ + patch.object(client._auth, "silent_reauth") as reauth: + assert client.has_weight_on_date(datetime(2026, 6, 10, tzinfo=timezone.utc)) is True + + reauth.assert_not_called() + call = cffi.calls[0] + assert call["method"] == "GET" + assert call["url"] == "https://connectapi.garmin.com/weight-service/weight/dateRange" + assert call["headers"]["Authorization"] == f"Bearer {TOKEN}" + assert set(call["params"]) == {"startDate", "endDate"} + + +def test_once_the_fingerprint_works_later_calls_skip_plain_requests(): + # The block belongs to the network, so after one fingerprint success the + # rest of the run goes straight through curl_cffi. A new client (the next + # run) starts on plain requests again. + held = {"samplePk": 1, "weight": 86200.0, "timestampGMT": _millis(BC_INSTANT)} + garmin, adapter = _real_garmin([JSON_403]) + client = _client_on(garmin) + cffi = _FakeCffi([ + _ok({"dateWeightList": []}), + _ok({"dateWeightList": [held]}), + _ok({}), + _ok({"detailedImportResult": {}}, status=202), + ]) + original_session = garmin.client._api_session + + with patch.object(garmin_client, "_new_impersonating_session", cffi.session), \ + patch.object(client._auth, "silent_reauth") as reauth: + assert client.has_weight_on_date(datetime(2026, 6, 10, tzinfo=timezone.utc)) is False + assert client.delete_weight_entry(BC_INSTANT, 86.2) is True + client.upload_body_composition(BC) + + reauth.assert_not_called() + assert len(adapter.sent) == 1 # only the first, refused, plain request + assert [c["method"] for c in cffi.calls] == ["GET", "GET", "DELETE", "POST"] + assert garmin.client._api_session is original_session + + next_run, next_adapter = _real_garmin([_ok({"dateWeightList": []})]) + fresh_client = _client_on(next_run) + with patch.object(garmin_client, "_new_impersonating_session", cffi.session): + fresh_client.has_weight_on_date(datetime(2026, 6, 10, tzinfo=timezone.utc)) + assert len(next_adapter.sent) == 1 and len(cffi.calls) == 4 + + +def test_a_failed_fingerprint_retry_does_not_make_it_sticky(): + garmin, adapter = _real_garmin([JSON_403, _ok({"dateWeightList": []})]) + client = _client_on(garmin) + client._reauth_attempted = True # isolate the fallback from the relogin + cffi = _FakeCffi([JSON_403]) + with patch.object(garmin_client, "_new_impersonating_session", cffi.session): + assert client.has_weight_on_date(datetime(2026, 6, 10, tzinfo=timezone.utc)) is False + assert client.has_weight_on_date(datetime(2026, 6, 10, tzinfo=timezone.utc)) is False + assert len(adapter.sent) == 2 and len(cffi.calls) == 1 + + +def test_a_cloudflare_block_that_survives_the_fingerprint_never_relogs_in(): + from eufy_sync.sync import PermanentSyncError + + garmin, _ = _real_garmin([CF_PAGE]) + client = _client_on(garmin) + cffi = _FakeCffi([CF_MITIGATED]) + + with patch.object(garmin_client, "_new_impersonating_session", cffi.session), \ + patch.object(client._auth, "silent_reauth") as reauth: + with pytest.raises(PermanentSyncError, match="Cloudflare"): + client.upload_body_composition(BC) + + reauth.assert_not_called() + assert len(cffi.calls) == 1 # one fallback per call, not a loop + + +def test_a_json_403_that_survives_the_fingerprint_relogs_in_once(): + # Both transports refuse the token, so it may really be dead: the run's one + # relogin happens, and the fresh session gets its own fingerprint retry. + stale, _ = _real_garmin([JSON_403]) + fresh, fresh_adapter = _real_garmin([JSON_403]) + client = _client_on(stale) + cffi = _FakeCffi([JSON_403, _ok({"dateWeightList": []})]) + + with patch.object(garmin_client, "_new_impersonating_session", cffi.session), \ + patch.object(client._auth, "silent_reauth", return_value=fresh) as reauth: + assert client.has_weight_on_date(datetime(2026, 6, 10, tzinfo=timezone.utc)) is False + + reauth.assert_called_once() + assert len(fresh_adapter.sent) == 1 and len(cffi.calls) == 2 + assert client._garmin is fresh + + +def test_a_401_relogs_in_without_the_fingerprint_retry(): + from garminconnect import GarminConnectConnectionError + + dead = MagicMock() + dead.get_body_composition.side_effect = GarminConnectConnectionError("API Error 401 - ") + fresh = MagicMock() + fresh.get_body_composition.return_value = {"dateWeightList": []} + client = _client_with_fake_garmin(dead) + client._allow_interactive = False + with patch.object(client, "_call_impersonating") as fallback, \ + patch.object(client._auth, "silent_reauth", return_value=fresh): + client.has_weight_on_date(datetime(2026, 6, 10, tzinfo=timezone.utc)) + fallback.assert_not_called() + dead.get_body_composition.assert_called_once() + + +def test_no_fallback_when_the_library_has_no_api_session(): + # A future library without _api_session loses the fallback, not the sync: + # a 403 goes straight to the relogin as before. + from garminconnect import GarminConnectConnectionError + + stale = MagicMock() + stale.client._api_session = None + stale.get_body_composition.side_effect = GarminConnectConnectionError("API Error 403 - ") + fresh = MagicMock() + fresh.get_body_composition.return_value = {"dateWeightList": [{"weight": 1}]} + client = _client_with_fake_garmin(stale) + client._allow_interactive = False + with patch.object(client._auth, "silent_reauth", return_value=fresh) as reauth: + assert client.has_weight_on_date(datetime(2026, 6, 10, tzinfo=timezone.utc)) is True + reauth.assert_called_once() + stale.get_body_composition.assert_called_once() + + +# --------------------------------------------------------------------------- +# Upload outcome classification +# --------------------------------------------------------------------------- + + +CONFLICT_409 = (409, {"Content-Type": "application/json"}, json.dumps({ + "detailedImportResult": {"failures": [{"messages": [{"content": "Duplicate Activity."}]}]}, +}).encode()) +BC_INSTANT = datetime.fromisoformat(BC.timestamp) + + +def test_upload_409_counts_as_uploaded_once_the_lookup_finds_the_weigh_in(): + held = {"samplePk": 1, "weight": 86200.0, "timestampGMT": _millis(BC_INSTANT)} + garmin, adapter = _real_garmin([CONFLICT_409, _ok({"dateWeightList": [held]})]) + client = _client_on(garmin) + with patch.object(client._auth, "silent_reauth") as reauth: + assert client.upload_body_composition(BC) == {"status": "duplicate"} + reauth.assert_not_called() + assert len(adapter.sent) == 2 + assert "/weight-service/weight/dayview/" in adapter.sent[1].url + + +@pytest.mark.parametrize("entries", [ + [], + # Right weight, but a manual weigh-in hours away is not our upload. + [{"samplePk": 1, "weight": 86200.0, "timestampGMT": _millis(BC_INSTANT + timedelta(hours=3))}], + # Right time, wrong weight. + [{"samplePk": 1, "weight": 90000.0, "timestampGMT": _millis(BC_INSTANT)}], +]) +def test_upload_409_without_the_weigh_in_on_garmin_is_permanent(entries): + from eufy_sync.sync import PermanentSyncError, _is_permanent + + garmin, _ = _real_garmin([CONFLICT_409, _ok({"dateWeightList": entries})]) + client = _client_on(garmin) + with pytest.raises(PermanentSyncError, match="409") as exc: + client.upload_body_composition(BC) + assert _is_permanent(exc.value) + + +@pytest.mark.parametrize("entry", [ + {"samplePk": 1, "weight": 86200.0}, + {"samplePk": 1, "weight": 86200.0, "timestampGMT": None, "date": "2026-06-10"}, + {"samplePk": 1, "weight": 86200.0, "timestampGMT": "08:00"}, +]) +def test_upload_409_is_not_confirmed_by_an_untimed_entry(entry): + """Same weight, but no parseable timestamp: it could be any weigh-in that + day, so it does not prove ours is there. delete_weight_entry still trusts + a lone untimed match; the 409 confirmation does not.""" + from eufy_sync.sync import PermanentSyncError + + garmin, _ = _real_garmin([CONFLICT_409, _ok({"dateWeightList": [entry]})]) + client = _client_on(garmin) + with pytest.raises(PermanentSyncError, match="409"): + client.upload_body_composition(BC) + + +@pytest.mark.parametrize("status", [500, 502, 503]) +def test_upload_409_lookup_that_keeps_failing_waits_for_the_next_run(status): + """Only the lookup is repeated, a bounded number of times, and the error + that escapes is one sync's _retry does not retry in-run.""" + from eufy_sync.sync import RetryNextRunError, _is_permanent + + failing = (status, {"Content-Type": "application/json"}, b"{}") + garmin, adapter = _real_garmin([CONFLICT_409] + [failing] * garmin_client._CONFIRM_LOOKUP_ATTEMPTS) + client = _client_on(garmin) + cffi = _FakeCffi([]) + with patch.object(garmin_client, "_new_impersonating_session", cffi.session), \ + patch.object(client._auth, "silent_reauth") as reauth, \ + patch.object(garmin_client.time, "sleep") as sleep: + with pytest.raises(RetryNextRunError) as exc: + client.upload_body_composition(BC) + assert not _is_permanent(exc.value) + reauth.assert_not_called() + # One POST, then only the GET is repeated; nothing replayed through curl_cffi. + assert [r.method for r in adapter.sent] == ["POST"] + ["GET"] * garmin_client._CONFIRM_LOOKUP_ATTEMPTS + assert cffi.calls == [] + assert sleep.call_count == garmin_client._CONFIRM_LOOKUP_ATTEMPTS - 1 + + +def test_upload_409_lookup_that_recovers_confirms_without_resending_the_upload(): + held = {"samplePk": 1, "weight": 86200.0, "timestampGMT": _millis(BC_INSTANT)} + failing = (503, {"Content-Type": "application/json"}, b"{}") + garmin, adapter = _real_garmin([CONFLICT_409, failing, _ok({"dateWeightList": [held]})]) + client = _client_on(garmin) + with patch.object(garmin_client.time, "sleep"): + assert client.upload_body_composition(BC) == {"status": "duplicate"} + assert [r.method for r in adapter.sent] == ["POST", "GET", "GET"] + + +def test_upload_409_lookup_network_failure_waits_for_the_next_run(): + from garminconnect import GarminConnectConnectionError + + from eufy_sync.sync import RetryNextRunError, _is_permanent + + fake = MagicMock() + fake.add_body_composition.side_effect = GarminConnectConnectionError("API Error 409 - Duplicate") + fake.get_daily_weigh_ins.side_effect = GarminConnectConnectionError("Connection error: timed out") + client = _client_with_fake_garmin(fake) + with patch.object(garmin_client.time, "sleep"): + with pytest.raises(RetryNextRunError, match="timed out") as exc: + client.upload_body_composition(BC) + assert not _is_permanent(exc.value) + fake.add_body_composition.assert_called_once() + assert fake.get_daily_weigh_ins.call_count == garmin_client._CONFIRM_LOOKUP_ATTEMPTS + + +def test_unconfirmed_409_posts_once_per_sync_run_and_queues_the_measurement(tmp_path): + """Through sync_user: a lookup that keeps failing produces exactly one + POST in the run, and the measurement waits in the retry queue.""" + from garminconnect import GarminConnectConnectionError + + from eufy_sync.config import EufyConfig, GarminConfig, UserConfig + from eufy_sync.eufy_client import EufyMeasurement + from eufy_sync.state import SyncState + from eufy_sync.sync import sync_user + + taken = datetime.now(timezone.utc).replace(microsecond=0) - timedelta(hours=2) + m = EufyMeasurement( + measurement_id="m1", customer_id="cust", device_id="dev", timestamp=taken, weight_kg=86.2, + ) + fake = MagicMock() + fake.add_body_composition.side_effect = GarminConnectConnectionError("API Error 409 - Duplicate") + fake.get_daily_weigh_ins.side_effect = GarminConnectConnectionError("Connection error: timed out") + fake.get_body_composition.return_value = {"dateWeightList": []} + client = _client_with_fake_garmin(fake) + client.authenticate = MagicMock() + eufy = MagicMock() + eufy.fetch_measurements.return_value = [m] + user = UserConfig( + name="default", + eufy=EufyConfig(email="e@example.com", password="pw"), + garmin=GarminConfig(email="g@example.com", password="pw"), + ) + state = SyncState(tmp_path / "s.db") + + with patch("eufy_sync.sync.EufyClient", return_value=eufy), \ + patch("eufy_sync.garmin_client.GarminClient", return_value=client), \ + patch("eufy_sync.sync.time.sleep"), \ + patch.object(garmin_client.time, "sleep"): + counts, errors = sync_user(user, state, headless=True) + + fake.add_body_composition.assert_called_once() + assert counts["garmin"] == 0 and "409" in errors["garmin"] + assert state.waiting_upload_retries("default") == {"garmin": 1} + state.close() + + +def test_upload_409_lookup_rate_limit_ends_garmin_for_the_run(): + from garminconnect import GarminConnectTooManyRequestsError + + garmin, adapter = _real_garmin([CONFLICT_409, (429, {"Content-Type": "application/json"}, b"{}")]) + client = _client_on(garmin) + with pytest.raises(GarminConnectTooManyRequestsError): + client.upload_body_composition(BC) + assert [r.method for r in adapter.sent] == ["POST", "GET"] + + +def test_upload_409_lookup_gets_the_fingerprint_fallback_without_resending_the_upload(): + held = {"samplePk": 1, "weight": 86200.0, "timestampGMT": _millis(BC_INSTANT)} + garmin, adapter = _real_garmin([CONFLICT_409, JSON_403]) + client = _client_on(garmin) + cffi = _FakeCffi([_ok({"dateWeightList": [held]})]) + with patch.object(garmin_client, "_new_impersonating_session", cffi.session), \ + patch.object(client._auth, "silent_reauth") as reauth: + assert client.upload_body_composition(BC) == {"status": "duplicate"} + reauth.assert_not_called() + assert [r.method for r in adapter.sent] == ["POST", "GET"] + assert [c["method"] for c in cffi.calls] == ["GET"] + assert client._impersonate_always is True + + +def test_upload_409_lookup_rides_the_sticky_fingerprint_path(): + held = {"samplePk": 1, "weight": 86200.0, "timestampGMT": _millis(BC_INSTANT)} + garmin, adapter = _real_garmin([]) + client = _client_on(garmin) + client._impersonate_always = True + cffi = _FakeCffi([CONFLICT_409, _ok({"dateWeightList": [held]})]) + with patch.object(garmin_client, "_new_impersonating_session", cffi.session): + assert client.upload_body_composition(BC) == {"status": "duplicate"} + assert adapter.sent == [] + assert [c["method"] for c in cffi.calls] == ["POST", "GET"] + + +def test_upload_409_lookup_relogs_in_once_without_resending_the_upload(): + from garminconnect import GarminConnectConnectionError + + held = {"samplePk": 1, "weight": 86200.0, "timestampGMT": _millis(BC_INSTANT)} + dead = MagicMock() + dead.add_body_composition.side_effect = GarminConnectConnectionError("API Error 409 - Duplicate") + dead.get_daily_weigh_ins.side_effect = GarminConnectConnectionError("API Error 401 - ") + fresh = MagicMock() + fresh.get_daily_weigh_ins.return_value = {"dateWeightList": [held]} + client = _client_with_fake_garmin(dead) + client._allow_interactive = False + with patch.object(client._auth, "silent_reauth", return_value=fresh) as reauth: + assert client.upload_body_composition(BC) == {"status": "duplicate"} + reauth.assert_called_once() + dead.add_body_composition.assert_called_once() + fresh.add_body_composition.assert_not_called() + + +def test_upload_409_lookup_refused_after_the_relogin_is_permanent(): + from garminconnect import GarminConnectConnectionError + + from eufy_sync.sync import PermanentSyncError + + dead = MagicMock() + dead.add_body_composition.side_effect = GarminConnectConnectionError("API Error 409 - Duplicate") + dead.get_daily_weigh_ins.side_effect = GarminConnectConnectionError("API Error 401 - ") + fresh = MagicMock() + fresh.get_daily_weigh_ins.side_effect = GarminConnectConnectionError("API Error 401 - ") + client = _client_with_fake_garmin(dead) + client._allow_interactive = False + with patch.object(client._auth, "silent_reauth", return_value=fresh): + with pytest.raises(PermanentSyncError, match="--reauth garmin"): + client.upload_body_composition(BC) + fresh.add_body_composition.assert_not_called() + + +def test_upload_429_becomes_a_rate_limit_that_sync_does_not_retry(): + from garminconnect import GarminConnectTooManyRequestsError + + from eufy_sync.sync import _is_permanent + + garmin, adapter = _real_garmin([(429, {"Content-Type": "application/json", "Retry-After": "120"}, b"{}")]) + client = _client_on(garmin) + with patch.object(client._auth, "silent_reauth") as reauth: + with pytest.raises(GarminConnectTooManyRequestsError) as exc: + client.upload_body_composition(BC) + reauth.assert_not_called() + assert len(adapter.sent) == 1 + assert _is_permanent(exc.value) + + +@pytest.mark.parametrize("status", [400, 404, 413, 422]) +def test_upload_other_4xx_is_a_permanent_bad_request(status): + from eufy_sync.sync import PermanentSyncError, _is_permanent + + garmin, adapter = _real_garmin([(status, {"Content-Type": "application/json"}, b'{"message":"bad file"}')]) + client = _client_on(garmin) + with patch.object(client._auth, "silent_reauth") as reauth: + with pytest.raises(PermanentSyncError, match=str(status)) as exc: + client.upload_body_composition(BC) + reauth.assert_not_called() + assert len(adapter.sent) == 1 + assert _is_permanent(exc.value) + + +@pytest.mark.parametrize("status", [408, 500, 502, 503]) +def test_upload_408_and_5xx_stay_transient_for_retry(status): + from garminconnect import GarminConnectConnectionError + + from eufy_sync.sync import _is_permanent + + garmin, _ = _real_garmin([(status, {"Content-Type": "application/json"}, b"{}")]) + client = _client_on(garmin) + with patch.object(client._auth, "silent_reauth") as reauth: + with pytest.raises(GarminConnectConnectionError) as exc: + client.upload_body_composition(BC) + reauth.assert_not_called() + assert not _is_permanent(exc.value) + + +@pytest.mark.parametrize("response", [ + CF_502, + (504, CF_502[1], CF_502[2].replace(b"502", b"504")), + (500, {"Content-Type": "text/html"}, b"Just a moment..."), +]) +def test_upload_after_a_gateway_error_is_never_replayed_through_curl_cffi(response): + """Garmin may have stored the upload behind a gateway error, so only the + ordinary retry policy may send it again, not the fingerprint fallback.""" + from garminconnect import GarminConnectConnectionError + + from eufy_sync.sync import _is_permanent + + garmin, adapter = _real_garmin([response]) + client = _client_on(garmin) + cffi = _FakeCffi([]) + with patch.object(garmin_client, "_new_impersonating_session", cffi.session), \ + patch.object(client._auth, "silent_reauth") as reauth: + with pytest.raises(GarminConnectConnectionError) as exc: + client.upload_body_composition(BC) + assert len(adapter.sent) == 1 and cffi.calls == [] + reauth.assert_not_called() + assert not _is_permanent(exc.value) + assert client._impersonate_always is False + + +def test_upload_is_not_replayed_when_the_recorded_answer_contradicts_a_403_message(): + # The recorded response is what Garmin answered; a stray "403" in the + # error text must not trigger a second POST after a 502. + from garminconnect import GarminConnectConnectionError + + from eufy_sync.sync import PermanentSyncError + + fake = MagicMock() + client = _client_with_fake_garmin(fake) + client._reauth_attempted = True # isolate the fallback from the relogin + + def gateway_failure(**kwargs): + client._last_response.record(_FakeCffiResponse(*CF_502)) + raise GarminConnectConnectionError("API Error 403 - proxied") + + fake.add_body_composition.side_effect = gateway_failure + with patch.object(client, "_call_impersonating") as fallback, \ + patch.object(client, "_can_impersonate", return_value=True), \ + patch.object(client._auth, "silent_reauth") as reauth: + with pytest.raises(PermanentSyncError): + client.upload_body_composition(BC) + fallback.assert_not_called() + reauth.assert_not_called() + fake.add_body_composition.assert_called_once() + + +def test_upload_network_failure_stays_transient(): + from garminconnect import GarminConnectConnectionError + + fake = MagicMock() + fake.add_body_composition.side_effect = GarminConnectConnectionError("Connection error: timed out") + client = _client_with_fake_garmin(fake) + with pytest.raises(GarminConnectConnectionError, match="timed out"): + client.upload_body_composition(BC) + + +def test_upload_401_after_a_successful_relogin_is_permanent(): + from garminconnect import GarminConnectConnectionError + + from eufy_sync.sync import PermanentSyncError + + dead = MagicMock() + dead.add_body_composition.side_effect = GarminConnectConnectionError("API Error 401 - ") + fresh = MagicMock() + fresh.add_body_composition.side_effect = GarminConnectConnectionError("API Error 401 - ") + client = _client_with_fake_garmin(dead) + client._allow_interactive = False + with patch.object(client._auth, "silent_reauth", return_value=fresh) as reauth: + with pytest.raises(PermanentSyncError, match="--reauth garmin"): + client.upload_body_composition(BC) + reauth.assert_called_once() + fresh.add_body_composition.assert_called_once() + + +# --------------------------------------------------------------------------- +# The pieces underneath +# --------------------------------------------------------------------------- + + +@pytest.mark.parametrize("message, expected", [ + ("API Error 403 - ", 403), + ("API Error 409 - Duplicate Activity.", 409), + ("API call client error (403): API Error 403", 403), + ("Connection error: timed out", None), + ("weight 4031 kg", None), +]) +def test_status_code_reads_the_library_messages(message, expected): + from garminconnect import GarminConnectConnectionError + assert garmin_client._status_code(GarminConnectConnectionError(message)) == expected + + +@pytest.mark.parametrize("response, blocked", [ + (CF_PAGE, True), + (CF_MITIGATED, True), + (JSON_403, False), # Cloudflare headers alone prove nothing + ((403, {"Content-Type": "text/html"}, b"Forbidden"), False), + ((200, {"Content-Type": "text/html"}, b"Just a moment"), False), + (CF_502, False), # branded gateway page: Garmin may have the request + ((504, {"Content-Type": "text/html", "cf-mitigated": "challenge"}, b""), False), + ((503, {"Content-Type": "text/html"}, b"Just a moment..."), True), + ((429, {"Content-Type": "text/html"}, b"
"), True), + ((429, {"Content-Type": "text/html"}, b"
Cloudflare Ray ID
"), False), + ((403, {"Content-Type": "text/html"}, b"Attention Required! | Cloudflare"), True), + ((403, {"Content-Type": "text/html"}, b"
Cloudflare Ray ID: 8f1
"), False), +]) +def test_cloudflare_block_detection(response, blocked): + recorder = garmin_client._LastResponse() + recorder.record(_FakeCffiResponse(*response)) + assert recorder.is_cloudflare_block() is blocked + + +def test_watch_responses_installs_its_hook_once(): + garmin, _ = _real_garmin([]) + client = _client_on(garmin) + client._watch_responses(garmin) + hooks = garmin.client._api_session.hooks["response"] + assert hooks.count(client._last_response.hook) == 1 + + +def test_curl_cffi_multipart_upload_reaches_a_local_server_intact(): + """The maintainer's reason not to use curl_cffi for data calls is that it + cannot do files=. It cannot, but CurlMime can: send a FIT-sized binary + through the real fallback transport to a local server and parse it back.""" + import threading + from email.parser import BytesParser + from http.server import BaseHTTPRequestHandler, HTTPServer + + received = {} + + class Handler(BaseHTTPRequestHandler): + def do_POST(self): + received["content_type"] = self.headers["Content-Type"] + received["auth"] = self.headers["Authorization"] + received["body"] = self.rfile.read(int(self.headers["Content-Length"])) + self.send_response(200) + self.send_header("Content-Type", "application/json") + self.end_headers() + self.wfile.write(b'{"ok": true}') + + def log_message(self, *args): + pass + + server = HTTPServer(("127.0.0.1", 0), Handler) + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + try: + payload = bytes(range(256)) * 8 + session = garmin_client._ImpersonatingSession(garmin_client._LastResponse()) + resp = session.request( + "POST", f"http://127.0.0.1:{server.server_port}/upload-service/upload", + headers={"Authorization": f"Bearer {TOKEN}"}, + files={"file": ("body_composition.fit", payload)}, + timeout=10, + ) + finally: + server.shutdown() + + assert resp.status_code == 200 and resp.json() == {"ok": True} + assert received["auth"] == f"Bearer {TOKEN}" + message = BytesParser().parsebytes( + b"Content-Type: " + received["content_type"].encode() + b"\r\n\r\n" + received["body"] + ) + (part,) = message.get_payload() + assert part.get_param("name", header="content-disposition") == "file" + assert part.get_filename() == "body_composition.fit" + assert part.get_payload(decode=True) == payload + + +def test_upload_passes_a_failed_relogins_own_error_through_unclassified(): + # A login that failed with an HTTP status is not Garmin refusing the + # upload, so it must not be dressed up as "refused after a fresh login". + from garminconnect import GarminConnectConnectionError + + dead = MagicMock() + dead.add_body_composition.side_effect = GarminConnectConnectionError("API Error 401 - ") + login_failure = GarminConnectConnectionError("Mobile login failed: HTTP 403") + client = _client_with_fake_garmin(dead) + client._allow_interactive = False + with patch.object(client._auth, "silent_reauth", side_effect=login_failure): + with pytest.raises(GarminConnectConnectionError) as exc: + client.upload_body_composition(BC) + assert exc.value is login_failure diff --git a/tests/test_update_check.py b/tests/test_update_check.py index 2533fd3..6482a44 100644 --- a/tests/test_update_check.py +++ b/tests/test_update_check.py @@ -318,3 +318,17 @@ def test_self_update_keeps_the_browser_extra_when_installed(monkeypatch): _self_update() assert mock_run.call_args.args[0] == ["uv", "tool", "install", "--force", "--refresh-package", "eufy-sync", "eufy-sync[browser]==9.9.9"] + + +def test_an_install_ahead_of_pypi_is_not_offered_a_downgrade(capsys): + from eufy_sync.cli import updater + + assert updater.is_newer("1.15.0", "1.14.0") + assert not updater.is_newer("1.14.0", "1.15.0") + assert not updater.is_newer("1.15.0", "1.15.0") + with patch.object(updater, "_latest_pypi_version", return_value="1.14.0"), \ + patch("eufy_sync.__version__", "1.15.0"), \ + patch.object(updater.subprocess, "run") as run: + updater._self_update() + run.assert_not_called() + assert "Already on the latest version" in capsys.readouterr().out diff --git a/tests/test_upload_retries.py b/tests/test_upload_retries.py new file mode 100644 index 0000000..7c21485 --- /dev/null +++ b/tests/test_upload_retries.py @@ -0,0 +1,599 @@ +"""Failed uploads: the fetch cursor brings them back, the retry queue counts +the attempts, lets newer ones past a measurement that keeps failing, and +gives it up only after a month in which newer ones kept landing.""" +from __future__ import annotations + +import logging +import sqlite3 +from datetime import datetime, timedelta, timezone +from pathlib import Path +from unittest.mock import MagicMock, patch + +import httpx + +from eufy_sync.cli.status import _retry_queue_line +from eufy_sync.config import EufyConfig, GarminConfig, StravaConfig, UserConfig, ZwiftConfig +from eufy_sync.eufy_client import EufyMeasurement +from eufy_sync.state import SyncState +from eufy_sync.sync import ( + MAX_RETRY_ATTEMPTS, + RETRY_ABANDON_DAYS, + RETRY_MAX_AGE_DAYS, + RETRY_REACH_BACK_MARGIN_DAYS, + PermanentSyncError, + UnsupportedMeasurementError, + sync_user, +) + +NOW = datetime.now(timezone.utc).replace(microsecond=0) + + +def _m(weight_kg: float, days_ago: float) -> EufyMeasurement: + dt = NOW - timedelta(days=days_ago) + return EufyMeasurement( + measurement_id=f"cust_{int(dt.timestamp())}", + customer_id="cust", + device_id="dev", + timestamp=dt, + weight_kg=weight_kg, + ) + + +def _user(garmin: bool = True, strava: bool = False) -> UserConfig: + return UserConfig( + name="default", + eufy=EufyConfig(email="e@example.com", password="pw"), + garmin=GarminConfig(email="g@example.com", password="pw") if garmin else None, + strava=StravaConfig(client_id="cid", client_secret="csec") if strava else None, + ) + + +def _run(user, state, history, fail_weights=(), garmin_error=None, strava_error=None, garmin_upload=None, + fetches=None, **kwargs): + """One sync run against a fake Eufy that honors the fetch cursor the way + the real one does (timestamp >= after). Uploads of a weight listed in + fail_weights raise the given error (a Garmin 503 by default). + garmin_upload, when given, replaces the fake Garmin upload outright. + fetches, when given, collects each fetch's after_timestamp.""" + def fetch(after_timestamp=None): + if fetches is not None: + fetches.append(after_timestamp) + return [m for m in history if after_timestamp is None or m.timestamp.timestamp() >= after_timestamp] + + fake_eufy = MagicMock() + fake_eufy.fetch_measurements.side_effect = fetch + + garmin_upload_override = garmin_upload + + def garmin_upload(body_comp): + if round(body_comp.weight, 2) in fail_weights: + raise garmin_error or RuntimeError("Garmin returned 503") + return {"ok": True} + + fake_garmin = MagicMock() + fake_garmin.has_weight_on_date.return_value = False + fake_garmin.upload_body_composition.side_effect = garmin_upload_override or garmin_upload + + def strava_update(weight_kg): + if round(weight_kg, 2) in fail_weights: + raise strava_error or RuntimeError("Strava returned 503") + return {"weight": weight_kg} + + fake_strava = MagicMock() + fake_strava.update_weight.side_effect = strava_update + + with patch("eufy_sync.sync.EufyClient", return_value=fake_eufy), \ + patch("eufy_sync.garmin_client.GarminClient", return_value=fake_garmin), \ + patch("eufy_sync.strava_client.StravaClient", return_value=fake_strava), \ + patch("eufy_sync.sync.time.sleep"): + counts, errors = sync_user(user, state, headless=True, **kwargs) + return counts, errors, fake_garmin, fake_strava + + +def _garmin_weights(fake_garmin) -> list[float]: + return [round(c.args[0].weight, 2) for c in fake_garmin.upload_body_composition.call_args_list] + + +def _rows(state, user_name="default") -> dict[tuple[str, str], dict]: + return {(r["target"], r["measurement_id"]): r for r in state.get_upload_retries(user_name)} + + +def test_garmin_failure_is_queued_and_retried_once_next_run(tmp_path: Path): + state = SyncState(tmp_path / "s.db") + user = _user() + m1, m2, m3 = _m(80.0, 3), _m(81.0, 2), _m(82.0, 1) + + counts, errors, garmin, _ = _run(user, state, [m1, m2, m3], fail_weights={81.0}) + assert counts["garmin"] == 1 and "503" in errors["garmin"] + assert _rows(state)[("garmin", m2.measurement_id)]["attempts"] == 1 + assert state.waiting_upload_retries("default") == {"garmin": 1} + + counts, errors, garmin, _ = _run(user, state, [m1, m2, m3]) + assert errors == {} + # The failed one is replayed once; the one that landed is not resent. + assert _garmin_weights(garmin) == [81.0, 82.0] + assert _rows(state) == {} + state.close() + + +def test_attempts_accumulate_across_runs(tmp_path: Path): + state = SyncState(tmp_path / "s.db") + user = _user() + m1 = _m(80.0, 1) + for _ in range(3): + _run(user, state, [m1], fail_weights={80.0}) + assert _rows(state)[("garmin", m1.measurement_id)]["attempts"] == 3 + state.close() + + +def test_permanent_unsupported_and_auth_failures_are_not_queued(tmp_path: Path): + state = SyncState(tmp_path / "s.db") + user = _user(strava=True) + m1 = _m(80.0, 1) + unauthorized = httpx.HTTPStatusError( + "401", request=httpx.Request("POST", "https://example.invalid"), + response=httpx.Response(401), + ) + + _run(user, state, [m1], fail_weights={80.0}, + garmin_error=PermanentSyncError("Run: eufy-sync --reauth garmin"), + strava_error=UnsupportedMeasurementError("out of range")) + _run(user, state, [m1], fail_weights={80.0}, garmin_error=unauthorized, + strava_error=PermanentSyncError("Strava rejected the session")) + assert _rows(state) == {} + state.close() + + +def _cap(state, measurement, *, attempts=MAX_RETRY_ATTEMPTS, first_failed=NOW): + state.record_upload_failure( + "default", "garmin", measurement.measurement_id, measurement.timestamp.isoformat(), + measurement.weight_kg, first_failed.isoformat(), + ) + state._conn.execute( + "UPDATE upload_retries SET attempts = ? WHERE measurement_id = ?", + (attempts, measurement.measurement_id), + ) + state._conn.commit() + + +def _backdate(state, measurement, days_ago): + """Make an entry look as if it first failed days_ago.""" + state._conn.execute( + "UPDATE upload_retries SET first_failed_at = ? WHERE measurement_id = ?", + ((NOW - timedelta(days=days_ago)).isoformat(), measurement.measurement_id), + ) + state._conn.commit() + + +def test_capped_entry_is_tried_each_run_without_blocking_newer_ones(tmp_path: Path): + """A measurement past its cap is tried once a run, in order; its failure + lets the newer ones through, and it stays queued with its attempts + counting. One short of the cap still stops the target, as before.""" + state = SyncState(tmp_path / "s.db") + user = _user() + poison, m2, m3 = _m(80.0, 3), _m(81.0, 2), _m(82.0, 1) + _cap(state, poison, attempts=MAX_RETRY_ATTEMPTS - 1) + + _, errors, garmin, _ = _run(user, state, [poison, m2, m3], fail_weights={80.0}) + assert _garmin_weights(garmin) == [80.0] * 3 # in-run retries, then stop + assert "503" in errors["garmin"] + + _, errors, garmin, _ = _run(user, state, [poison, m2, m3], fail_weights={80.0}) + assert errors == {} + # One try for the capped one, then the newer ones land. + assert _garmin_weights(garmin) == [80.0, 81.0, 82.0] + row = _rows(state)[("garmin", poison.measurement_id)] + assert row["gave_up"] is False + assert row["attempts"] == MAX_RETRY_ATTEMPTS + 1 + assert row["last_newer_success_at"] is not None + + # The cursor moved past it; the next run still reaches back for it alone. + _, errors, garmin, _ = _run(user, state, [poison, m2, m3], fail_weights={80.0}) + assert errors == {} + assert _garmin_weights(garmin) == [80.0] + assert _rows(state)[("garmin", poison.measurement_id)]["attempts"] == MAX_RETRY_ATTEMPTS + 2 + state.close() + + +def test_poison_entry_with_a_healthy_target_is_given_up_after_thirty_days(tmp_path: Path, caplog): + state = SyncState(tmp_path / "s.db") + user = _user() + poison, m2, m3 = _m(80.0, RETRY_ABANDON_DAYS + 1), _m(81.0, 2), _m(82.0, 1) + _cap(state, poison, attempts=1, first_failed=NOW - timedelta(days=RETRY_ABANDON_DAYS + 1)) + + # Old enough, but no newer weigh-in has landed since it first failed. + with caplog.at_level(logging.WARNING, logger="eufy_sync"): + _, errors, garmin, _ = _run(user, state, [poison, m2, m3], fail_weights={80.0}) + assert errors == {} + assert _garmin_weights(garmin) == [80.0, 81.0, 82.0] + assert "Giving up" not in caplog.text + assert _rows(state)[("garmin", poison.measurement_id)]["gave_up"] is False + + # Now both hold: given up before it is sent again. + with caplog.at_level(logging.WARNING, logger="eufy_sync"): + _, errors, garmin, _ = _run(user, state, [poison, m2, m3], fail_weights={80.0}) + assert errors == {} + assert _garmin_weights(garmin) == [] + assert "Giving up on the Garmin upload of 80.00 kg" in caplog.text + assert _rows(state)[("garmin", poison.measurement_id)]["gave_up"] is True + assert state.waiting_upload_retries("default") == {} + + # Given up for good: a backfill that reaches it skips it. + _, _, garmin, _ = _run(user, state, [poison, m2, m3], backfill_days=40) + assert _garmin_weights(garmin) == [] + state.close() + + +def test_intermittent_failures_on_a_capped_entry_do_not_give_it_up_before_thirty_days(tmp_path: Path, caplog): + """Newer weigh-ins keep landing while a capped one hits 5xx after 5xx: + under 30 days that is not enough to give it up, and when Garmin finally + takes it, it is delivered.""" + state = SyncState(tmp_path / "s.db") + user = _user() + head = _m(80.0, 21) + _cap(state, head, first_failed=NOW - timedelta(days=RETRY_ABANDON_DAYS - 1)) + newer = [_m(81.0, 3), _m(82.0, 2), _m(83.0, 1)] + + with caplog.at_level(logging.WARNING, logger="eufy_sync"): + for d in range(1, 4): + _, errors, garmin, _ = _run(user, state, [head] + newer[:d], fail_weights={80.0}) + assert errors == {} + assert _garmin_weights(garmin)[0] == 80.0 + assert "Giving up" not in caplog.text + row = _rows(state)[("garmin", head.measurement_id)] + assert row["gave_up"] is False and row["last_newer_success_at"] is not None + + _, errors, garmin, _ = _run(user, state, [head] + newer) + assert errors == {} + assert _garmin_weights(garmin) == [80.0] + assert state.is_synced("default", head.measurement_id, "garmin") + assert _rows(state) == {} + state.close() + + +def test_capped_entry_stays_pending_when_newer_ones_fail_too(tmp_path: Path): + """The newer measurements fail as well, so the target is down: the + capped entry keeps counting attempts and nothing is given up.""" + state = SyncState(tmp_path / "s.db") + user = _user() + head, m2, m3 = _m(80.0, RETRY_MAX_AGE_DAYS + 1), _m(81.0, 2), _m(82.0, 1) + _cap(state, head, first_failed=NOW - timedelta(days=RETRY_ABANDON_DAYS + 5)) + + _, errors, garmin, _ = _run(user, state, [head, m2, m3], fail_weights={80.0, 81.0, 82.0}, backfill_days=30) + + assert "503" in errors["garmin"] + # It moved past the capped head, then stopped at the next failure. + assert _garmin_weights(garmin) == [80.0] + [81.0] * 3 + rows = _rows(state) + assert rows[("garmin", head.measurement_id)]["gave_up"] is False + assert rows[("garmin", head.measurement_id)]["attempts"] == MAX_RETRY_ATTEMPTS + 1 + assert rows[("garmin", head.measurement_id)]["last_newer_success_at"] is None + assert rows[("garmin", m2.measurement_id)]["attempts"] == 1 + assert ("garmin", m3.measurement_id) not in rows + state.close() + + +def test_capped_entry_that_is_the_newest_stays_pending_and_reports(tmp_path: Path): + """Nothing newer to prove the target healthy: keep it, report the error.""" + state = SyncState(tmp_path / "s.db") + user = _user() + only = _m(80.0, RETRY_MAX_AGE_DAYS + 1) + _cap(state, only) + + _, errors, _, _ = _run(user, state, [only], fail_weights={80.0}, backfill_days=30) + + assert "503" in errors["garmin"] + assert _rows(state)[("garmin", only.measurement_id)]["gave_up"] is False + state.close() + + +def test_unconfirmed_409_is_not_reposted_in_the_same_run(tmp_path: Path): + """RetryNextRunError skips _retry's in-run retries: one POST, then the + measurement waits in the queue, and Garmin stops for the run as for any + uncapped failure.""" + from eufy_sync.sync import RetryNextRunError + + state = SyncState(tmp_path / "s.db") + user = _user() + m1, m2 = _m(80.0, 2), _m(81.0, 1) + + _, errors, garmin, _ = _run( + user, state, [m1, m2], fail_weights={80.0}, + garmin_error=RetryNextRunError("409 lookup kept failing"), + ) + assert _garmin_weights(garmin) == [80.0] + assert "409" in errors["garmin"] + assert _rows(state)[("garmin", m1.measurement_id)]["attempts"] == 1 + state.close() + + +def test_backfilled_old_measurement_is_not_capped_on_its_first_failures(tmp_path: Path): + """A weigh-in from a month ago reached by --backfill gets the full two + weeks from its first failure, not from when it was taken.""" + state = SyncState(tmp_path / "s.db") + user = _user() + old, new = _m(80.0, 30), _m(81.0, 1) + + for _ in range(2): + _, errors, garmin, _ = _run(user, state, [old, new], fail_weights={80.0}, backfill_days=40) + assert "503" in errors["garmin"] + assert _garmin_weights(garmin) == [80.0] * 3 # still holds the target + + row = _rows(state)[("garmin", old.measurement_id)] + assert row["attempts"] == 2 and row["gave_up"] is False + assert not state.is_synced("default", new.measurement_id, "garmin") + state.close() + + +def test_forty_day_outage_then_recovery_delivers_everything_and_gives_up_nothing(tmp_path: Path): + """Garmin is down for 40 days while a weigh-in arrives every day. Every + old entry passes both caps and the abandon age, but nothing newer ever + lands during the outage, so nothing is given up and recovery delivers it + all in order.""" + state = SyncState(tmp_path / "s.db") + user = _user() + before = _m(79.0, 41) + state.record_sync("default", before.measurement_id, before.timestamp.isoformat(), 79.0, NOW.isoformat()) + days = [_m(round(80.0 + d / 10, 2), 40 - d) for d in range(40)] + weights = [m.weight_kg for m in days] + + for d in range(1, 41): + _, errors, _, _ = _run(user, state, days[:d], fail_weights=set(weights)) + assert "503" in errors["garmin"] + # Each entry first failed on the day it was taken, so the old ones + # pass the caps and later runs move past them. + for m in days[:d]: + _backdate(state, m, (NOW - m.timestamp).days) + _, errors, _, _ = _run(user, state, days, fail_weights=set(weights)) + assert "503" in errors["garmin"] + rows = _rows(state).values() + assert sum(1 for r in rows if r["first_failed_at"] < (NOW - timedelta(days=RETRY_ABANDON_DAYS)).isoformat()) >= 9 + assert not any(r["gave_up"] for r in rows) + assert all(r["last_newer_success_at"] is None for r in rows) + + _, errors, garmin, _ = _run(user, state, days) + + assert errors == {} + assert _garmin_weights(garmin) == weights + assert _rows(state) == {} + state.close() + + +def test_reach_back_for_a_waiting_entry_is_bounded(tmp_path: Path, caplog): + """A queued failure inside the window pulls the fetch back to it; one + older than the abandon window plus the margin, whose target never took + anything newer, cannot be given up, stays queued, and does not drag the + fetch further back.""" + state = SyncState(tmp_path / "s.db") + user = _user() + recent = _m(80.0, 20) + _cap(state, recent, attempts=1) + fetches: list[int] = [] + + _run(user, state, [], fetches=fetches) + assert fetches[-1] == int(recent.timestamp.timestamp()) + + ancient = _m(79.0, 90) + _cap(state, ancient, attempts=MAX_RETRY_ATTEMPTS, first_failed=NOW - timedelta(days=90)) + with caplog.at_level(logging.INFO, logger="eufy_sync"): + _run(user, state, [], fetches=fetches) + + floor = NOW.timestamp() - (RETRY_ABANDON_DAYS + RETRY_REACH_BACK_MARGIN_DAYS) * 86400 + assert abs(fetches[-1] - floor) < 120 + assert caplog.text.count("still waiting") == 1 + assert _rows(state)[("garmin", ancient.measurement_id)]["gave_up"] is False + state.close() + + +def test_repair_overrides_a_given_up_entry_and_success_clears_it(tmp_path: Path): + state = SyncState(tmp_path / "s.db") + user = _user() + m1 = _m(80.0, 1) + state.record_upload_failure("default", "garmin", m1.measurement_id, m1.timestamp.isoformat(), 80.0, NOW.isoformat()) + state.give_up_upload_retry("default", "garmin", m1.measurement_id) + + _, _, garmin, _ = _run(user, state, [m1], repair_days=3) + assert _garmin_weights(garmin) == [80.0] + assert _rows(state) == {} + state.close() + + +def test_strava_replays_a_failed_weight_while_it_is_still_the_newest(tmp_path: Path): + state = SyncState(tmp_path / "s.db") + user = _user(garmin=False, strava=True) + latest = _m(80.0, 1) + + _run(user, state, [latest], fail_weights={80.0}) + assert ("strava", latest.measurement_id) in _rows(state) + + _, errors, _, strava = _run(user, state, [latest]) + assert errors == {} + assert [c.args[0] for c in strava.update_weight.call_args_list] == [80.0] + assert _rows(state) == {} + state.close() + + +def test_strava_never_replays_an_older_weight_once_a_newer_one_exists(tmp_path: Path): + state = SyncState(tmp_path / "s.db") + user = _user(garmin=False, strava=True) + older, newer = _m(80.0, 2), _m(79.5, 1) + + _run(user, state, [older], fail_weights={80.0}) + _, errors, _, strava = _run(user, state, [older, newer]) + + assert errors == {} + assert [c.args[0] for c in strava.update_weight.call_args_list] == [79.5] + assert _rows(state) == {} + state.close() + + +def test_newer_strava_failure_replaces_the_older_entry(tmp_path: Path): + state = SyncState(tmp_path / "s.db") + user = _user(garmin=False, strava=True) + older, newer = _m(80.0, 2), _m(79.5, 1) + + _run(user, state, [older], fail_weights={80.0}) + _run(user, state, [older, newer], fail_weights={79.5}) + + assert set(_rows(state)) == {("strava", newer.measurement_id)} + state.close() + + +def test_strava_entry_behind_its_current_weight_is_dropped_without_sending(tmp_path: Path): + """A newer weight reached Strava some other way (e.g. --target strava + with a different fetch); the queued older one must not overwrite it.""" + state = SyncState(tmp_path / "s.db") + user = _user(garmin=False, strava=True) + older, newer = _m(80.0, 2), _m(79.5, 1) + state.record_upload_failure("default", "strava", older.measurement_id, older.timestamp.isoformat(), 80.0, NOW.isoformat()) + state.record_sync("default", newer.measurement_id, newer.timestamp.isoformat(), 79.5, NOW.isoformat(), target="strava") + + _, _, _, strava = _run(user, state, [older], backfill_days=7) + strava.update_weight.assert_not_called() + assert _rows(state) == {} + state.close() + + +def test_dry_run_predicts_without_touching_the_queue(tmp_path: Path, capsys): + state = SyncState(tmp_path / "s.db") + user = _user() + aged, failed = _m(80.0, RETRY_MAX_AGE_DAYS + 1), _m(81.0, 1) + _cap(state, aged, attempts=1, first_failed=NOW - timedelta(days=RETRY_MAX_AGE_DAYS + 1)) + _cap(state, failed, attempts=1) + state.give_up_upload_retry("default", "garmin", failed.measurement_id) + state.record_upload_failure("default", "garmin", "synced", NOW.isoformat(), 79.0, NOW.isoformat()) + state.record_sync("default", "synced", NOW.isoformat(), 79.0, NOW.isoformat()) + before = _rows(state) + + _, _, garmin, _ = _run(user, state, [aged, failed], backfill_days=30, dry_run=True) + + garmin.upload_body_composition.assert_not_called() + assert _rows(state) == before + out = capsys.readouterr().out + # The capped one would be tried without holding back newer ones; the + # given-up one would not be sent at all. + assert "Would retry garmin: 80.0 kg" in out and "newer ones still upload" in out + assert "81.0 kg" not in out + state.close() + + +def test_existing_database_gains_the_retry_table(tmp_path: Path): + """A state.db from 1.14.0 has sync_log (with weight_only) and + pending_upgrades but no upload_retries. Opening it adds the table and + keeps every existing row.""" + db_path = tmp_path / "state.db" + conn = sqlite3.connect(str(db_path)) + conn.executescript(""" + CREATE TABLE sync_log ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + user_name TEXT NOT NULL, + eufy_measurement_id TEXT NOT NULL, + measurement_timestamp TEXT NOT NULL, + weight_kg REAL, + target TEXT NOT NULL DEFAULT 'garmin', + synced_at TEXT NOT NULL, + response TEXT, + weight_only INTEGER NOT NULL DEFAULT 0, + UNIQUE(user_name, eufy_measurement_id, target) + ); + CREATE TABLE pending_upgrades ( + user_name TEXT NOT NULL, + measurement_id TEXT NOT NULL, + previous_measurement_id TEXT NOT NULL, + measurement_json TEXT NOT NULL, + PRIMARY KEY(user_name, measurement_id) + ); + """) + conn.execute( + "INSERT INTO sync_log (user_name, eufy_measurement_id, measurement_timestamp, weight_kg, target, synced_at)" + " VALUES (?, ?, ?, ?, ?, ?)", + ("default", "m1", "2026-09-01T08:00:00+00:00", 80.0, "garmin", "2026-09-01T08:01:00+00:00"), + ) + conn.commit() + conn.close() + + state = SyncState(db_path) + assert state.is_synced("default", "m1", "garmin") + assert state.get_upload_retries("default") == [] + assert state.record_upload_failure("default", "garmin", "m2", "2026-09-02T08:00:00+00:00", 80.1, NOW.isoformat()) == 1 + state.close() + + # Reopening is a no-op for the schema and keeps the queued entry. + state = SyncState(db_path) + assert state.waiting_upload_retries("default") == {"garmin": 1} + state.close() + + +def test_retry_table_without_the_newer_success_column_gains_it(tmp_path: Path): + """An upload_retries table from an earlier 1.15 build lacks + last_newer_success_at. Opening the database adds it as NULL, which can + only delay a give-up, and keeps the queued rows.""" + db_path = tmp_path / "state.db" + conn = sqlite3.connect(str(db_path)) + conn.executescript(""" + CREATE TABLE upload_retries ( + user_name TEXT NOT NULL, + target TEXT NOT NULL, + measurement_id TEXT NOT NULL, + measurement_timestamp TEXT NOT NULL, + weight_kg REAL, + first_failed_at TEXT NOT NULL, + last_failed_at TEXT NOT NULL, + attempts INTEGER NOT NULL DEFAULT 1, + gave_up INTEGER NOT NULL DEFAULT 0, + PRIMARY KEY(user_name, target, measurement_id) + ); + """) + conn.execute( + "INSERT INTO upload_retries VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)", + ("default", "garmin", "m1", "2026-08-01T08:00:00+00:00", 80.0, + "2026-08-01T09:00:00+00:00", "2026-08-01T09:00:00+00:00", 90, 0), + ) + conn.commit() + conn.close() + + state = SyncState(db_path) + [row] = state.get_upload_retries("default") + assert row["attempts"] == 90 and row["last_newer_success_at"] is None + state.note_newer_upload("default", "garmin", NOW, NOW.isoformat()) + assert state.get_upload_retries("default")[0]["last_newer_success_at"] == NOW.isoformat() + state.close() + + +def test_status_line_counts_waiting_uploads(tmp_path: Path): + state = SyncState(tmp_path / "s.db") + user = _user(strava=True) + assert _retry_queue_line(state, user) is None + + for mid in ("a", "b"): + state.record_upload_failure("default", "garmin", mid, NOW.isoformat(), 80.0, NOW.isoformat()) + assert _retry_queue_line(state, user) == "2 uploads waiting to retry (Garmin)" + + state.record_upload_failure("default", "strava", "c", NOW.isoformat(), 80.0, NOW.isoformat()) + state.give_up_upload_retry("default", "garmin", "b") + assert _retry_queue_line(state, user) == "2 uploads waiting to retry (Garmin, Strava)" + + # A target removed from the config will never retry, so it is not shown. + garmin_only = _user() + assert _retry_queue_line(state, garmin_only) == "1 upload waiting to retry (Garmin)" + state.close() + + +def test_status_shows_the_retry_line(tmp_path: Path, capsys): + from eufy_sync.cli.status import _show_status + + state = SyncState(tmp_path / "s.db") + user = UserConfig( + name="default", + eufy=EufyConfig(email="e@example.com", password="pw"), + zwift=ZwiftConfig(email="z@example.com", password="pw"), + ) + state.record_upload_failure("default", "zwift", "a", NOW.isoformat(), 80.0, NOW.isoformat()) + + with patch("eufy_sync.eufy_client.EufyClient") as eufy, \ + patch("eufy_sync.cli.status._zwift_token_status", return_value={"state": "valid"}): + eufy.return_value.token_status.return_value = {"state": "valid", "days_remaining": 10} + _show_status(state, [user]) + + assert "1 upload waiting to retry (Zwift)" in capsys.readouterr().out + state.close()