Skip to content

Commit 6d172b1

Browse files
Address review: drop the fragments and the extra timeout branch
Nothing rate-limit has been released, so 12.feature.rst already covers the base class and the two #25 fragments go. The exhausted-budget case stops being handled separately. One message again, printing the real remaining rather than clamping a negative one to 0.000s, so an acquire() that overran still shows up in it. _BudgetEater goes too, its scenario folded into a parametrised case on _ElapsedAcquire that pins the negative budget, so re-adding the clamp still fails a test. asyncio.ensure_future becomes asyncio.create_task in the two places called out; the two that predate this branch are left alone. parametrize takes a tuple. The _fake_request docstring goes back to what it said, since 3.15 will date the replacement. The prose around the Redis sketch and release() was longer than what it had to say. Per-domain keying keeps the fact the deferred fix turns on -- that clone() takes no arguments, so it cannot learn the host.
1 parent f91ad76 commit 6d172b1

5 files changed

Lines changed: 64 additions & 114 deletions

File tree

‎CHANGES/25.breaking.rst‎

Lines changed: 0 additions & 7 deletions
This file was deleted.

‎CHANGES/25.bugfix.rst‎

Lines changed: 0 additions & 6 deletions
This file was deleted.

‎aiohttp_client_middlewares/rate_limit.py‎

Lines changed: 19 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -26,10 +26,9 @@ class RateLimiter(ABC):
2626
logic lives in :meth:`wait`, shared by every algorithm, so reserving a
2727
slot may perform I/O of its own -- against Redis or a database, say.
2828
29-
An :meth:`acquire` implementation must be cancellation-safe. If it is
30-
cancelled or raises before returning, it is responsible for ensuring
31-
that no reservation is left behind. Once it returns, :meth:`wait` owns
32-
the reservation and calls :meth:`release` if the slot cannot be used.
29+
Until :meth:`acquire` returns, cleaning up a half-made reservation is
30+
its own responsibility; once it returns, :meth:`wait` owns the slot and
31+
calls :meth:`release` if it cannot be used.
3332
3433
An async method that contains no suspension point still runs atomically
3534
when awaited directly. :class:`TokenBucket` relies on that property to
@@ -40,11 +39,9 @@ class RateLimiter(ABC):
4039
async def acquire(self) -> float:
4140
"""Reserve a slot and return the delay to sleep before sending.
4241
43-
The delay must be a non-negative, finite number of seconds.
44-
:meth:`wait` takes it on trust: a NaN compares false against both
45-
the budget and zero, so the request would go out unthrottled.
46-
47-
If cancellation or another exception prevents this method from
42+
The delay must be non-negative, finite seconds; :meth:`wait` takes that
43+
on trust, and a NaN would send the request through unthrottled. If
44+
cancellation or another exception prevents this method from
4845
returning, it must not leave a reservation behind.
4946
"""
5047

