From 0e433417215a11d4df02cfee62d719b8ca7e67d9 Mon Sep 17 00:00:00 2001 From: Thomas Grainger Date: Mon, 5 Oct 2026 10:54:06 +0100 Subject: [PATCH 1/6] Don't report queued datagrams to the protocol after UDPTransport.abort() Aborting a UDPTransport while a datagram was queued reported a fatal CancelledError and a resume_writing() AttributeError to the loop's exception handler, because libuv cancels queued sends when the handle closes and the send callback treated that as an error after the protocol had been detached. Also, a queued datagram that was sent between abort() and the transport closing still resumed the protocol. Like asyncio, stop reporting send results to the protocol once the connection is lost. Fixes #771 Co-Authored-By: Claude Opus 5.5 --- tests/test_udp.py | 81 ++++++++++++++++++++++++++++++++++++++++++ uvloop/handles/udp.pyx | 10 ++++++ 2 files changed, 91 insertions(+) diff --git a/tests/test_udp.py b/tests/test_udp.py index e529b9e33..f1b46de99 100644 --- a/tests/test_udp.py +++ b/tests/test_udp.py @@ -251,6 +251,87 @@ def connection_lost(self, exc): self.assertIn(tmp_file2, pr.addrs) + def _close_with_queued_datagram(self, method, drain): + # Close the transport while a datagram is queued (the OS refused + # it with EAGAIN), optionally letting the peer make room for it + # before the loop runs again, and return the protocol events. + # Any call to the loop's exception handler fails the test. + + class Proto(asyncio.DatagramProtocol): + def __init__(self, loop): + self.events = [] + self.done = asyncio.Future(loop=loop) + + def connection_made(self, transport): + transport.set_write_buffer_limits(0) + + def pause_writing(self): + self.events.append('pause_writing') + + def resume_writing(self): + self.events.append('resume_writing') + + def error_received(self, exc): + self.events.append(('error_received', exc)) + + def connection_lost(self, exc): + self.events.append(('connection_lost', exc)) + self.done.set_result(None) + + async def run(peer, path): + sock = socket.socket(socket.AF_UNIX, socket.SOCK_DGRAM) + sock.connect(path) + sock.setblocking(False) + pr = Proto(self.loop) + tr, _ = await self.loop.create_datagram_endpoint( + lambda: pr, sock=sock) + + while not tr.get_write_buffer_size(): + tr.sendto(b'x' * 64) + + getattr(tr, method)() + pr.events.append(method) + if drain: + try: + while True: + peer.recv(64) + except BlockingIOError: + pass + + await pr.done + await asyncio.sleep(0.1) + return pr.events + + with tempfile.TemporaryDirectory() as tmp_dir: + path = os.path.join(tmp_dir, 'peer.sock') + # A UNIX datagram socket is used because UDP over loopback + # never refuses a datagram, so the sender's queue never fills. + with socket.socket(socket.AF_UNIX, socket.SOCK_DGRAM) as peer: + peer.bind(path) + peer.setblocking(False) + return self.loop.run_until_complete(run(peer, path)) + + def test_abort_with_queued_datagram(self): + # gh-771: the queued datagram is discarded silently. + self.assertEqual( + self._close_with_queued_datagram('abort', drain=False), + ['pause_writing', 'abort', ('connection_lost', None)]) + + def test_abort_with_queued_datagram_then_writable(self): + # The queued datagram can be sent before connection_lost() runs; + # the protocol must not hear about it after abort(). + self.assertEqual( + self._close_with_queued_datagram('abort', drain=True), + ['pause_writing', 'abort', ('connection_lost', None)]) + + def test_close_with_queued_datagram_then_writable(self): + # close() flushes the queue, so the protocol is resumed before + # connection_lost(). + self.assertEqual( + self._close_with_queued_datagram('close', drain=True), + ['pause_writing', 'close', 'resume_writing', + ('connection_lost', None)]) + def test_create_datagram_1(self): server_addr = ('127.0.0.1', 8888) client_addr = ('127.0.0.1', 0) diff --git a/uvloop/handles/udp.pyx b/uvloop/handles/udp.pyx index eac1bca51..e6863ca75 100644 --- a/uvloop/handles/udp.pyx +++ b/uvloop/handles/udp.pyx @@ -267,6 +267,11 @@ cdef class UDPTransport(UVBaseTransport): run_in_context1(self.context, self._protocol.error_received, exc) cdef _on_sent(self, object exc, object context=None): + if self._conn_lost: + # abort() or a fatal error discarded the write buffer; like + # asyncio, don't report anything more to the protocol. + return + if exc is not None: if isinstance(exc, OSError): if context is None: @@ -400,6 +405,11 @@ cdef void __uv_udp_on_send( ctx.close() + if status == uv.UV_ECANCELED and udp._closed: + # The handle is being closed (e.g. by abort()) and libuv + # cancelled the queued send; the datagram is simply discarded. + return + if status < 0: exc = convert_error(status) print(exc) From b4c4762f055420b963caa1e7ee807db9091c7e5c Mon Sep 17 00:00:00 2001 From: Thomas Grainger Date: Mon, 5 Oct 2026 11:16:42 +0100 Subject: [PATCH 2/6] Test loop.close() with a queued datagram; drop unneeded sleep The UV_ECANCELED check in the send callback was only exercised by closing the loop with a UDP transport still open, as abort() is now handled by the _conn_lost check in _on_sent. Co-Authored-By: Claude Opus 5.5 --- tests/test_udp.py | 36 +++++++++++++++++++++++++++++++++++- 1 file changed, 35 insertions(+), 1 deletion(-) diff --git a/tests/test_udp.py b/tests/test_udp.py index f1b46de99..74fcf4a95 100644 --- a/tests/test_udp.py +++ b/tests/test_udp.py @@ -299,7 +299,6 @@ async def run(peer, path): pass await pr.done - await asyncio.sleep(0.1) return pr.events with tempfile.TemporaryDirectory() as tmp_dir: @@ -432,6 +431,41 @@ def test_create_datagram_endpoint_reuse_address_warning(self): class Test_UV_UDP(_TestUDP, tb.UVTestCase): + def test_loop_close_with_queued_datagram(self): + # Closing the loop closes the transport's handle directly, which + # cancels the queued datagram; that must not be reported to the + # protocol or the exception handler. + events = [] + + class Proto(asyncio.DatagramProtocol): + def connection_made(self, transport): + transport.set_write_buffer_limits(0) + + def resume_writing(self): + events.append('resume_writing') + + def error_received(self, exc): + events.append(('error_received', exc)) + + with tempfile.TemporaryDirectory() as tmp_dir: + path = os.path.join(tmp_dir, 'peer.sock') + with socket.socket(socket.AF_UNIX, socket.SOCK_DGRAM) as peer: + peer.bind(path) + + sock = socket.socket(socket.AF_UNIX, socket.SOCK_DGRAM) + sock.connect(path) + sock.setblocking(False) + tr, _ = self.loop.run_until_complete( + self.loop.create_datagram_endpoint(Proto, sock=sock)) + + while not tr.get_write_buffer_size(): + tr.sendto(b'x' * 64) + + with self.assertWarnsRegex(ResourceWarning, 'unclosed'): + self.loop.close() + + self.assertEqual(events, []) + def test_create_datagram_endpoint_wrong_sock(self): sock = socket.socket(socket.AF_INET) with sock: From fd02ac30c53e8c5e5d9e7b2650c08e2603446b52 Mon Sep 17 00:00:00 2001 From: Thomas Grainger Date: Mon, 5 Oct 2026 11:21:37 +0100 Subject: [PATCH 3/6] Close the test sockets with a with block Co-Authored-By: Claude Opus 5.5 --- tests/test_udp.py | 55 ++++++++++++++++++++++++----------------------- 1 file changed, 28 insertions(+), 27 deletions(-) diff --git a/tests/test_udp.py b/tests/test_udp.py index 74fcf4a95..beb7d8044 100644 --- a/tests/test_udp.py +++ b/tests/test_udp.py @@ -279,27 +279,27 @@ def connection_lost(self, exc): self.done.set_result(None) async def run(peer, path): - sock = socket.socket(socket.AF_UNIX, socket.SOCK_DGRAM) - sock.connect(path) - sock.setblocking(False) - pr = Proto(self.loop) - tr, _ = await self.loop.create_datagram_endpoint( - lambda: pr, sock=sock) + with socket.socket(socket.AF_UNIX, socket.SOCK_DGRAM) as sock: + sock.connect(path) + sock.setblocking(False) + pr = Proto(self.loop) + tr, _ = await self.loop.create_datagram_endpoint( + lambda: pr, sock=sock) - while not tr.get_write_buffer_size(): - tr.sendto(b'x' * 64) + while not tr.get_write_buffer_size(): + tr.sendto(b'x' * 64) - getattr(tr, method)() - pr.events.append(method) - if drain: - try: - while True: - peer.recv(64) - except BlockingIOError: - pass + getattr(tr, method)() + pr.events.append(method) + if drain: + try: + while True: + peer.recv(64) + except BlockingIOError: + pass - await pr.done - return pr.events + await pr.done + return pr.events with tempfile.TemporaryDirectory() as tmp_dir: path = os.path.join(tmp_dir, 'peer.sock') @@ -452,17 +452,18 @@ def error_received(self, exc): with socket.socket(socket.AF_UNIX, socket.SOCK_DGRAM) as peer: peer.bind(path) - sock = socket.socket(socket.AF_UNIX, socket.SOCK_DGRAM) - sock.connect(path) - sock.setblocking(False) - tr, _ = self.loop.run_until_complete( - self.loop.create_datagram_endpoint(Proto, sock=sock)) + with socket.socket( + socket.AF_UNIX, socket.SOCK_DGRAM) as sock: + sock.connect(path) + sock.setblocking(False) + tr, _ = self.loop.run_until_complete( + self.loop.create_datagram_endpoint(Proto, sock=sock)) - while not tr.get_write_buffer_size(): - tr.sendto(b'x' * 64) + while not tr.get_write_buffer_size(): + tr.sendto(b'x' * 64) - with self.assertWarnsRegex(ResourceWarning, 'unclosed'): - self.loop.close() + with self.assertWarnsRegex(ResourceWarning, 'unclosed'): + self.loop.close() self.assertEqual(events, []) From 6bd03ac733f8672aaa6bb711e6981c86f0d57afd Mon Sep 17 00:00:00 2001 From: Thomas Grainger Date: Mon, 5 Oct 2026 12:05:29 +0100 Subject: [PATCH 4/6] Stop receiving when a UDPTransport handle is closed directly UDPTransport keeps a reference to itself while receiving, which was only dropped by _stop_reading(). loop.close() closes leftover handles with _close(), so a UDP transport still open at that point was never freed. Stop reading in _close(), as UVStream does. Co-Authored-By: Claude Opus 5.5 --- uvloop/handles/udp.pxd | 2 ++ uvloop/handles/udp.pyx | 8 ++++++++ 2 files changed, 10 insertions(+) diff --git a/uvloop/handles/udp.pxd b/uvloop/handles/udp.pxd index daa9a1bee..e1302154d 100644 --- a/uvloop/handles/udp.pxd +++ b/uvloop/handles/udp.pxd @@ -13,6 +13,8 @@ cdef class UDPTransport(UVBaseTransport): cdef open(self, int family, int sockfd) cdef _set_broadcast(self, bint on) + cdef _close(self) + cdef inline __receiving_started(self) cdef inline __receiving_stopped(self) diff --git a/uvloop/handles/udp.pyx b/uvloop/handles/udp.pyx index e6863ca75..61a2c699c 100644 --- a/uvloop/handles/udp.pyx +++ b/uvloop/handles/udp.pyx @@ -171,6 +171,14 @@ cdef class UDPTransport(UVBaseTransport): else: self.__receiving_stopped() + cdef _close(self): + try: + # Drop the reference held while receiving, even when the + # handle is closed directly (e.g. by loop.close()). + self._stop_reading() + finally: + UVSocketHandle._close(self) + cdef inline __receiving_started(self): if self.__receiving: return From 9a0faf838c4979491a2769cd30e3e19d20114a55 Mon Sep 17 00:00:00 2001 From: Thomas Grainger Date: Mon, 5 Oct 2026 12:17:27 +0100 Subject: [PATCH 5/6] Make the queued datagram tests terminate on macOS The tests sent datagrams until the transport queued one, which never happened on macOS CI and spun until the job was killed. Use a small send buffer, as anyio's UNIX datagram tests do, and give up after a bounded number of sends. Co-Authored-By: Claude Opus 5.5 --- tests/test_udp.py | 23 +++++++++++++++++++---- 1 file changed, 19 insertions(+), 4 deletions(-) diff --git a/tests/test_udp.py b/tests/test_udp.py index beb7d8044..71728ec38 100644 --- a/tests/test_udp.py +++ b/tests/test_udp.py @@ -251,6 +251,15 @@ def connection_lost(self, exc): self.assertIn(tmp_file2, pr.addrs) + def _fill_send_queue(self, tr): + # Send until the OS refuses a datagram and the transport has to + # queue it. + for _ in range(10000): + if tr.get_write_buffer_size(): + return + tr.sendto(b'x' * 64) + self.fail('the OS never refused a datagram') + def _close_with_queued_datagram(self, method, drain): # Close the transport while a datagram is queued (the OS refused # it with EAGAIN), optionally letting the peer make room for it @@ -280,14 +289,17 @@ def connection_lost(self, exc): async def run(peer, path): with socket.socket(socket.AF_UNIX, socket.SOCK_DGRAM) as sock: + # Keep the send buffer small (as anyio's tests do) so that + # macOS refuses datagrams with EAGAIN too. + sock.setsockopt( + socket.SOL_SOCKET, socket.SO_SNDBUF, 1024) sock.connect(path) sock.setblocking(False) pr = Proto(self.loop) tr, _ = await self.loop.create_datagram_endpoint( lambda: pr, sock=sock) - while not tr.get_write_buffer_size(): - tr.sendto(b'x' * 64) + self._fill_send_queue(tr) getattr(tr, method)() pr.events.append(method) @@ -454,13 +466,16 @@ def error_received(self, exc): with socket.socket( socket.AF_UNIX, socket.SOCK_DGRAM) as sock: + # Keep the send buffer small (as anyio's tests do) so that + # macOS refuses datagrams with EAGAIN too. + sock.setsockopt( + socket.SOL_SOCKET, socket.SO_SNDBUF, 1024) sock.connect(path) sock.setblocking(False) tr, _ = self.loop.run_until_complete( self.loop.create_datagram_endpoint(Proto, sock=sock)) - while not tr.get_write_buffer_size(): - tr.sendto(b'x' * 64) + self._fill_send_queue(tr) with self.assertWarnsRegex(ResourceWarning, 'unclosed'): self.loop.close() From 26651d5fd31dd6c25c54daefd7dadbfaf4d726f8 Mon Sep 17 00:00:00 2001 From: Thomas Grainger Date: Mon, 5 Oct 2026 12:30:46 +0100 Subject: [PATCH 6/6] Skip the asyncio queued datagram tests on macOS macOS refuses UNIX datagrams with ENOBUFS rather than EAGAIN. libuv treats ENOBUFS like EAGAIN and queues the datagram, but asyncio reports it to error_received() and drops it, so the send queue never fills. The small send buffer did not change that, so drop it again. Co-Authored-By: Claude Opus 5.5 --- tests/test_udp.py | 14 ++++++-------- 1 file changed, 6 insertions(+), 8 deletions(-) diff --git a/tests/test_udp.py b/tests/test_udp.py index 71728ec38..94018053c 100644 --- a/tests/test_udp.py +++ b/tests/test_udp.py @@ -265,6 +265,12 @@ def _close_with_queued_datagram(self, method, drain): # it with EAGAIN), optionally letting the peer make room for it # before the loop runs again, and return the protocol events. # Any call to the loop's exception handler fails the test. + if self.implementation == 'asyncio' and sys.platform == 'darwin': + # macOS refuses UNIX datagrams with ENOBUFS rather than + # EAGAIN; asyncio reports that to error_received() and drops + # the datagram instead of queueing it. + raise unittest.SkipTest( + 'asyncio does not queue datagrams refused with ENOBUFS') class Proto(asyncio.DatagramProtocol): def __init__(self, loop): @@ -289,10 +295,6 @@ def connection_lost(self, exc): async def run(peer, path): with socket.socket(socket.AF_UNIX, socket.SOCK_DGRAM) as sock: - # Keep the send buffer small (as anyio's tests do) so that - # macOS refuses datagrams with EAGAIN too. - sock.setsockopt( - socket.SOL_SOCKET, socket.SO_SNDBUF, 1024) sock.connect(path) sock.setblocking(False) pr = Proto(self.loop) @@ -466,10 +468,6 @@ def error_received(self, exc): with socket.socket( socket.AF_UNIX, socket.SOCK_DGRAM) as sock: - # Keep the send buffer small (as anyio's tests do) so that - # macOS refuses datagrams with EAGAIN too. - sock.setsockopt( - socket.SOL_SOCKET, socket.SO_SNDBUF, 1024) sock.connect(path) sock.setblocking(False) tr, _ = self.loop.run_until_complete(