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()