@@ -64,40 +61,31 @@ def release(self) -> None:
6461
cancelled while sleeping. The default is a no-op for algorithms
6562
that have nothing to return.
6663
67-
Runs from an ``except asyncio.CancelledError`` block, so it must
68-
not await: a second cancellation, or the loop shutting down, would
69-
truncate it part-way and lose the slot for good. A limiter that has
70-
to reach its backend to hand a slot back can schedule that round
71-
trip as a task from here.
72-
73-
It must not raise, either. :meth:`wait` calls it while unwinding, so
74-
an exception here would replace the :exc:`asyncio.TimeoutError` the
75-
caller is owed -- or the :exc:`asyncio.CancelledError`, leaving a
76-
cancelled request reporting an ordinary failure.
64+
Must neither await nor raise, since one of those calls is from an
65+
``except asyncio.CancelledError`` block: awaiting there can be
66+
truncated part-way, and raising would replace the exception the
67+
caller is owed. A limiter that has to reach its backend to hand a
68+
slot back can schedule that round trip as a task from here.
7769
"""
7870

7971
async def wait(self, timeout: float | None = None) -> None:
8072
"""Reserve a slot and wait until the request may be sent.
8173
82-
Time spent in :meth:`acquire` is charged against *timeout* once it
83-
returns, but is not bounded by it, so an implementation that can
84-
hang needs a deadline of its own. When the delay left to serve
85-
exceeds what is left of the budget, the slot is handed back and
86-
:exc:`asyncio.TimeoutError` is raised without sleeping.
74+
Time in :meth:`acquire` is charged against *timeout* once it
75+
returns, though not bounded by it, so an implementation that can
76+
hang needs its own deadline. When the delay exceeds what is left,
77+
the slot is handed back and :exc:`asyncio.TimeoutError` raised
78+
without sleeping.
8779
"""
8880
started = time.monotonic()
8981
delay = await self.acquire()
9082

9183
if timeout is not None:
92-
elapsed = time.monotonic() - started
93-
remaining = timeout - elapsed
84+
# Goes negative when acquiring alone outlasted the timeout; the
85+
# message reports it as such rather than clamping it to zero.
86+
remaining = timeout - (time.monotonic() - started)
9487
if delay > remaining:
9588
self.release()
96-
if remaining <= 0.0:
97-
raise asyncio.TimeoutError(
98-
f"reserving a rate-limit slot took {elapsed:.3f}s, "
99-
f"exhausting the {timeout:.3f}s timeout"
100-
)
10189
raise asyncio.TimeoutError(
10290
f"rate limiter would delay the request {delay:.3f}s, "
10391
f"beyond the {remaining:.3f}s remaining timeout"

‎docs/api.rst‎

Lines changed: 14 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -103,29 +103,22 @@ Rate limiting
103103
def clone(self):
104104
return RedisLimiter(self._redis, self._key, self._script)
105105

106-
Three shortcuts in that sketch matter. Its ``clone()`` reuses one key, so
107-
under ``per_domain=True`` every host would draw on a single shared limit --
108-
and since ``clone()`` is called as a no-argument factory it never learns the
109-
host, so a shared backend cannot key itself per host that way. Give each host
110-
its own limiter, built with its own key, rather than using ``per_domain=True``.
111-
It leaves ``release()`` at the default no-op, so a slot reserved just before
112-
a timeout bail is not handed back. And a cancellation between the script
113-
running and its reply arriving leaves a reservation nobody holds. The last
114-
two are what an expiry on the reservation is for: give every slot one, and
115-
an abandoned slot lapses on its own.
106+
The sketch keeps one key across clones, and ``clone()`` takes no arguments
107+
so it cannot learn the host; give each host its own limiter rather than
108+
using ``per_domain=True``. It also leaves ``release()`` at the default
109+
no-op, and a cancelled round trip can leave a reservation nobody holds;
110+
an expiry on each reservation covers both.
116111

117112
``wait(timeout=None)`` is supplied by the base class. It charges async
118-
acquisition against *timeout* once ``acquire()`` returns -- it does not
119-
bound the call itself, so an implementation that can hang needs a deadline
120-
of its own -- then fails fast when the delay left to serve exceeds what is
121-
left of the budget, sleeps otherwise, and calls ``release()`` if a
122-
successfully acquired slot cannot be used. ``release()`` is synchronous
123-
and defaults to a no-op for algorithms that have nothing to return. It stays
124-
synchronous because ``wait()`` also calls it from a cancellation handler,
125-
where an awaiting implementation can be truncated part-way and lose the slot
126-
for good; a limiter that has to reach its backend to hand one back can
127-
schedule that round trip as a task. It must not raise either, since
128-
``wait()`` calls it while unwinding.
113+
acquisition against *timeout* once ``acquire()`` returns -- without bounding
114+
the call itself, so an implementation that can hang needs its own deadline --
115+
then fails fast when the delay exceeds what is left, sleeps otherwise, and
116+
calls ``release()`` if an acquired slot cannot be used. ``release()`` defaults
117+
to a no-op for algorithms with nothing to return, and stays synchronous: one
118+
of those calls is from a cancellation handler, where awaiting can be truncated
119+
part-way and lose the slot for good, and raising would replace the exception
120+
the caller is owed. A limiter that has to reach its backend to hand a slot
121+
back can schedule that round trip as a task.
129122

130123
``acquire()`` must be cancellation-safe: if cancellation or another
131124
exception prevents it from returning, it must leave no reservation behind.

‎tests/test_rate_limit.py‎

Lines changed: 31 additions & 49 deletions
Original file line numberDiff line numberDiff line change
@@ -35,9 +35,8 @@ def _fake_request(
3535
"""A stand-in ``ClientRequest`` exposing just ``.url`` and ``.timeout``.
3636
3737
The timeout is set here, while the mock is still untyped, rather than on
38-
the returned value: the supported aiohttp releases have no
39-
``ClientRequest.timeout`` at all, and 3.15 adds it read-only, so assigning
40-
it through the ``ClientRequest`` annotation type-checks on neither.
38+
the returned value: ``ClientRequest.timeout`` is a read-only property, so
39+
assigning it through the ``ClientRequest`` annotation does not type-check.
4140
"""
4241
req = mock.create_autospec(ClientRequest, instance=True)
4342
req.url = URL(f"http://{host}")
@@ -428,28 +427,45 @@ async def handler(req: ClientRequest) -> ClientResponse:
428427

429428

430429
class _ElapsedAcquire(RateLimiter):
431-
"""Limiter that advances a fake clock while reserving a slot."""
430+
"""Limiter that spends fake time reserving a slot, then asks for a delay."""
432431

433-
def __init__(self, clock: _FakeClock) -> None:
432+
def __init__(self, clock: _FakeClock, spend: float, delay: float) -> None:
434433
self._clock = clock
434+
self._spend = spend
435+
self._delay = delay
435436
self.releases = 0
436437

437438
async def acquire(self) -> float:
438-
self._clock.advance(0.075)
439-
return 0.05
439+
self._clock.advance(self._spend)
440+
return self._delay
440441

441442
def release(self) -> None:
442443
self.releases += 1
443444

444445
def clone(self) -> "_ElapsedAcquire":
445-
return _ElapsedAcquire(self._clock)
446+
return _ElapsedAcquire(self._clock, self._spend, self._delay)
446447

447448

448-
async def test_acquisition_time_reduces_remaining_timeout(clock: _FakeClock) -> None:
449-
"""The delay is checked against the budget left after acquisition."""
450-
limiter = _ElapsedAcquire(clock)
449+
@pytest.mark.parametrize(
450+
("spend", "delay", "budget_left"),
451+
(
452+
(0.075, 0.05, r"0\.025s"), # some budget survived acquisition
453+
(0.5, 0.0, r"-0\.400s"), # acquiring alone outlasted the whole timeout
454+
),
455+
)
456+
async def test_acquisition_time_reduces_remaining_timeout(
457+
clock: _FakeClock, spend: float, delay: float, budget_left: str
458+
) -> None:
459+
"""The delay is checked against the budget left after acquisition.
451460
452-
with pytest.raises(asyncio.TimeoutError, match="remaining timeout"):
461+
The overrun case reports a negative budget rather than clamping it: a
462+
0.000s delay said to exceed a 0.000s timeout explains nothing.
463+
"""
464+
limiter = _ElapsedAcquire(clock, spend, delay)
465+
466+
with pytest.raises(
467+
asyncio.TimeoutError, match=f"beyond the {budget_left} remaining"
468+
):
453469
await limiter.wait(timeout=0.1)
454470

455471
assert limiter.releases == 1
@@ -480,7 +496,7 @@ def clone(self) -> "_CancelledAcquire":
480496
async def test_acquisition_cancellation_does_not_call_release() -> None:
481497
"""Cleanup before acquire returns belongs to the acquire implementation."""
482498
limiter = _CancelledAcquire()
483-
task = asyncio.ensure_future(limiter.wait(timeout=60.0))
499+
task = asyncio.create_task(limiter.wait(timeout=60.0))
484500
await asyncio.sleep(0)
485501

486502
await _cancel_and_join(task)
@@ -489,40 +505,6 @@ async def test_acquisition_cancellation_does_not_call_release() -> None:
489505
assert limiter.releases == 0
490506

491507

492-
class _BudgetEater(RateLimiter):
493-
"""Limiter whose reservation alone outlasts the caller's whole timeout."""
494-
495-
def __init__(self, clock: _FakeClock) -> None:
496-
self._clock = clock
497-
self.releases = 0
498-
499-
async def acquire(self) -> float:
500-
self._clock.advance(0.5)
501-
return 0.0 # no delay left to serve -- the budget went on acquiring
502-
503-
def release(self) -> None:
504-
self.releases += 1
505-
506-
def clone(self) -> "_BudgetEater":
507-
return _BudgetEater(self._clock)
508-
509-
510-
async def test_acquisition_alone_can_exhaust_the_timeout(clock: _FakeClock) -> None:
511-
"""A zero delay still fails once acquiring has spent the whole budget.
512-
513-
The message has to name acquisition as the cause: reporting a 0.000s delay
514-
as exceeding a 0.000s budget would be self-contradictory.
515-
"""
516-
limiter = _BudgetEater(clock)
517-
518-
with pytest.raises(
519-
asyncio.TimeoutError, match=r"took 0\.500s, exhausting the 0\.100s"
520-
):
521-
await limiter.wait(timeout=0.1)
522-
523-
assert limiter.releases == 1
524-
525-
526508
async def test_non_positive_total_timeout_means_no_deadline() -> None:
527509
"""aiohttp arms no deadline for total <= 0, so the limiter must not either.
528510
@@ -566,7 +548,7 @@ async def test_token_bucket_acquire_never_yields_to_the_loop(clock: _FakeClock)
566548
async def competitor() -> None:
567549
ran.append("competitor")
568550

569-
task = asyncio.ensure_future(competitor())
551+
task = asyncio.create_task(competitor())
570552
await bucket.acquire()
571553
assert ran == [], "acquire() yielded, so callers are no longer ordered"
572554

@@ -594,7 +576,7 @@ def clone(self) -> "_RecordsReleases":
594576
return _RecordsReleases(self._delay, self._error)
595577

596578

597-
@pytest.mark.parametrize("delay", [0.0, 0.001])
579+
@pytest.mark.parametrize("delay", (0.0, 0.001))
598580
async def test_wait_keeps_a_slot_it_actually_used(delay: float) -> None:
599581
"""A slot the caller goes on to use must never be handed back.
600582

0 commit comments

Comments
 (0)