Skip to content

Commit e56211f

Browse files
committed
gh-157025: Fix transport cleanup when cancelling asyncio native sendfile
Ensure cancellation while native sendfile waits for the write buffer to drain resets the empty waiter, restores the prior reading state, and restores the selector transport registry entry.
1 parent 57594aa commit e56211f

4 files changed

Lines changed: 47 additions & 2 deletions

File tree

‎Lib/asyncio/proactor_events.py‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -766,8 +766,9 @@ async def _sock_sendfile_native(self, sock, file, offset, count):
766766
async def _sendfile_native(self, transp, file, offset, count):
767767
resume_reading = transp.is_reading()
768768
transp.pause_reading()
769-
await transp._make_empty_waiter()
769+
empty_waiter = transp._make_empty_waiter()
770770
try:
771+
await empty_waiter
771772
return await self.sock_sendfile(transp._sock, file, offset, count,
772773
fallback=False)
773774
finally:

‎Lib/asyncio/selector_events.py‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -739,8 +739,9 @@ async def _sendfile_native(self, transp, file, offset, count):
739739
del self._transports[transp._sock_fd]
740740
resume_reading = transp.is_reading()
741741
transp.pause_reading()
742-
await transp._make_empty_waiter()
742+
empty_waiter = transp._make_empty_waiter()
743743
try:
744+
await empty_waiter
744745
return await self.sock_sendfile(transp._sock, file, offset, count,
745746
fallback=False)
746747
finally:

‎Lib/test/test_asyncio/test_sendfile.py‎

Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -376,6 +376,47 @@ def test_sendfile(self):
376376
self.assertEqual(srv_proto.data, self.DATA)
377377
self.assertEqual(self.file.tell(), len(self.DATA))
378378

379+
def test_sendfile_cancel_empty_waiter(self):
380+
for reading in (True, False):
381+
with self.subTest(reading=reading):
382+
srv_proto, cli_proto = self.prepare_sendfile()
383+
transport = cli_proto.transport
384+
if not reading:
385+
transport.pause_reading()
386+
waiter = self.loop.create_future()
387+
388+
def make_empty_waiter():
389+
transport._empty_waiter = waiter
390+
return waiter
391+
392+
with mock.patch.object(transport, '_make_empty_waiter',
393+
side_effect=make_empty_waiter):
394+
task = self.loop.create_task(
395+
self.loop.sendfile(transport, self.file))
396+
test_utils.run_briefly(self.loop)
397+
self.assertIs(transport._empty_waiter, waiter)
398+
self.assertFalse(waiter.done())
399+
self.assertFalse(transport.is_reading())
400+
task.cancel()
401+
with self.assertRaises(asyncio.CancelledError):
402+
self.run_loop(task)
403+
404+
try:
405+
self.assertIsNone(transport._empty_waiter)
406+
self.assertEqual(transport.is_reading(), reading)
407+
if isinstance(self.loop, asyncio.SelectorEventLoop):
408+
self.assertIs(
409+
self.loop._transports[transport._sock_fd],
410+
transport)
411+
finally:
412+
transport._reset_empty_waiter()
413+
414+
ret = self.run_loop(self.loop.sendfile(transport, self.file))
415+
transport.close()
416+
self.run_loop(srv_proto.done)
417+
self.assertEqual(ret, len(self.DATA))
418+
self.assertEqual(srv_proto.data, self.DATA)
419+
379420
def test_sendfile_force_fallback(self):
380421
srv_proto, cli_proto = self.prepare_sendfile()
381422

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,2 @@
1+
Fix transport cleanup when cancelling :meth:`asyncio.loop.sendfile` while
2+
waiting for the write buffer to drain in the native implementation.

0 commit comments

Comments
 (0)