Keep stream deadlines from cancelling the caller between messages - #380
Conversation
|
(In draft because I still need to understand this a little more fully myself; feel free to review I suppose...) |
9b41377 to
aa5c73c
Compare
| ) | ||
| else: | ||
| timeout_s = None | ||
| deadline = _StreamDeadline(None) |
There was a problem hiding this comment.
I think the state machine is to eagerly close a stream that isn't being read with a separate async task - it feels too heavyweight to me and I think there is a better alternative on the server side.
I think this simpler form that makes sure timeout is only around the await in our code here and not around the response generator should work more simply enough for this PR
deadline = loop.time() + timeout_ms / 1000.0 # or None
async with contextlib.AsyncExitStack() as stack:
async with asyncio_timeout_at(deadline):
resp = await stack.enter_async_context(self._http_client.stream(...))
...
it = aiter(resp.content)
while True:
async with asyncio_timeout_at(deadline):
try:
chunk = await anext(it)
except StopAsyncIteration:
break
for message in reader.feed(bytes(chunk)):
yield message
await sleep(0)There was a problem hiding this comment.
Thanks, that's much simpler. I've dropped the close-at-deadline commit and gone with your version, plus one addition: a deadline check before each read. When a chunk is already buffered, anext(it) returns without suspending, so asyncio_timeout_at with a deadline that has already passed never fires and the stream runs to completion after the deadline. The regression test caught this with your snippet as-is (DID NOT RAISE ConnectError). With if deadline is not None and loop.time() >= deadline: raise TimeoutError at the top of the loop, it passes on 3.10 and 3.14. The error-response path never yields, so it keeps a single asyncio_timeout_at around its whole read.
On the server side: our server parses connect-timeout-ms / grpc-timeout into ctx.timeout_ms but doesn't end the call when it runs out. Is that what you had in mind? If so, I'll open an issue for it as a follow-up.
There was a problem hiding this comment.
Yeah - to be honest I had thought about that a while ago but forgot about it, so the issue would have helped I guess 😅 The weird things that needed more thought
- How could conformance pass in the first place without a server timeout
- Relatively hard/annoying for WSGI, especially on Windows
But probably time to sort this out
The async client entered `asyncio.timeout` around the whole streaming response, including each `yield`. When the deadline passed while the caller was handling a message, asyncio cancelled the caller's task: a bare `CancelledError` surfaced in the caller's own `await`, the task was left with `cancelling() == 1`, and the `DEADLINE_EXCEEDED` error was raised later in an unobserved cleanup task. The deadline now applies to opening the stream and to each read of the response body, and a read after the deadline raises `DEADLINE_EXCEEDED` from the iterator. Each read checks the deadline first, since reading a chunk that is already buffered doesn't suspend and so wouldn't trigger the timeout. The Python 3.10 timeout backport gains `timeout_at` for this. This is the pattern flake8-async's `ASYNC101` describes; ruff doesn't implement it, and its preview `ASYNC119` also flags harmless context managers, so the regression test guards it instead. Signed-off-by: Stefan VanBuren <stefan@vanburen.xyz>
aa5c73c to
9997498
Compare
With the async client, a server-streaming or bidi call with
timeout_mscould cancel the caller's own task. If the deadline passed while the caller was doing its own work between messages (a slowasync forbody, or holding the iterator while doing something else), the caller got a bareCancelledErrorin unrelated code instead ofConnectError(DEADLINE_EXCEEDED), and the task was left with a pending cancellation that can confuse an enclosingTaskGrouporasyncio.timeout.The deadline is now enforced only while the client is waiting on the network, so it surfaces as
DEADLINE_EXCEEDEDfrom the iterator. The new test fails onmainon 3.10 and 3.14. I looked at enabling ruff's previewASYNC119as a guard, but it also flags harmless code (_request_streamon the server holds only a validating reader across itsyield), andASYNC101, which targets this case exactly, isn't implemented in ruff, so the test is the guard.