diff --git a/CHANGELOG.md b/CHANGELOG.md index ebec53e..a968da6 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,62 @@ # Changelog +## [3.2.0] - 2026-09-26 + +### Added + +- **The client reads the rate-limit headers, and acts on them.** Both meters, + the per-minute REST one and the separate hourly history budget, are exposed + on `client.rate_limit`: + + ```python + tp.rate_limit.rest.remaining # 299 + tp.rate_limit.history.remaining # on a history call + tp.rate_limit.rest.seconds_until_reset() + ``` + + Every field is optional, and `None` means the server did not say rather than + "nothing left". Use `.exhausted`, which is true only when the server said + zero. The per-minute figures ride most responses; the hourly history ones are + withheld from anything a shared cache may store, because they are per-caller; + an unmetered plan advertises nothing. A response served from a cache is + ignored entirely, because its figures belong to whoever populated the entry. `reset` is a + relative countdown frozen when it was read, so `seconds_until_reset()` ages + it rather than returning a stale number. + + When a response says the window is spent, the next request now waits for the + advertised reset instead of sending one that is certain to be refused, and to + spend a unit of budget being refused. `RetryConfig(respect_remaining=False)` + turns it off. + + The hourly history budget is new on the wire; before it there was nothing to + read. + +### Changed + +- **Calls may now block before sending.** When the server has said your window + is spent, or has issued a 429 that is still in force, the client waits rather + than sending a request that is certain to be refused. A call that used to + return in 200ms can now take up to `retry.max_retry_after` (120s) first. That + is a TOTAL across the call, not per wait: the shared 429 gate and the + spent-window wait stack, and before the budget existed a 429 carrying both a + `Retry-After` and a spent window blocked for 180 seconds under a 120 second + cap. Turn the two halves off with `RetryConfig(respect_remaining=False)` and + `RetryConfig(respect_429=False)`. + +### Fixed + +- **`respect_429=False` did not opt out.** It raised the error the caller asked + for and then held their NEXT call for the full `Retry-After` anyway, because + the shared gate was closed regardless of the setting. + +- **A 429 was waited out once per in-flight request.** The wait belongs to the + caller, not to whichever request met it, so ten concurrent requests each + slept their own `Retry-After` and then retried at the same instant, + re-tripping the limit together. It is now taken once, on a gate shared by the + whole client, with a little jitter so the waiters do not wake in unison. A + shorter wait arriving while a longer one is in force no longer brings the + gate forward. + ## [3.1.0] - 2026-09-23 ### Added diff --git a/README.md b/README.md index 416da33..0c58210 100644 --- a/README.md +++ b/README.md @@ -93,7 +93,7 @@ Both `ThemeParks` and `AsyncThemeParks` take the same keyword-only options: | `api_key` | `str \| None` | `None` | Sent as the `x-api-key` header. Needed for anything beyond the free tier: deeper history, higher rate limits. | | `user_agent` | `str \| None` | `themeparks-sdk-py/` | Sent as the `User-Agent` header. Set this to identify your app. | | `timeout` | `float` (seconds) | `10.0` | Per-request timeout. | -| `retry` | `RetryConfig \| None` | `RetryConfig(max_retries=3, respect_429=True, max_retry_after=120.0)` | Retry/backoff behavior. `max_retries` is N retries beyond the first attempt (so N+1 total calls). `max_retry_after` is the longest `Retry-After` the client will sleep through; past it you get `RateLimitError` instead of a silent wait. | +| `retry` | `RetryConfig \| None` | `RetryConfig(max_retries=3, respect_429=True, max_retry_after=120.0, respect_remaining=True)` | Retry/backoff behavior. `max_retries` is N retries beyond the first attempt (so N+1 total calls). `max_retry_after` is the TOTAL the client will block for within one call, across both the shared 429 gate and any spent-window wait. Past a single `Retry-After` that long you get `RateLimitError` instead of a silent wait. | | `cache` | `Cache \| CacheConfig \| bool \| None` | `True` (in-memory LRU) | See **Caching** below. `False` disables caching entirely. | Example: @@ -216,6 +216,50 @@ remaining keys are whatever fields that variant carries. timezone-aware `datetime`, honoring the entity's IANA timezone for naive inputs. +## Rate limits + +`client.rate_limit` is a read-only property, not a constructor option. + +The API meters requests per minute, and history requests again per hour. Both +are advertised on every response that can carry them, and the client reads +them: + +```python +with ThemeParks(api_key=KEY) as tp: + tp.entity(park_id).live() + + print(tp.rate_limit.rest.remaining) # 299 + print(tp.rate_limit.rest.seconds_until_reset()) + print(tp.rate_limit.history.remaining) # on a history call +``` + +**`None` means the server did not say, never "nothing left".** Use +`.exhausted`, which is true only when the server actually said zero. + +Which figures you get depends on the response: + +- The **per-minute** figures ride most responses, anonymous ones included. +- The **hourly history** figures are withheld from anything a shared cache may + store, because they are per-caller and a cache would hand one caller's budget + to another. In practice you get them on calls made with a key. +- An **unmetered plan** advertises nothing at all. + +A response served from a cache is ignored entirely. Its figures belong to +whoever populated the entry and its countdown is already wrong: a cached +`remaining: 0` would otherwise make the client sleep out someone else's +window. + +The client also acts on what it reads. When a response says the window is +spent, the next request waits for the advertised reset rather than sending a +request that is certain to be refused, and to cost a unit of budget being +refused. Turn that off with `RetryConfig(respect_remaining=False)`. + +**A 429 is held once for the whole client.** The wait belongs to the caller, +not to whichever request happened to meet it, so it goes on a shared gate with +a little jitter. Without that, ten concurrent requests each sleep their own +copy of `Retry-After` and then all retry at the same instant, re-tripping the +limit together. + ## History `tp.entity(id).history` reads the archive. Both methods page for you and yield diff --git a/docs/api/ratelimits.md b/docs/api/ratelimits.md new file mode 100644 index 0000000..539f8c5 --- /dev/null +++ b/docs/api/ratelimits.md @@ -0,0 +1,17 @@ +# Rate limits + +The API meters requests per minute, and history requests again per hour. Both +are read off every response that carries them and exposed on the client. + +`None` means the server did not say, never "nothing left". An unmetered plan +advertises no figures, and neither does a publicly cacheable response, because +the numbers are per-caller and a shared cache would hand one caller's budget to +another. Anonymous calls therefore carry nothing; calls made with a key do. + +::: themeparks.RateLimits + options: + heading_level: 2 + +::: themeparks.RateLimit + options: + heading_level: 2 diff --git a/mkdocs.yml b/mkdocs.yml index bb3006b..5ae916e 100644 --- a/mkdocs.yml +++ b/mkdocs.yml @@ -59,6 +59,7 @@ nav: - Client: api/client.md - Entity: api/entity.md - History: api/history.md + - Rate limits: api/ratelimits.md - Destinations: api/destinations.md - Raw client: api/raw.md - Helpers: api/helpers.md diff --git a/pyproject.toml b/pyproject.toml index 297b2e0..018dabf 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "hatchling.build" [project] name = "themeparks" -version = "3.1.0" +version = "3.2.0" description = "Official SDK for the ThemeParks.wiki API" readme = "README.md" requires-python = ">=3.9" diff --git a/tests/unit/test_history.py b/tests/unit/test_history.py index 9b31a88..148701e 100644 --- a/tests/unit/test_history.py +++ b/tests/unit/test_history.py @@ -437,4 +437,8 @@ def test_an_ordinary_rest_429_is_still_ridden_out(self): with pytest.raises(RateLimitError) as caught: list(tp.entity("park-1").history.days("2026-09-01", "2026-09-02")) assert not isinstance(caught.value, BudgetExhaustedError) - assert slept == [2.0, 2.0, 2.0] + # Jittered: the gate is shared, so without a little spread every + # waiter would wake at the same instant and re-trip the limit + # together. One wait per retry, taken once, never doubled. + assert len(slept) == 3 + assert all(2.0 <= s < 2.3 for s in slept), slept diff --git a/tests/unit/test_ratelimit.py b/tests/unit/test_ratelimit.py new file mode 100644 index 0000000..1e3ca4c --- /dev/null +++ b/tests/unit/test_ratelimit.py @@ -0,0 +1,419 @@ +"""Reading the two budgets, and staying inside them. + +The SDK used to read exactly one header, `Retry-After`, and only after a 429 +had already happened. It could tell you that you had run out, never that you +were about to. These cover the three things that changed: the figures are +read, absence is not confused with zero, and the wait a 429 imposes is taken +ONCE for the whole client rather than once per in-flight request. +""" + +import time + +import httpx +import pytest + +import themeparks._ratelimit as rl +from themeparks import RateLimitError, RateLimits, RetryConfig, ThemeParks +from themeparks._ratelimit import Gate, RateLimit, read_rate_limits + +REST = { + "RateLimit-Limit": "300", + "RateLimit-Policy": "300;w=60", + "RateLimit-Remaining": "299", + "RateLimit-Reset": "60", +} +HISTORY = { + "RateLimit-History-Limit": "600", + "RateLimit-History-Policy": "600;w=3600", + "RateLimit-History-Remaining": "599", + "RateLimit-History-Reset": "3412", +} + + +class TestReadingHeaders: + def test_reads_the_rest_meter(self): + out = read_rate_limits(REST, RateLimits()) + assert (out.rest.limit, out.rest.remaining, out.rest.reset) == (300, 299, 60) + assert out.rest.policy == "300;w=60" + + def test_reads_the_history_meter(self): + out = read_rate_limits(HISTORY, RateLimits()) + assert (out.history.limit, out.history.remaining, out.history.reset) == (600, 599, 3412) + + def test_the_two_meters_do_not_bleed_into_each_other(self): + # "ratelimit-history-limit" also starts with "ratelimit-", so a naive + # prefix match reads the hourly figure as the per-minute one and a + # client paces itself against the wrong window. + out = read_rate_limits({**REST, **HISTORY}, RateLimits()) + assert out.rest.limit == 300 + assert out.history.limit == 600 + assert out.rest.reset == 60 + assert out.history.reset == 3412 + + def test_history_headers_alone_do_not_invent_a_rest_meter(self): + out = read_rate_limits(HISTORY, RateLimits()) + assert out.rest.limit is None + + def test_a_response_mentioning_neither_keeps_what_we_knew(self): + # Most responses mention one meter or, if publicly cacheable, neither. + # Overwriting with blanks would mean the last cacheable response + # erased everything the client had learned. + known = read_rate_limits({**REST, **HISTORY}, RateLimits()) + after = read_rate_limits({"content-type": "application/json"}, known) + assert after.rest.remaining == 299 + assert after.history.remaining == 599 + + def test_header_case_does_not_matter(self): + out = read_rate_limits({k.lower(): v for k, v in REST.items()}, RateLimits()) + assert out.rest.remaining == 299 + + def test_a_malformed_value_is_unknown_rather_than_zero(self): + # An int() that fails must not become 0, or the client would hold + # forever waiting for a window it invented. + out = read_rate_limits({**REST, "RateLimit-Remaining": "lots"}, RateLimits()) + assert out.rest.remaining is None + assert out.rest.exhausted is False + + +class TestACachedResponseSaysNothing: + """Its figures belong to whoever populated the entry. + + The server withholds the HISTORY figures from anything a shared cache may + store, but the per-minute ones ride those responses. Measured against + production: three consecutive calls returning `age: 9` and an unmoving + `remaining: 285`. A cached `remaining: 0` would make the client sleep out + a window belonging to someone else. + """ + + def test_a_cache_hit_is_ignored(self): + known = read_rate_limits(REST, RateLimits()) + after = read_rate_limits({**REST, "RateLimit-Remaining": "0", "Age": "1713"}, known) + assert after.rest.remaining == 299, "took a cached caller's figures" + + def test_a_fresh_response_is_recorded(self): + # A cache MISS carries no Age at all, which is the path that matters: + # confirmed against production, a MISS returns the figures and a HIT + # returns them stale. + out = read_rate_limits(REST, RateLimits()) + assert out.rest.remaining == 299 + + def test_age_zero_is_fresh(self): + out = read_rate_limits({**REST, "Age": "0"}, RateLimits()) + assert out.rest.remaining == 299 + + def test_a_cache_hit_does_not_erase_what_we_knew(self): + known = read_rate_limits({**REST, **HISTORY}, RateLimits()) + after = read_rate_limits({"Age": "60"}, known) + assert after.rest.remaining == 299 + assert after.history.remaining == 599 + + +class TestAbsenceIsNotZero: + def test_unknown_remaining_is_not_exhausted(self): + # Anonymous responses carry no figures at all, because they are + # publicly cacheable and the figures are per-caller. Reading that as + # "nothing left" would stall every anonymous client permanently. + assert RateLimit().exhausted is False + + def test_zero_remaining_is_exhausted(self): + assert RateLimit(remaining=0).exhausted is True + + # Both of these pass an explicit `now` rather than reading a clock. The + # method takes one precisely so these can be exact; asserting a range + # around real wall-clock time makes a gate test flaky under CPU load. + def test_reset_counts_down_from_when_it_was_read(self): + # `reset` is relative and frozen at observed_at. Using it later + # without ageing it is how a client waits far longer than it needs to. + meter = RateLimit(reset=60, observed_at=1000.0) + assert meter.seconds_until_reset(now=1050.0) == 10.0 + + def test_an_expired_window_never_reports_negative(self): + meter = RateLimit(reset=5, observed_at=1000.0) + assert meter.seconds_until_reset(now=1100.0) == 0.0 + + def test_unknown_reset_has_no_countdown(self): + assert RateLimit(remaining=0).seconds_until_reset() is None + + +class TestGate: + """One wait for the whole client, not one per in-flight request.""" + + def test_an_open_gate_costs_nothing(self): + assert Gate().wait_seconds() == 0.0 + + def test_a_closed_gate_holds_every_caller(self): + gate = Gate() + gate.close_for(5) + first, second = gate.wait_seconds(), gate.wait_seconds() + assert first > 5.0 and second > 5.0 + + def test_waiters_are_jittered_so_they_do_not_wake_together(self): + # Waking in unison is the other half of the thundering herd: the + # sleeps are shared, then everyone retries at the same instant and + # re-trips the limit. + gate = Gate() + gate.close_for(5) + waits = {gate.wait_seconds() for _ in range(20)} + assert len(waits) > 1 + + def test_a_shorter_wait_never_brings_the_gate_forward(self): + # A 2-second Retry-After arriving while a 60-second one is in force + # would otherwise release the herd early. + gate = Gate() + gate.close_for(60) + gate.close_for(2) + assert gate.wait_seconds() > 55.0 + + +class TestThroughTheClient: + def _client(self, headers, slept=None): + def handler(request): + return httpx.Response(200, headers=headers, json={"destinations": []}) + + tp = ThemeParks(transport=httpx.MockTransport(handler), cache=False) + if slept is not None: + tp.raw._t._sleep = slept.append + return tp + + def test_a_real_call_records_both_meters(self): + tp = self._client({**REST, **HISTORY}) + tp.destinations.list() + assert tp.rate_limit.rest.remaining == 299 + assert tp.rate_limit.history.remaining == 599 + + def test_before_any_call_everything_is_unknown(self): + tp = self._client(REST) + assert tp.rate_limit.rest.limit is None + assert tp.rate_limit.rest.exhausted is False + + def test_a_spent_window_is_waited_out_rather_than_walked_into(self): + # Sending into a window the server said is spent is a guaranteed 429 + # that also costs a unit of budget to refuse. + slept: list[float] = [] + tp = self._client({**REST, "RateLimit-Remaining": "0", "RateLimit-Reset": "7"}, slept) + tp.destinations.list() # learns remaining 0 + tp.destinations.list() # should hold first + assert slept, "walked straight into a window the server said was spent" + # Upper bound allows the spread now applied to this path: without it + # every waiter woke at the same absolute instant. + assert 0 < slept[0] <= 7.0 + 0.25 + + def test_an_unknown_remaining_never_holds(self): + slept: list[float] = [] + tp = self._client({"content-type": "application/json"}, slept) + tp.destinations.list() + tp.destinations.list() + assert slept == [] + + def test_respect_remaining_can_be_turned_off(self): + slept: list[float] = [] + headers = {**REST, "RateLimit-Remaining": "0", "RateLimit-Reset": "7"} + + def handler(request): + return httpx.Response(200, headers=headers, json={"destinations": []}) + + tp = ThemeParks( + transport=httpx.MockTransport(handler), + cache=False, + retry=RetryConfig(respect_remaining=False), + ) + tp.raw._t._sleep = slept.append + tp.destinations.list() + tp.destinations.list() + assert slept == [] + + +class TestThroughTheCache: + """Caching is ON by default, and it wraps the transport. + + Every other test here passes `cache=False`, so none of them touched the + wrapper. The first real call against production raised AttributeError: + the caching transport had no `rate_limit` to forward. Tests that all take + the same non-default path cover a shape the users do not have. + """ + + def _client(self, headers): + def handler(request): + return httpx.Response(200, headers=headers, json={"destinations": []}) + + return ThemeParks(transport=httpx.MockTransport(handler)) # cache default: ON + + def test_the_figures_survive_the_cache_wrapper(self): + tp = self._client(REST) + tp.destinations.list() + assert tp.rate_limit.rest.remaining == 299 + + def test_a_cache_hit_keeps_the_last_known_figures(self): + # A hit sends no request and so learns nothing, which is right: it + # spent no budget either, so the previous figures still stand. + calls: list[int] = [] + + def handler(request): + calls.append(1) + return httpx.Response(200, headers=REST, json={"destinations": []}) + + tp = ThemeParks(transport=httpx.MockTransport(handler)) + tp.destinations.list() + tp.destinations.list() + assert len(calls) == 1, "expected the second call to be a cache hit" + assert tp.rate_limit.rest.remaining == 299 + + +class TestTheOptOutsActuallyOptOut: + """An advertised switch that does not switch anything is worse than none. + + `respect_429=False` raised the RateLimitError the caller asked for, and + then closed the shared gate anyway, so their NEXT call blocked for the + full Retry-After with no way to stop it. The setting says "do not wait on + a 429"; the gate is a wait on a 429. + """ + + def _client(self, retry, slept): + def handler(request): + return httpx.Response(429, headers={"retry-after": "45"}, json={}) + + tp = ThemeParks(transport=httpx.MockTransport(handler), cache=False, retry=retry) + tp.raw._t._sleep = slept.append + return tp + + def test_respect_429_false_never_sleeps_even_on_a_later_call(self): + slept: list[float] = [] + tp = self._client(RetryConfig(respect_429=False), slept) + for _ in range(3): + with pytest.raises(RateLimitError): + tp.destinations.list() + assert slept == [], "opted out of 429 waiting and waited anyway" + + def test_respect_429_true_still_holds_the_gate(self): + # The opt-out must not have disabled the feature for everyone else. + slept: list[float] = [] + tp = self._client(RetryConfig(max_retries=0), slept) + with pytest.raises(RateLimitError): + tp.destinations.list() + with pytest.raises(RateLimitError): + tp.destinations.list() + assert slept, "the gate stopped holding for callers who did want it" + + +class TestTheCapBoundsTheWholeCall: + """`max_retry_after` promises a ceiling. It has to be a real one. + + Two separate self-initiated holds exist -- the shared gate, and the + spent-window wait -- and they stack. A 429 carrying BOTH a `Retry-After` + and `RateLimit-Remaining: 0` slept 5s at the gate and then 55s for the + window, three times over: 180 seconds inside one call whose cap was 120. + Each leg was under the cap, so the per-leg check never fired. + + Nothing caught it because no test sent a 429 carrying rate-limit headers, + and because a fake sleep that does not advance the clock cannot show a + cumulative total at all. This one advances an injected clock, which is + what the real world does. + """ + + def _run(self, cap, headers): + now = [1000.0] + original = rl.time.monotonic + rl.time.monotonic = lambda: now[0] + try: + slept: list[float] = [] + + def sleep(seconds): + slept.append(seconds) + now[0] += seconds + + def handler(request): + return httpx.Response(429, headers=headers, json={}) + + tp = ThemeParks( + transport=httpx.MockTransport(handler), + cache=False, + retry=RetryConfig(max_retry_after=cap), + ) + tp.raw._t._sleep = sleep + tp.raw._t._gate._until = 0.0 + with pytest.raises(RateLimitError): + tp.destinations.list() + return slept + finally: + rl.time.monotonic = original + + REAL_429 = { + "retry-after": "5", + "RateLimit-Limit": "300", + "RateLimit-Remaining": "0", + "RateLimit-Reset": "60", + } + + def test_the_total_never_exceeds_the_cap(self): + slept = self._run(120.0, self.REAL_429) + assert sum(slept) <= 120.0 + 0.01, f"blocked {sum(slept):.0f}s under a 120s cap" + + def test_a_smaller_cap_binds_harder(self): + slept = self._run(30.0, self.REAL_429) + assert sum(slept) <= 30.0 + 0.01, f"blocked {sum(slept):.0f}s under a 30s cap" + + def test_it_still_waits_when_there_is_budget(self): + # The cap must bound the feature, not disable it. + slept = self._run(120.0, self.REAL_429) + assert slept, "stopped waiting altogether" + assert sum(slept) > 5.0, "only paid the gate, never the window" + + def test_no_meaningless_micro_sleeps(self): + # Floating-point residue was producing a trailing sleep of ~1e-14. + slept = self._run(120.0, self.REAL_429) + assert all(s > 0.001 for s in slept), slept + + +class TestWaitersDoNotWakeAsOne: + """The gate exists for concurrency, and nothing measured concurrency. + + Every other test here is sequential, so the one property the gate is for + -- N waiters released without re-tripping the limit together -- was never + asserted. Both release paths are covered, because the spent-window one had + no spread at all: ten waiters derived the same deadline from the same + observed_at and left inside the same millisecond, the tightest burst in + the client, on the branch that exists to avoid a 429. + """ + + def test_gate_waiters_are_spread(self): + gate = rl.Gate() + gate.close_for(5.0) + waits = [gate.wait_seconds() for _ in range(20)] + assert len(set(waits)) > 1, "every waiter would wake at the same instant" + assert all(5.0 <= w < 5.3 for w in waits), waits + + def test_the_spent_window_path_is_spread_too(self): + # Against a FROZEN clock. With a live one `left` varies by itself as + # time passes between runs, so the set is distinct with or without + # spread and the test measures nothing -- the same vacuity that hid + # the missing spread here in the first place. + original = rl.time.monotonic + rl.time.monotonic = lambda: 1000.0 + try: + seen = set() + for _ in range(20): + slept: list[float] = [] + + def handler(request): + return httpx.Response( + 200, + headers={**REST, "RateLimit-Remaining": "0", "RateLimit-Reset": "7"}, + json={"destinations": []}, + ) + + tp = ThemeParks(transport=httpx.MockTransport(handler), cache=False) + tp.raw._t._sleep = slept.append + tp.destinations.list() + tp.destinations.list() + if slept: + seen.add(round(slept[0], 6)) + assert len(seen) > 1, f"all waiters left at the same instant: {seen}" + assert all(7.0 <= w < 7.3 for w in seen), seen + finally: + rl.time.monotonic = original + + def test_the_spread_is_bounded(self): + # It must not be mistaken for the wait itself. + gate = rl.Gate() + gate.close_for(1.0) + assert max(gate.wait_seconds() for _ in range(50)) < 1.3 diff --git a/tests/unit/test_ratelimit_async.py b/tests/unit/test_ratelimit_async.py new file mode 100644 index 0000000..355f074 --- /dev/null +++ b/tests/unit/test_ratelimit_async.py @@ -0,0 +1,110 @@ +"""The async client's half of the rate-limit feature. + +It had NO tests. Every assertion in test_ratelimit.py drives the sync client, +so the async transport's hold, the async caching wrapper's `rate_limit` +property and AsyncThemeParks.rate_limit were all reachable only by reading. + +That matters here specifically: the caching wrapper with no `rate_limit` to +forward is the exact bug that crashed on the first real call against +production. It was fixed in both wrappers and tested in one. +""" + +import httpx +import pytest + +from themeparks import AsyncThemeParks, RateLimitError, RetryConfig + +REST = { + "RateLimit-Limit": "300", + "RateLimit-Policy": "300;w=60", + "RateLimit-Remaining": "299", + "RateLimit-Reset": "60", +} + + +def _client(headers, *, cache=True, retry=None, status=200): + def handler(request): + return httpx.Response(status, headers=headers, json={"destinations": []}) + + kwargs = {"transport": httpx.MockTransport(handler), "cache": cache} + if retry is not None: + kwargs["retry"] = retry + return AsyncThemeParks(**kwargs) + + +class TestAsyncReadsTheMeters: + async def test_a_real_call_records_them(self): + tp = _client(REST, cache=False) + await tp.destinations.list() + assert tp.rate_limit.rest.remaining == 299 + + async def test_it_survives_the_caching_wrapper(self): + # The sync wrapper's missing property raised AttributeError on the + # first production call. The async one is the same shape. + tp = _client(REST) + await tp.destinations.list() + assert tp.rate_limit.rest.remaining == 299 + + async def test_before_any_call_everything_is_unknown(self): + tp = _client(REST, cache=False) + assert tp.rate_limit.rest.limit is None + assert tp.rate_limit.rest.exhausted is False + + async def test_a_cached_response_says_nothing(self): + tp = _client({**REST, "Age": "1713"}, cache=False) + await tp.destinations.list() + assert tp.rate_limit.rest.remaining is None + + +class TestAsyncHolds: + async def _slept(self, headers, retry=None, status=200): + slept: list[float] = [] + tp = _client(headers, cache=False, retry=retry, status=status) + + async def sleep(seconds): + slept.append(seconds) + + tp.raw._t._sleep = sleep + return tp, slept + + async def test_a_spent_window_is_waited_out(self): + tp, slept = await self._slept({**REST, "RateLimit-Remaining": "0", "RateLimit-Reset": "7"}) + await tp.destinations.list() + await tp.destinations.list() + assert slept, "walked into a window the server said was spent" + # Upper bound allows the spread now applied to this path: without it + # every waiter woke at the same absolute instant. + assert 0 < slept[0] <= 7.0 + 0.25 + + async def test_an_unknown_remaining_never_holds(self): + tp, slept = await self._slept({"content-type": "application/json"}) + await tp.destinations.list() + await tp.destinations.list() + assert slept == [] + + async def test_respect_remaining_can_be_turned_off(self): + tp, slept = await self._slept( + {**REST, "RateLimit-Remaining": "0", "RateLimit-Reset": "7"}, + retry=RetryConfig(respect_remaining=False), + ) + await tp.destinations.list() + await tp.destinations.list() + assert slept == [] + + async def test_respect_429_false_never_sleeps_even_on_a_later_call(self): + tp, slept = await self._slept( + {"retry-after": "45"}, retry=RetryConfig(respect_429=False), status=429 + ) + for _ in range(3): + with pytest.raises(RateLimitError): + await tp.destinations.list() + assert slept == [], "opted out of 429 waiting and waited anyway" + + async def test_the_gate_still_holds_for_callers_who_want_it(self): + tp, slept = await self._slept( + {"retry-after": "45"}, retry=RetryConfig(max_retries=0), status=429 + ) + for _ in range(2): + with pytest.raises(RateLimitError): + await tp.destinations.list() + assert slept diff --git a/tests/unit/test_transport.py b/tests/unit/test_transport.py index 9f0d5ac..98005e3 100644 --- a/tests/unit/test_transport.py +++ b/tests/unit/test_transport.py @@ -1,3 +1,6 @@ +import time +from email.utils import formatdate + import httpx import pytest @@ -153,7 +156,10 @@ def test_retry_after_http_date_header_is_parsed(): # None assert _parse_retry_after(None) is None # HTTP-date (in the past -> clamped to 0.0) - assert _parse_retry_after("Wed, 21 Oct 2015 07:28:00 GMT") == 0.0 + # A date in the past is not a wait. It used to come back as 0.0, and only + # None reaches the exponential backoff, so 0.0 meant no wait at all: + # four requests in 3ms against a server that had just said 429. + assert _parse_retry_after("Wed, 21 Oct 2015 07:28:00 GMT") is None # Unparseable assert _parse_retry_after("not-a-date-ever") is None @@ -213,12 +219,22 @@ def test_a_long_wait_is_not_slept_through(self): def test_a_short_wait_is_still_honoured(self): _, slept, calls = self._run("5") - assert slept == [5.0, 5.0, 5.0] + # Jittered: the gate is shared, so without a little spread every + # waiter would wake at the same instant and re-trip the limit + # together. One wait per retry, taken once, never doubled. + assert len(slept) == 3 + assert all(5.0 <= s < 5.3 for s in slept), slept assert len(calls) == 4 - def test_the_cap_is_configurable(self): + def test_the_cap_is_configurable_and_bounds_the_whole_call(self): + # This used to assert three sleeps of 3000s: 9000 seconds of blocking + # under a 3600s cap, because the cap was checked per leg and never + # against the total. The cap is a per-CALL budget now, so the sum is + # what it bounds. _, slept, _ = self._run("3000", RetryConfig(max_retry_after=3600.0)) - assert slept == [3000.0, 3000.0, 3000.0] + assert sum(slept) <= 3600.0 + 0.01, slept + assert slept, "stopped waiting altogether" + assert slept[0] >= 3000.0 def test_no_retry_after_header_still_backs_off(self): # The cap is about the server's stated wait. With no header we fall @@ -226,3 +242,53 @@ def test_no_retry_after_header_still_backs_off(self): _, slept, _ = self._run(None) assert len(slept) == 3 assert all(s > 0 for s in slept) + + +class TestRetryAfterNeverMeansNoWait: + """`None` and `0` are different answers, and conflating them hammers us. + + Only `None` reaches the exponential backoff. A header that parsed to zero + -- `Retry-After: 0`, which RFC 9110 permits, or a negative, or an + already-past date -- therefore produced no wait at all. Measured before + this: four requests in 3ms against a server actively refusing them, and + 204 requests a second across ten threads. + """ + + @pytest.mark.parametrize("raw", ["0", "-5", "0.0", "Wed, 21 Oct 2015 07:28:00 GMT"]) + def test_a_non_positive_wait_is_no_wait_at_all(self, raw): + assert _parse_retry_after(raw) is None + + @pytest.mark.parametrize("raw", ["1", "45", "0.5"]) + def test_a_real_wait_is_honoured(self, raw): + assert _parse_retry_after(raw) == float(raw) + + def test_a_spin_falls_back_to_backoff(self): + # The behaviour that matters: a zero must not skip the backoff. + slept: list[float] = [] + + def handler(request: httpx.Request) -> httpx.Response: + return httpx.Response(429, headers={"retry-after": "0"}, json={}) + + client = httpx.Client(transport=httpx.MockTransport(handler), base_url="https://x/v1") + transport = SyncTransport( + client=client, + base_url="https://x/v1", + user_agent="test/1", + retry=RetryConfig(max_retries=3), + sleep=slept.append, + ) + with pytest.raises(RateLimitError): + transport.get("/anything") + assert len(slept) == 3 + assert all(s > 0 for s in slept), slept + # And it grows, rather than retrying at a fixed rate. + assert slept[-1] > slept[0] + + def test_a_naive_http_date_is_read_as_utc(self): + # parsedate_to_datetime returns a NAIVE datetime for the RFC 5322 + # `-0000` form, and .timestamp() then read it as local time: wrong by + # the host's UTC offset, and negative enough to become the spin above. + naive = _parse_retry_after(formatdate(time.time() + 120)) + gmt = _parse_retry_after(formatdate(time.time() + 120, usegmt=True)) + assert naive is not None and gmt is not None + assert abs(naive - gmt) < 2.0, (naive, gmt) diff --git a/tests/unit/test_transport_async.py b/tests/unit/test_transport_async.py index 98888ba..53b2480 100644 --- a/tests/unit/test_transport_async.py +++ b/tests/unit/test_transport_async.py @@ -153,5 +153,9 @@ async def test_a_short_wait_is_still_honoured(self): transport, slept, calls = self._run_setup("5", RetryConfig()) with pytest.raises(RateLimitError): await transport.get("/anything") - assert slept == [5.0, 5.0, 5.0] + # Jittered: the gate is shared, so without a little spread every + # waiter would wake at the same instant and re-trip the limit + # together. One wait per retry, taken once, never doubled. + assert len(slept) == 3 + assert all(5.0 <= s < 5.3 for s in slept), slept assert len(calls) == 4 diff --git a/themeparks/__init__.py b/themeparks/__init__.py index 48d148c..bf65055 100644 --- a/themeparks/__init__.py +++ b/themeparks/__init__.py @@ -10,6 +10,7 @@ ThemeParksError, TimeoutError, ) +from themeparks._ratelimit import RateLimit, RateLimits from themeparks._transport import RetryConfig __all__ = [ @@ -17,6 +18,8 @@ "AsyncThemeParks", "BudgetExhaustedError", "HistorySpan", + "RateLimit", + "RateLimits", "Cache", "CacheConfig", "InMemoryLRUCache", diff --git a/themeparks/_client.py b/themeparks/_client.py index 4b2a136..1cdb04f 100644 --- a/themeparks/_client.py +++ b/themeparks/_client.py @@ -10,6 +10,7 @@ from themeparks._cache import Cache, CacheConfig, InMemoryLRUCache, ttl_for_path from themeparks._ergonomic.destinations import AsyncDestinationsApi, DestinationsApi from themeparks._ergonomic.entity import AsyncEntityHandle, EntityHandle +from themeparks._ratelimit import RateLimits from themeparks._raw import AsyncRawClient, RawClient from themeparks._transport import AsyncTransport, RetryConfig, SyncTransport @@ -52,6 +53,16 @@ def __init__(self, inner: SyncTransport, cache: Cache) -> None: self._inner = inner self._cache = cache + @property + def rate_limit(self) -> RateLimits: + """Whatever the inner transport last learned. + + A cache HIT sends no request and so learns nothing, which is correct: + the figures then keep saying what the last real response said. They + are not invalidated by a hit, because a hit spent no budget either. + """ + return self._inner.rate_limit + def get(self, path: str) -> Any: ttl = ttl_for_path(path) if ttl > 0: @@ -69,6 +80,16 @@ def __init__(self, inner: AsyncTransport, cache: Cache) -> None: self._inner = inner self._cache = cache + @property + def rate_limit(self) -> RateLimits: + """Whatever the inner transport last learned. + + A cache HIT sends no request and so learns nothing, which is correct: + the figures then keep saying what the last real response said. They + are not invalidated by a hit, because a hit spent no budget either. + """ + return self._inner.rate_limit + async def get(self, path: str) -> Any: ttl = ttl_for_path(path) if ttl > 0: @@ -125,6 +146,24 @@ def __init__( # noqa: PLR0913 ) self.destinations = DestinationsApi(raw=self.raw) + @property + def rate_limit(self) -> RateLimits: + """What the server last said about your two budgets. + + `rate_limit.rest` is the per-minute REST meter; `rate_limit.history` + is the separate hourly history budget. Every field can be None, + because every field can be legitimately absent: an unmetered plan + advertises nothing, and neither does a publicly cacheable response, + since the figures belong to whoever populated the cache. + + None therefore means "the server did not say", never "nothing left". + + with ThemeParks(api_key=KEY) as tp: + tp.entity(park).live() + print(tp.rate_limit.rest.remaining) # e.g. 299 + """ + return self.raw._t.rate_limit + def entity(self, entity_id: str) -> EntityHandle: return self._entity_ctor(entity_id) @@ -184,6 +223,24 @@ def __init__( # noqa: PLR0913 ) self.destinations = AsyncDestinationsApi(raw=self.raw) + @property + def rate_limit(self) -> RateLimits: + """What the server last said about your two budgets. + + `rate_limit.rest` is the per-minute REST meter; `rate_limit.history` + is the separate hourly history budget. Every field can be None, + because every field can be legitimately absent: an unmetered plan + advertises nothing, and neither does a publicly cacheable response, + since the figures belong to whoever populated the cache. + + None therefore means "the server did not say", never "nothing left". + + with ThemeParks(api_key=KEY) as tp: + tp.entity(park).live() + print(tp.rate_limit.rest.remaining) # e.g. 299 + """ + return self.raw._t.rate_limit + def entity(self, entity_id: str) -> AsyncEntityHandle: return self._entity_ctor(entity_id) diff --git a/themeparks/_ratelimit.py b/themeparks/_ratelimit.py new file mode 100644 index 0000000..24725a5 --- /dev/null +++ b/themeparks/_ratelimit.py @@ -0,0 +1,182 @@ +"""What the server says about your budget, and how to stay inside it. + +TWO BUDGETS, SEPARATELY METERED. The API meters requests per minute, and +history requests again per hour. They are different windows over different +counters, so the server advertises them in two sets of headers: + + RateLimit-Limit / -Policy / -Remaining / -Reset per minute + RateLimit-History-Limit / -Policy / -Remaining / -Reset per hour + +Both were being thrown away. The SDK only ever read `Retry-After`, and only +after a 429 had already happened -- so it could tell you that you had run out, +never that you were about to. + +ABSENCE IS NOT ZERO. A response a shared cache may store carries no per-caller +figures at all, because they belong to whoever populated the cache entry. That +is every anonymous response. So `None` here means "the server did not say", +which is a different thing from "nothing left", and nothing in this module may +confuse the two. +""" + +from __future__ import annotations + +import random +import threading +import time +from collections.abc import Mapping +from dataclasses import dataclass, field + +#: Header prefixes for the two meters. +_REST_PREFIX = "ratelimit" +_HISTORY_PREFIX = "ratelimit-history" + + +def _int_or_none(raw: str | None) -> int | None: + """These fields are integers per the draft spec; anything else is unknown. + + `int()` alone was too generous and differed from the JavaScript sibling on + the same input: PEP 515 means it reads "1_0" as 10. The server only ever + sends a non-negative integer, so it does not fire today, but a pair of + libraries whose selling point is parity should not disagree on it. + """ + if raw is None: + return None + trimmed = raw.strip() + if not trimmed.isdigit(): + return None + return int(trimmed) + + +@dataclass(frozen=True) +class RateLimit: + """One meter's state, as of the last response that mentioned it. + + Every field is optional because every field can be legitimately absent: + an unmetered plan advertises nothing, and neither does a publicly + cacheable response. + """ + + #: Requests allowed per window, or None if the server did not say. + limit: int | None = None + #: Requests left in the current window, or None if the server did not say. + remaining: int | None = None + #: Seconds until the window resets, as of `observed_at`. + reset: int | None = None + #: The raw policy string, e.g. "300;w=60". + policy: str | None = None + #: `time.monotonic()` when this was read, so `reset` can be aged. + observed_at: float | None = None + + @property + def exhausted(self) -> bool: + """True only when the server SAID there is nothing left. + + An unknown remaining is not exhaustion. Treating it as such would make + an anonymous caller, whose responses never carry figures, wait forever. + """ + return self.remaining == 0 + + def seconds_until_reset(self, now: float | None = None) -> float | None: + """How long is left of the window, counting down from when we read it. + + `reset` is a relative value frozen at `observed_at`; using it later + without ageing it is how a client waits far longer than it needs to. + """ + if self.reset is None or self.observed_at is None: + return None + elapsed = (time.monotonic() if now is None else now) - self.observed_at + return max(0.0, self.reset - elapsed) + + +@dataclass(frozen=True) +class RateLimits: + """Both meters. Reached as `client.rate_limit`.""" + + rest: RateLimit = field(default_factory=RateLimit) + history: RateLimit = field(default_factory=RateLimit) + + +def _read_one(headers: Mapping[str, str], prefix: str, now: float) -> RateLimit: + limit = _int_or_none(headers.get(f"{prefix}-limit")) + remaining = _int_or_none(headers.get(f"{prefix}-remaining")) + reset = _int_or_none(headers.get(f"{prefix}-reset")) + policy = headers.get(f"{prefix}-policy") + if limit is None and remaining is None and reset is None and policy is None: + return RateLimit() + return RateLimit(limit=limit, remaining=remaining, reset=reset, policy=policy, observed_at=now) + + +def read_rate_limits(headers: Mapping[str, str], previous: RateLimits) -> RateLimits: + """Merge whatever this response said into what we already knew. + + A response that mentions neither meter leaves both alone, and so does a + response served from a cache. That matters + because most responses mention only one: the history headers appear on + history routes, and on a cacheable response neither appears. Overwriting + with blanks would mean the last cacheable response erased everything the + SDK had learned. + + Header lookup is case-insensitive via httpx's own mapping, but a plain dict + is accepted too, so the keys are compared lowercased. + """ + lowered = {k.lower(): v for k, v in headers.items()} + # A cache HIT carries the figures of whoever populated the entry, frozen + # at that moment. They are not ours and the countdown is already wrong, so + # the honest reading is that this response said nothing. + age = _int_or_none(lowered.get("age")) + if age is not None and age > 0: + return previous + # No prefix filtering needed: _read_one looks up EXACT keys, so + # "ratelimit-limit" and "ratelimit-history-limit" cannot collide. An + # earlier version filtered the history keys out before reading the REST + # meter; removing that filter changed no behaviour and no test, which is + # what a redundant guard looks like. The bleed it guarded against is + # covered by a test either way. + now = time.monotonic() + history = _read_one(lowered, _HISTORY_PREFIX, now) + rest = _read_one(lowered, _REST_PREFIX, now) + return RateLimits( + rest=rest if rest.observed_at is not None else previous.rest, + history=history if history.observed_at is not None else previous.history, + ) + + +class Gate: + """One shared "not before" instant for a whole client. + + WHY SHARED. A 429 applies to the CALLER, not to the request that happened + to meet it. With a per-request backoff, ten concurrent requests each sleep + their own Retry-After and then all retry at the same instant, re-tripping + the limit together -- a thundering herd the client inflicts on itself, and + on us. One gate means the wait is taken once. + + Each waiter adds its own small jitter on the way out, because waking + together is the other half of the same problem. + """ + + def __init__(self, jitter: float = 0.25) -> None: + self._until = 0.0 + self._jitter = jitter + self._lock = threading.Lock() + + @property + def deadline(self) -> float: + """When the gate opens, on the monotonic clock. 0 if it is open.""" + with self._lock: + return self._until + + def close_for(self, seconds: float) -> None: + """Hold every request on this client for at least `seconds`.""" + deadline = time.monotonic() + max(0.0, seconds) + with self._lock: + # Never bring the gate forward: a shorter Retry-After arriving + # while a longer one is in force would release the herd early. + self._until = max(self._until, deadline) + + def wait_seconds(self) -> float: + """How long this caller should hold off, jitter included. 0 if open.""" + with self._lock: + remaining = self._until - time.monotonic() + if remaining <= 0: + return 0.0 + return remaining + random.random() * self._jitter diff --git a/themeparks/_transport.py b/themeparks/_transport.py index 4023e50..7a8f5f9 100644 --- a/themeparks/_transport.py +++ b/themeparks/_transport.py @@ -8,16 +8,22 @@ import time from collections.abc import Awaitable, Callable from dataclasses import dataclass +from datetime import timezone from typing import Any import httpx from themeparks._errors import APIError, NetworkError, RateLimitError, TimeoutError +from themeparks._ratelimit import Gate, RateLimits, read_rate_limits _STATUS_TOO_MANY_REQUESTS = 429 _STATUS_SERVER_ERROR = 500 _ERROR_BODY_EXCERPT_LIMIT = 200 _ERROR_MESSAGE_LIMIT = 300 +#: Below this, a computed wait is floating-point residue rather than a wait. +_MIN_SLEEP_SECONDS = 0.001 +#: Spread applied to a synchronised release, so waiters do not wake as one. +_SPREAD_SECONDS = 0.25 @dataclass @@ -33,20 +39,57 @@ class RetryConfig: #: not sleep at all, and raise `RateLimitError` carrying `retry_after` so #: the caller can checkpoint and come back. max_retry_after: float = 120.0 + #: Wait out a window the server has already told us is spent. + #: + #: When a response says `remaining: 0`, the next request is a guaranteed + #: 429 that also costs us a unit of the caller's budget to refuse. Waiting + #: for the reset it advertised is strictly better than sending it. Off + #: turns the client back into a purely reactive one. + respect_remaining: bool = True def _parse_retry_after(raw: str | None) -> float | None: + """Seconds to wait, or None when the header gives us nothing usable. + + NONE AND ZERO ARE DIFFERENT ANSWERS, and conflating them turned the client + into a hammer. Only `None` reaches the exponential backoff, so a header + that parsed to 0 -- which `Retry-After: 0` is, legally, per RFC 9110, and + which a negative or already-past date also produces -- meant no wait at + all. Measured: four requests in 3ms against a server that had just said + 429, where an absent header correctly took 1962ms. Ten threads made that + 204 requests a second at a server actively refusing them. + + So a non-positive wait is not a wait, and we say None. + """ if raw is None: return None + seconds: float | None try: - return max(0.0, float(raw)) + seconds = float(raw) except ValueError: - pass + seconds = _parse_http_date_delta(raw) + if seconds is None or seconds <= 0: + return None + return seconds + + +def _parse_http_date_delta(raw: str) -> float | None: + """An HTTP-date Retry-After, as seconds from now. + + RFC 9110 requires the IMF-fixdate (GMT) form, but RFC 5322 `-0000` and a + bare date both appear in the wild, and `parsedate_to_datetime` returns a + NAIVE datetime for them. `.timestamp()` then reads it as local time, so on + a host an hour off UTC the answer was wrong by exactly that hour -- and + because it came out negative it became 0, landing in the no-backoff spin + above. Assume UTC when the sender did not say. + """ try: parsed = email.utils.parsedate_to_datetime(raw) - return max(0.0, parsed.timestamp() - time.time()) except Exception: return None + if parsed.tzinfo is None: + parsed = parsed.replace(tzinfo=timezone.utc) + return parsed.timestamp() - time.time() def _wait_too_long(retry_after: float | None, retry: RetryConfig) -> bool: @@ -125,11 +168,75 @@ def __init__( # noqa: PLR0913 self._retry = retry self._headers = _headers(user_agent, api_key) self._sleep = sleep + self.rate_limit = RateLimits() + self._gate = Gate() + + def _hold(self, budget: float) -> float: + """Wait before sending, if we already know this request would fail. + + Returns how long it slept, so the caller can keep a running total. The + TOTAL is what `max_retry_after` bounds, not each leg: the gate wait and + the spent-window wait are both self-initiated holds, and they stack. + Measured before this budget existed: a 429 carrying `Retry-After: 5` + and `RateLimit-Reset: 60` slept 5 then 55, three times over -- 180 + seconds inside one call whose cap was 120. Each leg was under the cap, + so the per-leg check never fired, and the promise the cap makes was + reachable around. + + Two reasons to hold, and they are different. The GATE is a 429 the + server has already issued to this caller: the wait belongs to them, + not to whichever request met it, so it is shared and taken once. The + REMAINING check is a window the server told us is spent -- sending + into it is a guaranteed 429 that also costs a unit of budget to + refuse, so waiting for the advertised reset is strictly better. + + A remaining we were never told is not a spent one. Anonymous + responses carry no figures at all, so an unknown must never hold. + """ + spent = 0.0 + # Re-read the gate after waiting. It slept once and returned, so a + # waiter that woke while someone else's 429 had pushed the gate + # further out sent anyway. Only loop when the deadline actually + # MOVED: re-reading unconditionally spins against any clock that does + # not advance. + while True: + before = self._gate.deadline + wait = min(self._gate.wait_seconds(), budget - spent) + if wait <= _MIN_SLEEP_SECONDS: + break + self._sleep(wait) + spent += wait + if self._gate.deadline <= before: + break + if not self._retry.respect_remaining: + return spent + for meter in (self.rate_limit.rest, self.rate_limit.history): + if not meter.exhausted: + continue + left = meter.seconds_until_reset() + if left is None or left <= 0 or left > self._retry.max_retry_after: + # Past the cap we do not sit on it: the caller gets the 429 + # and its Retry-After, and can decide. Same rule the retry + # path follows. + continue + # Jittered like the gate. Without it every waiter derived `left` + # from the same observed_at and woke at the same absolute + # instant -- the tightest burst in the client, on the very branch + # that exists to avoid a 429. + left = min(left + random.random() * _SPREAD_SECONDS, budget - spent) + if left <= _MIN_SLEEP_SECONDS: + break + self._sleep(left) + spent += left + return spent def get(self, path: str) -> Any: url = self._base_url + path attempt = 0 + # One budget for the whole call, because that is what the cap promises. + budget = self._retry.max_retry_after while True: + budget -= self._hold(budget) try: response = self._client.get( path, @@ -144,6 +251,8 @@ def get(self, path: str) -> Any: continue raise NetworkError(f"network error calling {url}") from exc + self.rate_limit = read_rate_limits(response.headers, self.rate_limit) + if response.is_success: return _parse_body(response) @@ -151,13 +260,36 @@ def get(self, path: str) -> Any: status = response.status_code retry_after = _parse_retry_after(response.headers.get("retry-after")) + if ( + status == _STATUS_TOO_MANY_REQUESTS + and self._retry.respect_429 + and retry_after is not None + and not _wait_too_long(retry_after, self._retry) + ): + # The wait belongs to the CALLER, not to whichever request met + # it, so it goes on the shared gate and _hold() serves it once. + # + # respect_429=False means "do not wait on a 429", so it must + # gate the gate too. Without this the caller got the exception + # they asked for and then their NEXT call silently blocked, + # which is an opt-out that does not opt out. + # + # Past the cap the gate is left OPEN on purpose: we raise + # instead, and blocking the caller's next call for most of an + # hour is the opposite of letting them checkpoint and resume. + self._gate.close_for(retry_after) if ( status == _STATUS_TOO_MANY_REQUESTS and self._retry.respect_429 and attempt < self._retry.max_retries and not _wait_too_long(retry_after, self._retry) ): - self._sleep(retry_after if retry_after is not None else _backoff(attempt)) + # No sleep here: the gate above holds the wait and _hold() + # at the top of the loop serves it once. Paying it here too + # would double every backoff, and ten concurrent requests + # would each pay their own and then retry in unison. + if retry_after is None: + self._sleep(_backoff(attempt)) attempt += 1 continue if status == _STATUS_TOO_MANY_REQUESTS: @@ -197,11 +329,51 @@ def __init__( # noqa: PLR0913 self._retry = retry self._headers = _headers(user_agent, api_key) self._sleep: Callable[..., Awaitable[None]] = sleep if sleep is not None else asyncio.sleep + self.rate_limit = RateLimits() + self._gate = Gate() + + async def _hold(self, budget: float) -> float: + """Asynchronous mirror of :meth:`SyncTransport._hold`.""" + spent = 0.0 + # Re-read the gate after waiting. It slept once and returned, so a + # waiter that woke while someone else's 429 had pushed the gate + # further out sent anyway. Only loop when the deadline actually + # MOVED: re-reading unconditionally spins against any clock that does + # not advance. + while True: + before = self._gate.deadline + wait = min(self._gate.wait_seconds(), budget - spent) + if wait <= _MIN_SLEEP_SECONDS: + break + await self._sleep(wait) + spent += wait + if self._gate.deadline <= before: + break + if not self._retry.respect_remaining: + return spent + for meter in (self.rate_limit.rest, self.rate_limit.history): + if not meter.exhausted: + continue + left = meter.seconds_until_reset() + if left is None or left <= 0 or left > self._retry.max_retry_after: + continue + # Jittered like the gate. Without it every waiter derived `left` + # from the same observed_at and woke at the same absolute + # instant -- the tightest burst in the client, on the very branch + # that exists to avoid a 429. + left = min(left + random.random() * _SPREAD_SECONDS, budget - spent) + if left <= _MIN_SLEEP_SECONDS: + break + await self._sleep(left) + spent += left + return spent async def get(self, path: str) -> Any: url = self._base_url + path attempt = 0 + budget = self._retry.max_retry_after while True: + budget -= await self._hold(budget) try: response = await self._client.get( path, @@ -216,6 +388,8 @@ async def get(self, path: str) -> Any: continue raise NetworkError(f"network error calling {url}") from exc + self.rate_limit = read_rate_limits(response.headers, self.rate_limit) + if response.is_success: return _parse_body(response) @@ -223,13 +397,36 @@ async def get(self, path: str) -> Any: status = response.status_code retry_after = _parse_retry_after(response.headers.get("retry-after")) + if ( + status == _STATUS_TOO_MANY_REQUESTS + and self._retry.respect_429 + and retry_after is not None + and not _wait_too_long(retry_after, self._retry) + ): + # The wait belongs to the CALLER, not to whichever request met + # it, so it goes on the shared gate and _hold() serves it once. + # + # respect_429=False means "do not wait on a 429", so it must + # gate the gate too. Without this the caller got the exception + # they asked for and then their NEXT call silently blocked, + # which is an opt-out that does not opt out. + # + # Past the cap the gate is left OPEN on purpose: we raise + # instead, and blocking the caller's next call for most of an + # hour is the opposite of letting them checkpoint and resume. + self._gate.close_for(retry_after) if ( status == _STATUS_TOO_MANY_REQUESTS and self._retry.respect_429 and attempt < self._retry.max_retries and not _wait_too_long(retry_after, self._retry) ): - await self._sleep(retry_after if retry_after is not None else _backoff(attempt)) + # No sleep here: the gate above holds the wait and _hold() + # at the top of the loop serves it once. Paying it here too + # would double every backoff, and ten concurrent requests + # would each pay their own and then retry in unison. + if retry_after is None: + await self._sleep(_backoff(attempt)) attempt += 1 continue if status == _STATUS_TOO_MANY_REQUESTS: diff --git a/uv.lock b/uv.lock index 2fd649b..386cada 100644 --- a/uv.lock +++ b/uv.lock @@ -2114,7 +2114,7 @@ wheels = [ [[package]] name = "themeparks" -version = "3.1.0" +version = "3.2.0" source = { editable = "." } dependencies = [ { name = "eval-type-backport", marker = "python_full_version < '3.10'" },