diff --git a/routstr/wallet.py b/routstr/wallet.py index 8828f04c..cbf0552a 100644 --- a/routstr/wallet.py +++ b/routstr/wallet.py @@ -58,6 +58,8 @@ class TokenConsumedError(Exception): # httpx base classes cover their subclasses. HTTPStatusError is excluded on # purpose — that means the mint answered, just with an error status. _MINT_TRANSPORT_COOLDOWN_SECONDS = 30.0 +_MINT_RATE_LIMIT_BASE_COOLDOWN_SECONDS = 60.0 +_MINT_RATE_LIMIT_MAX_COOLDOWN_SECONDS = 7 * 60 * 60 _TRANSPORT_EXC_TYPES: tuple[type[BaseException], ...] = ( httpx.NetworkError, @@ -90,6 +92,7 @@ class _MintRateGuard: ) self._cooldown_until = 0.0 self._cooldown_reason: str | None = None + self._consecutive_rate_limits = 0 self._needs_probe = False self._probe_lock = asyncio.Lock() @@ -103,6 +106,25 @@ class _MintRateGuard: self._cooldown_reason = reason self._needs_probe = True + def apply_rate_limit_cooldown(self, retry_after: float | None = None) -> float: + remaining = self.cooldown_remaining() + if remaining > 0 and self._cooldown_reason == "rate_limited": + minimum = min( + _MINT_RATE_LIMIT_MAX_COOLDOWN_SECONDS, + max(_MINT_RATE_LIMIT_BASE_COOLDOWN_SECONDS, retry_after or 0.0), + ) + if minimum > remaining: + self.apply_cooldown(minimum, reason="rate_limited") + return minimum + return remaining + + self._consecutive_rate_limits += 1 + base = max(_MINT_RATE_LIMIT_BASE_COOLDOWN_SECONDS, retry_after or 0.0) + multiplier = 2 ** min(self._consecutive_rate_limits - 1, 10) + delay = min(_MINT_RATE_LIMIT_MAX_COOLDOWN_SECONDS, base * multiplier) + self.apply_cooldown(delay, reason="rate_limited") + return delay + def cooldown_remaining(self) -> float: return max(0.0, self._cooldown_until - time.monotonic()) @@ -132,9 +154,16 @@ class _MintRateGuard: try: result = await factory() except Exception as error: - # Keep queued callers behind the probe while _mint_operation applies - # the precise Retry-After/backoff from this failure. - self.apply_cooldown(1.0) + # Keep queued callers behind the probe. Handle rate limits here so + # the next exponential step is recorded before another waiter can + # acquire the probe lock. + if _is_mint_rate_limited(error): + retry_after = None + if isinstance(error, httpx.HTTPStatusError): + retry_after = _parse_retry_after(error.response.headers) + self.apply_rate_limit_cooldown(retry_after) + else: + self.apply_cooldown(1.0) logger.warning( "Mint cooldown probe failed", extra={ @@ -142,6 +171,8 @@ class _MintRateGuard: "mint_url": self._mint_url, "error": str(error), "error_type": type(error).__name__, + "cooldown_seconds": round(self.cooldown_remaining(), 2), + "consecutive_rate_limits": self._consecutive_rate_limits, }, ) raise @@ -149,6 +180,7 @@ class _MintRateGuard: self._needs_probe = False self._cooldown_until = 0.0 self._cooldown_reason = None + self._consecutive_rate_limits = 0 logger.warning( "Mint cooldown probe succeeded; restoring normal concurrency", extra={ @@ -260,8 +292,9 @@ async def _mint_operation( retry_after = _parse_retry_after(exc.response.headers) if retry_after is not None: backoff = max(retry_after, backoff) + cooldown = backoff if guard is not None: - guard.apply_cooldown(backoff, reason="rate_limited") + cooldown = guard.apply_rate_limit_cooldown(backoff) # When the caller has a fallback strategy (trusted-mint # list), re-raise immediately so the caller can try the next @@ -272,7 +305,10 @@ async def _mint_operation( extra={ "op_name": op_name, "mint_url": mint_url, - "cooldown_seconds": round(backoff, 2), + "cooldown_seconds": round(cooldown, 2), + "consecutive_rate_limits": guard._consecutive_rate_limits + if guard is not None + else attempt + 1, }, ) raise @@ -285,11 +321,14 @@ async def _mint_operation( "op_name": op_name, "mint_url": mint_url, "attempt": attempt + 1, - "cooldown_seconds": round(backoff, 2), + "cooldown_seconds": round(cooldown, 2), + "consecutive_rate_limits": guard._consecutive_rate_limits + if guard is not None + else attempt + 1, }, ) if guard is None: - await asyncio.sleep(backoff) + await asyncio.sleep(cooldown) raise RuntimeError(f"{op_name}: exhausted retries unexpectedly") @@ -1648,7 +1687,11 @@ async def fetch_all_balances( if connection_failure else "mint_error" ) - if connection_failure or rate_limited: + if rate_limited: + _MintRateGuard.get(mint_url).apply_rate_limit_cooldown( + _BALANCE_FETCH_RETRY_SECONDS + ) + elif connection_failure: _MintRateGuard.get(mint_url).apply_cooldown( _BALANCE_FETCH_RETRY_SECONDS, reason=error_code ) diff --git a/tests/unit/test_wallet.py b/tests/unit/test_wallet.py index 39d87899..46339065 100644 --- a/tests/unit/test_wallet.py +++ b/tests/unit/test_wallet.py @@ -1509,6 +1509,33 @@ async def test_mint_rate_guard_waits_for_adaptive_cooldown() -> None: operation.assert_awaited_once() +@pytest.mark.asyncio +async def test_mint_rate_guard_exponentially_backs_off_repeated_429s() -> None: + from routstr.wallet import _MintRateGuard + + guard = _MintRateGuard("http://mint:3338", 4) + expected_delays = [60, 120, 240, 480, 960, 1920, 3840, 7680, 15360, 25200] + now = 0.0 + + with patch("routstr.wallet.time.monotonic") as monotonic: + for index, expected in enumerate(expected_delays, start=1): + monotonic.return_value = now + assert guard.apply_rate_limit_cooldown(60) == expected + assert guard._consecutive_rate_limits == index + if index == 1: + # Concurrent responses from the same 429 wave do not escalate + # the retry count before the first cooldown probe. + assert guard.apply_rate_limit_cooldown(60) == expected + assert guard._consecutive_rate_limits == 1 + now += expected + 1 + + monotonic.return_value = now + operation = AsyncMock(return_value="ok") + assert await guard.run(operation) == "ok" + assert guard._consecutive_rate_limits == 0 + assert guard.apply_rate_limit_cooldown(60) == 60 + + @pytest.mark.asyncio async def test_mint_rate_guard_allows_one_probe_after_cooldown() -> None: from routstr.wallet import _MintRateGuard