From dca1b5fe2d1be3dc4ece5ffc51afca902821ec38 Mon Sep 17 00:00:00 2001 From: Cory Lowe Date: Sun, 30 Aug 2026 21:09:10 -0500 Subject: [PATCH] Close the QUIC datagram transports on shutdown worker_serve keeps its TCP servers in `servers` and closes them in its finally, but the transport returned by create_datagram_endpoint was discarded, so shutdown could not close what it never kept. Each QUIC-serving worker leaked its datagram transport, visible as a GC-time "unclosed transport" ResourceWarning (python -W always). Keep the datagram transports and close them after the graceful drain, alongside the existing server close. Closing after the drain rather than before it preserves in-flight QUIC traffic during graceful_timeout: unlike a TCP Server, the datagram transport is both the listener and the data path for established connections. Co-Authored-By: Claude Fable 5 --- src/hypercorn/asyncio/run.py | 10 +++++++- tests/asyncio/test_run.py | 45 ++++++++++++++++++++++++++++++++++++ 2 files changed, 54 insertions(+), 1 deletion(-) create mode 100644 tests/asyncio/test_run.py diff --git a/src/hypercorn/asyncio/run.py b/src/hypercorn/asyncio/run.py index 09fec25..246daae 100644 --- a/src/hypercorn/asyncio/run.py +++ b/src/hypercorn/asyncio/run.py @@ -135,13 +135,15 @@ async def _server_callback(reader: asyncio.StreamReader, writer: asyncio.StreamW bind = repr_socket_addr(sock.family, sock.getsockname()) await config.log.info(f"Running on http://{bind} (CTRL + C to quit)") + datagram_transports: list[asyncio.DatagramTransport] = [] for sock in sockets.quic_sockets: if config.workers > 1 and platform.system() == "Windows": sock = _share_socket(sock) - _, protocol = await loop.create_datagram_endpoint( + transport, protocol = await loop.create_datagram_endpoint( lambda: UDPServer(app, loop, config, context, lifespan_state), sock=sock ) + datagram_transports.append(transport) task = loop.create_task(protocol.run()) server_tasks.add(task) task.add_done_callback(server_tasks.discard) @@ -175,6 +177,12 @@ async def _server_callback(reader: asyncio.StreamReader, writer: asyncio.StreamW # prevent a warning that this hasn't been done. gathered_server_tasks.exception() + # Datagram transports are not Servers and hence aren't closed + # by the server close above; without this they outlive the + # worker and are reported unclosed. + for transport in datagram_transports: + transport.close() + await lifespan.wait_for_shutdown() lifespan_task.cancel() await lifespan_task diff --git a/tests/asyncio/test_run.py b/tests/asyncio/test_run.py new file mode 100644 index 0000000..9771277 --- /dev/null +++ b/tests/asyncio/test_run.py @@ -0,0 +1,45 @@ +from __future__ import annotations + +import asyncio + +import pytest + +from hypercorn.app_wrappers import ASGIWrapper +from hypercorn.asyncio.run import worker_serve +from hypercorn.config import Config +from ..helpers import sanity_framework + + +@pytest.mark.asyncio +async def test_worker_serve_closes_datagram_transports() -> None: + pytest.importorskip("aioquic") + event_loop = asyncio.get_running_loop() + config = Config() + config.bind = [] + config.quic_bind = ["127.0.0.1:0"] + config.certfile = "tests/assets/cert.pem" + config.keyfile = "tests/assets/key.pem" + config.graceful_timeout = 0.1 + sockets = config.create_sockets() + + transports: list[asyncio.BaseTransport] = [] + create_datagram_endpoint = event_loop.create_datagram_endpoint + + async def _capturing(*args: object, **kwargs: object) -> tuple: + transport, protocol = await create_datagram_endpoint(*args, **kwargs) # type: ignore + transports.append(transport) + return transport, protocol + + async def _shutdown() -> None: + return None + + event_loop.create_datagram_endpoint = _capturing # type: ignore + try: + await worker_serve( + ASGIWrapper(sanity_framework), config, sockets=sockets, shutdown_trigger=_shutdown + ) + finally: + event_loop.create_datagram_endpoint = create_datagram_endpoint # type: ignore + + assert transports != [] + assert all(transport.is_closing() for transport in transports)