Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 9 additions & 1 deletion src/hypercorn/asyncio/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

One cleanup edge remains here: if a later QUIC endpoint fails to start, control never reaches the try below, so transports already appended to datagram_transports stay open.

I reproduced this on dca1b5f with two fake QUIC sockets: the first create_datagram_endpoint() returned a tracked transport, the second raised RuntimeError, and after worker_serve() exited the first transport still reported is_closing() == False. The new test covers normal shutdown only. Moving endpoint setup under the cleanup try (and closing partially created transports on startup failure) would cover this path too.

Disclosure: I ran the changed worker through Lumi Trace at NOQT for review context, then verified this behavior with the focused runtime repro above.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for checking this. Agreed that a failure while starting a later QUIC endpoint leaves earlier transports open — but that is the existing shape of worker_serve rather than something this change adds: the TCP servers list a few lines up has the same gap, since setup for both runs before the cleanup try. I kept this PR to the one observable defect (transports discarded on the normal shutdown path). Happy to do the startup-failure cleanup for both the TCP servers and the datagram transports as a follow-up if @pgjones would like it.

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)
Expand Down Expand Up @@ -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
Expand Down
45 changes: 45 additions & 0 deletions tests/asyncio/test_run.py
Original file line number Diff line number Diff line change
@@ -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)