SWE-Race › Tasks › mitmproxy-mitmproxy-5589 ← prevnext →

mitmproxy-mitmproxy-5589

mitmproxy/mitmproxyhardsinglemerged 2022-09-22MITfix: 2 files, +17 −51 fail-to-pass · 1 pass-to-pass
Results
Modelsolved / attemptsmedian stepsmedian costattempts
GPT-5.6 Luna3/616$0.0171✗ 2✓ 3✗ 4✓ 5✓ 6✗
DeepSeek V4 Flash0/283$0.1811✗ 2✗
GLM-5.3 Flash0/288$0.0641✗ 2✗
The prompt the agent sees

UDP client/server connections can be established and used, but closing them may fail to complete cleanly because the connection-handling code attempts to write to streams that are already closing. When a caller closes both sides of a UDP client/server setup, the corresponding client and server writers should transition to the closing state without errors or hangs.

Hidden tests · 1 fail-to-pass, 1 pass-to-passrun after the agent submits, in a clean verifier
test_client_server
Test patch · 21 lines
diff --git a/test/mitmproxy/net/test_udp.py b/test/mitmproxy/net/test_udp.py
index 6e5a9f6237..f90538f559 100644
--- a/test/mitmproxy/net/test_udp.py
+++ b/test/mitmproxy/net/test_udp.py
@@ -45,11 +45,16 @@ def handle_datagram(
     server.resume_writing()
     await server.drain()
 
+    assert not client_writer.is_closing()
+    assert not server_writer.is_closing()
+
     assert await client_reader.read(MAX_DATAGRAM_SIZE) == b"msg4"
     client_writer.close()
+    assert client_writer.is_closing()
     await client_writer.wait_closed()
 
     server_writer.close()
+    assert server_writer.is_closing()
     await server_writer.wait_closed()
 
     server.close()
Reference fix · 2 files, +17 −5the upstream merge, used only for grading calibration

The agent could not see this: the repository holds one commit and the sandbox has no network. Leak audit.

mitmproxy/net/udp.py, mitmproxy/proxy/server.py

diff --git a/mitmproxy/net/udp.py b/mitmproxy/net/udp.py
index 328fcba751..3590d42843 100644
--- a/mitmproxy/net/udp.py
+++ b/mitmproxy/net/udp.py
@@ -175,8 +175,12 @@ def __init__(
         """
         self._transport = transport
         self._remote_addr = remote_addr
-        self._reader = reader
-        self._closed = asyncio.Event() if reader is not None else None
+        if reader is not None:
+            self._reader = reader
+            self._closed = asyncio.Event()
+        else:
+            self._reader = None
+            self._closed = None
 
     @property
     def _protocol(self) -> DrainableDatagramProtocol | udp_wireguard.WireGuardDatagramTransport:
@@ -199,9 +203,15 @@ def close(self) -> None:
             self._transport.close()
         else:
             self._closed.set()
-        if self._reader is not None:
+            assert self._reader
             self._reader.feed_eof()
 
+    def is_closing(self) -> bool:
+        if self._closed is None:
+            return self._transport.is_closing()
+        else:
+            return self._closed.is_set()
+
     async def wait_closed(self) -> None:
         if self._closed is None:
             await self._protocol.wait_closed()
diff --git a/mitmproxy/proxy/server.py b/mitmproxy/proxy/server.py
index 06033e4272..1341564ada 100644
--- a/mitmproxy/proxy/server.py
+++ b/mitmproxy/proxy/server.py
@@ -367,7 +367,8 @@ def server_event(self, event: events.Event) -> None:
                 elif isinstance(command, commands.SendData):
                     writer = self.transports[command.connection].writer
                     assert writer
-                    writer.write(command.data)
+                    if not writer.is_closing():
+                        writer.write(command.data)
                 elif isinstance(command, commands.CloseConnection):
                     self.close_connection(command.connection, command.half_close)
                 elif isinstance(command, commands.StartHook):
@@ -393,7 +394,8 @@ def close_connection(
             try:
                 writer = self.transports[connection].writer
                 assert writer
-                writer.write_eof()
+                if not writer.is_closing():
+                    writer.write_eof()
             except OSError:
                 # if we can't write to the socket anymore we presume it completely dead.
                 connection.state = ConnectionState.CLOSED

What each model did

Submitted diff first, then every command. Outputs show 12 lines; expand for the rest.
GPT-5.6 Luna · attempt 3 · not solved · 16 steps · $0.018
Submitted diff · 60 lines
diff --git a/mitmproxy/net/udp.py b/mitmproxy/net/udp.py
index 328fcba..562d54e 100644
--- a/mitmproxy/net/udp.py
+++ b/mitmproxy/net/udp.py
@@ -55,6 +55,8 @@ class DrainableDatagramProtocol(asyncio.DatagramProtocol):
 
     def connection_lost(self, exc: Exception | None) -> None:
         self._closed.set()
+        # A paused transport must not leave drain() waiting forever after it closes.
+        self._can_write.set()
         if exc:
             logger.warning(f"Connection lost on {self!r}: {exc!r}")  # pragma: no cover
 
@@ -162,6 +164,7 @@ class DatagramWriter:
     _remote_addr: Address
     _reader: DatagramReader | None
     _closed: asyncio.Event | None
+    _protocol_instance: DrainableDatagramProtocol | udp_wireguard.WireGuardDatagramTransport
 
     def __init__(
         self,
@@ -174,13 +177,16 @@ class DatagramWriter:
         Specify a reader to prevent closing the transport and instead only feed EOF to the reader.
         """
         self._transport = transport
+        self._protocol_instance = transport.get_protocol()  # type: ignore
         self._remote_addr = remote_addr
         self._reader = reader
         self._closed = asyncio.Event() if reader is not None else None
 
     @property
     def _protocol(self) -> DrainableDatagramProtocol | udp_wireguard.WireGuardDatagramTransport:
-        return self._transport.get_protocol()  # type: ignore
+        # asyncio clears the protocol from a transport while it is closing. Keep
+        # the protocol captured at construction so wait_closed/drain remain safe.
+        return self._protocol_instance
 
     def write(self, data: bytes) -> None:
         self._transport.sendto(data, self._remote_addr)
diff --git a/mitmproxy/proxy/server.py b/mitmproxy/proxy/server.py
index 06033e4..73869eb 100644
--- a/mitmproxy/proxy/server.py
+++ b/mitmproxy/proxy/server.py
@@ -365,9 +365,13 @@ class ConnectionHandler(metaclass=abc.ABCMeta):
                 ):
                     pass  # The connection has already been closed.
                 elif isinstance(command, commands.SendData):
-                    writer = self.transports[command.connection].writer
-                    assert writer
-                    writer.write(command.data)
+                    # A peer may close while a layer still has a queued send
+                    # command (UDP in particular has no half-close handshake).
+                    # Do not write to a transport that is already closing.
+                    if command.connection.state & ConnectionState.CAN_WRITE:
+                        writer = self.transports[command.connection].writer
+                        assert writer
+                        writer.write(command.data)
                 elif isinstance(command, commands.CloseConnection):
                     self.close_connection(command.connection, command.half_close)
                 elif isinstance(command, commands.StartHook):
step 1

Plan:

$ pwd && ls -la && find . -maxdepth 2 -type f | head -80 && git status --short
output · 40 lines
/repo
total 92
drwxr-xr-x 1 root root  4096 Oct  1 00:26 .
drwxr-xr-x 1 root root  4096 Oct  1 00:25 ..
drwxr-xr-x 1 root root  4096 Oct  1 00:25 .git
-rw-r--r-- 1 root root   113 Oct  1 00:24 .gitattributes
drwxr-xr-x 1 root root  4096 Oct  1 00:24 .github
-rw-r--r-- 1 root root   348 Oct  1 00:24 .gitignore
-rw-r--r-- 1 root root 54466 Oct  1 00:24 CHANGELOG.md
-rw-r--r-- 1 root root  3532 Oct  1 00:24 CONTRIBUTING.md
-rw-r--r-- 1 root root  1079 Oct  1 00:24 LICENSE
-rw-r--r-- 1 root root    66 Oct  1 00:24 MANIFEST.in
28 more lines
-rw-r--r-- 1 root root  2218 Oct  1 00:24 README.md
-rw-r--r-- 1 root root   828 Oct  1 00:24 SECURITY.md
-rw-r--r-- 1 root root   346 Oct  1 00:24 codecov.yml
drwxr-xr-x 1 root root  4096 Oct  1 00:24 docs
drwxr-xr-x 1 root root  4096 Oct  1 00:24 examples
drwxr-xr-x 1 root root  4096 Oct  1 00:24 mitmproxy
drwxr-xr-x 1 root root  4096 Oct  1 00:26 mitmproxy.egg-inf
…[1568 chars omitted]…
release/build.py
./release/deploy-microsoft-store.py
./release/deploy.py
./release/README.md
./SECURITY.md
./setup.cfg
./setup.py
./test/conftest.py
./test/filename_matching.py
./test/full_coverage_plugin.py
./test/individual_coverage.py
./test/__init__.py
./tox.ini
./web/.editorconfig
./web/gulpfile.js
./web/jest.config.js
./web/package-lock.json
./web/package.json
./web/README.md
./web/tsconfig.json
step 2
$ find test -iname '*udp*' -o -iname '*connection*' | sort
output · 5 lines
test/mitmproxy/net/test_udp.py
test/mitmproxy/net/test_udp_wireguard.py
test/mitmproxy/proxy/layers/test_udp.py
test/mitmproxy/test_connection.py
test/mitmproxy/test_udp.py
step 3
$ cat test/mitmproxy/net/test_udp.py && grep -R "start_server\|open_connection" -n test/mitmproxy/net test/mitmproxy/proxy | head -100
output · 31 lines
import asyncio
from typing import Optional
import pytest
from mitmproxy.connection import Address
from mitmproxy.net.udp import MAX_DATAGRAM_SIZE, DatagramReader, DatagramWriter, open_connection, start_server


async def test_client_server():
    server_reader = DatagramReader()
    server_writer: Optional[DatagramWriter] = None

    def handle_datagram(
19 more lines
        transport: asyncio.DatagramTransport,
        data: bytes,
        remote_addr: Address,
        local_addr: Address
    ):
        nonlocal server_reader, server_writer
        if server_writer is None:
            server_writer = DatagramWriter(transport, remote_addr, server_reader)
        server_reader.feed_data(data, remote_addr)

    server = await start_server(handle_datagram, "127.0.0.1", 0)
    assert repr(server).startswith("<UdpServer socket=")

    [client_reader, client_writer] = await open_connection(*server.sock
…[4036 chars omitted]…
t asyncio.open_connection(host, port)
test/mitmproxy/proxy/test_mode_servers.py:98:        reader, writer = await asyncio.open_connection(host, port)
test/mitmproxy/proxy/test_mode_servers.py:202:    server = await asyncio.start_server(MagicMock(), host="127.0.0.1", port=0, reuse_address=False)
test/mitmproxy/proxy/test_mode_servers.py:237:        reader, writer = await udp.open_connection(host, port)
step 4
$ sed -n '200,280p' test/mitmproxy/proxy/test_mode_servers.py; sed -n '1,130p' test/mitmproxy/proxy/layers/test_udp.py
output · 37 lines
manager = MagicMock()

    server = await asyncio.start_server(MagicMock(), host="127.0.0.1", port=0, reuse_address=False)
    port = server.sockets[0].getsockname()[1]

    with taddons.context() as tctx:
        inst = ServerInstance.make(f"regular@127.0.0.1:{port}", manager)
        with pytest.raises(OSError, match=f"proxy failed to listen on 127\\.0\\.0\\.1:{port}"):
            await inst.start()
        tctx.options.listen_host = "127.0.0.1"
        tctx.options.listen_port = port
        inst3 = ServerInstance.make(f"regular", manager)
25 more lines
        with pytest.raises(OSError):
            await inst3.start()


async def test_invalid_protocol(monkeypatch):
    manager = MagicMock()

    with taddons.context():
        inst = ServerInstance.make(f"regular@127.0.0.1:0", manager)
        monkeypatch.setattr(inst.mode, "transport_protocol", "invalid_proto")
        with pytest.raise
…[4953 chars omitted]…
(f)
        >> reply()
        << SendData(tctx.server, b"hello!")
        # and the other way...
        >> UdpMessageInjected(
            f, UDPMessage(False, b"I have already done the greeting for you.")
        )
        << udp.UdpMessageHook(f)
        >> reply()
        << SendData(tctx.client, b"I have already done the greeting for you.")
        << None
    )
    assert len(f().messages) == 2
step 5
$ cat mitmproxy/proxy/commands.py; sed -n '1,290p' mitmproxy/proxy/server.py; grep -R "ConnectionState" -n mitmproxy/proxy | head -80
output · 35 lines
"""
Commands make it possible for layers to communicate with the "outer world",
e.g. to perform IO or to ask the master.
A command is issued by a proxy layer and is then passed upwards to the proxy server, and from there
possibly to the master and addons.

The counterpart to commands are events.
"""
import logging
import warnings
from typing import Union, TYPE_CHECKING
23 more lines
import mitmproxy.hooks
from mitmproxy.connection import Connection, Server

if TYPE_CHECKING:
    import mitmproxy.proxy.layer


class Command:
    """
    Base class for all commands
    """

    blocking: Union[bool, "mitmproxy.proxy.layer.Layer"] = False
    """
    Determines if the command blocks until it has been completed.
    For practical purposes, this attribute should be thought of as a boolean value,
    layers may swap out `True` with a reference to themselves to signal to outer layers
    that they do not
…[8629 chars omitted]…
                       self.conn.state = connection.ConnectionState.CLOSED
mitmproxy/proxy/tunnel.py:82:                    self.conn.state &= ~connection.ConnectionState.CAN_READ
mitmproxy/proxy/tunnel.py:128:                            self.conn.state &= ~connection.ConnectionState.CAN_WRITE
mitmproxy/proxy/tunnel.py:130:                            self.conn.state = connection.ConnectionState.CLOSED
step 6
$ sed -n '180,245p' mitmproxy/proxy/server.py; sed -n '350,420p' mitmproxy/proxy/server.py; sed -n '1,170p' mitmproxy/proxy/layers/tcp.py; grep -R "sendto\|DatagramWriter" -n mitmproxy/proxy/server.py mitmproxy/proxy | head -50
output · 24 lines
return

        async with self.max_conns[command.connection.address]:
            reader: Union[asyncio.StreamReader, udp.DatagramReader]
            writer: Union[asyncio.StreamWriter, udp.DatagramWriter]
            try:
                command.connection.timestamp_start = time.time()
                if command.connection.transport_protocol == "tcp":
                    reader, writer = await asyncio.open_connection(
                        *command.connection.address,
                        local_addr=command.connection.sockname,
                    )
12 more lines
                elif command.connection.transport_protocol == "udp":
                    reader, writer = await udp.open_connection(
                        *command.connection.address,
                        local_addr=command.connection.sockname,
                    )
                else:
                    raise As
…[8629 chars omitted]…
ver.py:184:            writer: Union[asyncio.StreamWriter, udp.DatagramWriter]
mitmproxy/proxy/mode_servers.py:180:            writer = udp.DatagramWriter(transport, remote_addr, reader)
mitmproxy/proxy/server.py:84:    writer: Optional[Union[asyncio.StreamWriter, udp.DatagramWriter, wg.TcpStream]] = None
mitmproxy/proxy/server.py:184:            writer: Union[asyncio.StreamWriter, udp.DatagramWriter]
step 7
$ sed -n '375,420p' mitmproxy/proxy/server.py; sed -n '130,210p' mitmproxy/proxy/mode_servers.py; grep -R "def close" -n mitmproxy/net mitmproxy/proxy | head -30
output · 30 lines
self.hook_task(command),
                        name=f"handle_hook({command.name})",
                        client=self.client.peername,
                    )
                elif isinstance(command, commands.Log):
                    self.log(command.message, command.level)
                else:
                    raise RuntimeError(f"Unexpected command: {command}")
        except Exception:
            self.log(f"mitmproxy has crashed!\n{traceback.format_exc()}", logging.ERROR)

    def close_connection(
18 more lines
        self, connection: Connection, half_close: bool = False
    ) -> None:
        if half_close:
            if not connection.state & ConnectionState.CAN_WRITE:
                return
            self.log(f"half-closing {connection}", logging.DEBUG)
            try:
                writer = self.transports[connection].writer
                assert wri
…[4438 chars omitted]…
  @property
    def is_running(self) -> bool:
        return self._server is not None
mitmproxy/net/tls.py:70:    def close(self):
mitmproxy/net/udp.py:111:    def close(self) -> None:
mitmproxy/net/udp.py:197:    def close(self) -> None:
mitmproxy/proxy/layers/http/_http2.py:280:    def close_connection(self, msg: str) -> CommandGenerator[None]:
mitmproxy/proxy/server.py:386:    def close_connection(
step 8
$ pytest -q test/mitmproxy/net/test_udp.py test/mitmproxy/proxy/layers/test_udp.py test/mitmproxy/proxy/test_mode_servers.py -q --disable-warnings --maxfail=1 --timeout=20 2>&1 | tail -40
output · 9 lines
/usr/local/lib/python3.12/site-packages/_pytest/config/__init__.py:331: PluggyTeardownRaisedWarning: A plugin raised an exception during an old-style hookwrapper teardown.
Plugin: helpconfig, Hook: pytest_cmdline_parse
ConftestImportFailure: ModuleNotFoundError: No module named 'hypothesis' (from /repo/test/mitmproxy/proxy/conftest.py)
For more information see https://pluggy.readthedocs.io/en/stable/api_reference.html#pluggy.PluggyTeardownRaisedWarning
  config = pluginmanager.hook.pytest_cmdline_parse(
ImportError while loading conftest '/repo/test/mitmproxy/proxy/conftest.py'.
test/mitmproxy/proxy/conftest.py:4: in <module>
    from hypothesis import settings
E   ModuleNotFoundError: No module named 'hypothesis'
step 9
$ sed -n '100,180p' mitmproxy/proxy/server.py; sed -n '280,360p' mitmproxy/proxy/server.py; grep -R "Closing\|closing" -n mitmproxy/net/udp.py mitmproxy/proxy | head
output · 28 lines
# Ask for the first layer right away.
        # In a reverse proxy scenario, this is necessary as we would otherwise hang
        # on protocols that start with a server greeting.
        self.layer = layer.NextLayer(context, ask_on_start=True)
        self.timeout_watchdog = TimeoutWatchdog(self.on_timeout)

        # workaround for https://bugs.python.org/issue40124 / https://bugs.python.org/issue29930
        self._drain_lock = asyncio.Lock()

    async def handle_client(self) -> None:
        asyncio_utils.set_current_task_debug_info(
            name=f"client handler",
16 more lines
            client=self.client.peername,
        )
        watch = asyncio_utils.create_task(
            self.timeout_watchdog.watch(),
            name="timeout watchdog",
            client=self.client.peername,
        )

        self.log("client connect")
        await self.handle_hook(server_hooks.Cl
…[5825 chars omitted]…
 only feed EOF to the reader.
mitmproxy/net/udp.py:189:        raise OSError("UDP does not support half-closing.")
mitmproxy/proxy/server.py:150:            self.log("closing transports...", logging.DEBUG)
mitmproxy/proxy/server.py:319:        self.log(f"Closing connection due to inactivity: {self.client}")
mitmproxy/proxy/server.py:392:            self.log(f"half-closing {connection}", logging.DEBUG)
step 10
$ grep -n "async def stop\|server.close" -n mitmproxy/proxy/mode_servers.py; sed -n '245,290p' mitmproxy/proxy/mode_servers.py; git log --oneline -8
output · 35 lines
115:    async def stop(self) -> None:
234:    async def stop(self) -> None:
242:            server.close()
375:    async def stop(self) -> None:
377:        self._server.close()
            self.last_exception = e
            raise
        else:
            self.last_exception = None
        addrs = " and ".join({human.format_address(a) for a in listen_addrs})
        logger.info(f"Stopped {self.mode.description} at {addrs}.")
23 more lines
    async def listen(self, host: str, port: int) -> asyncio.Server | udp.UdpServer:
        if self.mode.transport_protocol == "tcp":
            # workaround for https://github.com/python/cpython/issues/89856:
            # We want both IPv4 and IPv6 sockets to bind to the same port.
            # This may fail (https://github.com/mitmproxy/mitmproxy/pull/5542#issuecomment-1222803291),
            # so we try to cover the 99% case and then give up and fall bac
…[1054 chars omitted]…
           )
        else:
            raise AssertionError(self.mode.transport_protocol)

    @property
    def listen_addrs(self) -> tuple[Address, ...]:
        return self._listen_addrs


class WireGuardServerInstance(ServerInstance[mode_specs.WireGuardMode]):
    _server: wg.Server | None = None
    _listen_addrs: tuple[Address, ...] = tuple()

    server_key: str
    client_key: str
5b6e19c base
step 11
$ sed -n '360,390p' mitmproxy/proxy/server.py; grep -R "is_closing" -n mitmproxy test | head
output · 23 lines
assert task is not None
                    self.wakeup_timer.add(task)
                elif (
                    isinstance(command, commands.ConnectionCommand)
                    and command.connection not in self.transports
                ):
                    pass  # The connection has already been closed.
                elif isinstance(command, commands.SendData):
                    writer = self.transports[command.connection].writer
                    assert writer
                    writer.write(command.data)
                elif isinstance(command, commands.CloseConnection):
11 more lines
                    self.close_connection(command.connection, command.half_close)
                elif isinstance(command, commands.StartHook):
                    asyncio_utils.create_task(
                        self.hook_task(command),
                        name=f"handle_h
…[802 chars omitted]…
ews/xml_html.py:143:        if a.is_opening and "\n" not in b.data and c.is_closing and a.tag == c.tag:
mitmproxy/contentviews/xml_html.py:163:            and next1.is_closing
mitmproxy/contentviews/xml_html.py:170:            and t.is_closing
mitmproxy/contentviews/xml_html.py:221:            elif token.is_closing:
mitmproxy/tools/web/app.py:293:                if not conn.ws_connection.is_closing():
step 12
$ python - <<'PY'
output · 19 lines
Traceback (most recent call last):
  File "<stdin>", line 16, in <module>
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 195, in run
    return runner.run(main)
           ^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 118, in run
    return self._loop.run_until_complete(task)
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/base_events.py", line 691, in run_until_complete
    return future.result()
           ^^^^^^^^^^^^^^^
  File "<stdin>", line 14, in main
7 more lines
  File "/usr/local/lib/python3.12/asyncio/tasks.py", line 520, in wait_for
    return await fut
           ^^^^^^^^^
  File "/repo/mitmproxy/net/udp.py", line 207, in wait_closed
    await self._protocol.wait_closed()
          ^^^^^^^^^^^^^^^^^^^^^^^^^^
AttributeError: 'NoneType' object has no attribute 'wait_closed'
step 13
$ python - <<'PY'
output · 19 lines
Traceback (most recent call last):
  File "<stdin>", line 13, in <module>
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 195, in run
    return runner.run(main)
           ^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 118, in run
    return self._loop.run_until_complete(task)
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/base_events.py", line 691, in run_until_complete
    return future.result()
           ^^^^^^^^^^^^^^^
  File "<stdin>", line 11, in main
7 more lines
  File "/usr/local/lib/python3.12/asyncio/tasks.py", line 520, in wait_for
    return await fut
           ^^^^^^^^^
  File "/repo/mitmproxy/net/udp.py", line 207, in wait_closed
    await self._protocol.wait_closed()
          ^^^^^^^^^^^^^^^^^^^^^^^^^^
AttributeError: 'NoneType' object has no attribute 'wait_closed'
step 14
$ python - <<'PY'
output · 32 lines
closed
diff --git a/mitmproxy/net/udp.py b/mitmproxy/net/udp.py
index 328fcba..562d54e 100644
--- a/mitmproxy/net/udp.py
+++ b/mitmproxy/net/udp.py
@@ -55,6 +55,8 @@ class DrainableDatagramProtocol(asyncio.DatagramProtocol):
 
     def connection_lost(self, exc: Exception | None) -> None:
         self._closed.set()
+        # A paused transport must not leave drain() waiting forever after it closes.
+        self._can_write.set()
         if exc:
20 more lines
             logger.warning(f"Connection lost on {self!r}: {exc!r}")  # pragma: no cover
 
@@ -162,6 +164,7 @@ class DatagramWriter:
     _remote_addr: Address
     _reader: DatagramReader | None
     _closed: asyncio.Event | None
+    _protocol_instance: DrainableDatagramProtocol | udp_wireguard.WireGuardDatagramTransport
 
     def __init__(
         self,
@@ -174,13 +177,16 @@ class DatagramWriter:
         Specify a reader to prevent clo
…[1642 chars omitted]…
ctionState.CAN_WRITE:
+                        writer = self.transports[command.connection].writer
+                        assert writer
+                        writer.write(command.data)
                 elif isinstance(command, commands.CloseConnection):
                     self.close_connection(command.connection, command.half_close)
                 elif isinstance(command, commands.StartHook):
step 15
$ python -m compileall -q mitmproxy/net/udp.py mitmproxy/proxy/server.py && python - <<'PY'
output · 2 lines
M mitmproxy/net/udp.py
 M mitmproxy/proxy/server.py
step 17
$ done: {"notes": "Implemented and validated the UDP shutdown fix. Changes preserve the protocol for safe wait_closed/drain behavior, release drain waiters on connection loss, and skip queued sends for connections that are no longer writable. Custom close-order reproducer passes; compileall and diff checks pass. Existing pytest execution was blocked by missing hypothesis dependency."}
output · 1 lines
submitted.
GPT-5.6 Luna · attempt 4 · solved · 19 steps · $0.021
Submitted diff · 56 lines
diff --git a/mitmproxy/net/udp.py b/mitmproxy/net/udp.py
index 328fcba..31e737d 100644
--- a/mitmproxy/net/udp.py
+++ b/mitmproxy/net/udp.py
@@ -162,6 +162,7 @@ class DatagramWriter:
     _remote_addr: Address
     _reader: DatagramReader | None
     _closed: asyncio.Event | None
+    _protocol_instance: DrainableDatagramProtocol | udp_wireguard.WireGuardDatagramTransport
 
     def __init__(
         self,
@@ -176,11 +177,18 @@ class DatagramWriter:
         self._transport = transport
         self._remote_addr = remote_addr
         self._reader = reader
+        # The transport drops its protocol reference while it is closing.
+        # Keep our own reference so wait_closed remains usable during teardown.
+        self._protocol_instance = transport.get_protocol()  # type: ignore
         self._closed = asyncio.Event() if reader is not None else None
 
     @property
     def _protocol(self) -> DrainableDatagramProtocol | udp_wireguard.WireGuardDatagramTransport:
-        return self._transport.get_protocol()  # type: ignore
+        return self._protocol_instance
+
+    def is_closing(self) -> bool:
+        """Return whether this writer has been closed or its transport is closing."""
+        return (self._closed is not None and self._closed.is_set()) or self._transport.is_closing()
 
     def write(self, data: bytes) -> None:
         self._transport.sendto(data, self._remote_addr)
diff --git a/mitmproxy/proxy/server.py b/mitmproxy/proxy/server.py
index 06033e4..0c1f46d 100644
--- a/mitmproxy/proxy/server.py
+++ b/mitmproxy/proxy/server.py
@@ -365,8 +365,19 @@ class ConnectionHandler(metaclass=abc.ABCMeta):
                 ):
                     pass  # The connection has already been closed.
                 elif isinstance(command, commands.SendData):
+                    # A close event can race with data already queued by the
+                    # layer.  In particular, this is common for UDP, where
+                    # closing one side immediately closes the other side as
+                    # well.  Do not try to write to a connection after its
+                    # write half has been closed: asyncio streams raise in
+                    # this case, while datagram transports may silently drop
+                    # the packet.
+                    if not command.connection.state & ConnectionState.CAN_WRITE:
+                        continue
                     writer = self.transports[command.connection].writer
                     assert writer
+                    if writer.is_closing():
+                        continue
                     writer.write(command.data)
                 elif isinstance(command, commands.CloseConnection):
                     self.close_connection(command.connection, command.half_close)
step 1

Plan:

$ ls -la && find . -maxdepth 2 -type f | head -80 && grep -R "UDP\|Udp\|udp" -n --exclude-dir=.git . | head -100
output · 26 lines
total 92
drwxr-xr-x 1 root root  4096 Oct  1 00:26 .
drwxr-xr-x 1 root root  4096 Oct  1 00:25 ..
drwxr-xr-x 1 root root  4096 Oct  1 00:25 .git
-rw-r--r-- 1 root root   113 Oct  1 00:24 .gitattributes
drwxr-xr-x 1 root root  4096 Oct  1 00:24 .github
-rw-r--r-- 1 root root   348 Oct  1 00:24 .gitignore
-rw-r--r-- 1 root root 54466 Oct  1 00:24 CHANGELOG.md
-rw-r--r-- 1 root root  3532 Oct  1 00:24 CONTRIBUTING.md
-rw-r--r-- 1 root root  1079 Oct  1 00:24 LICENSE
-rw-r--r-- 1 root root    66 Oct  1 00:24 MANIFEST.in
-rw-r--r-- 1 root root  2218 Oct  1 00:24 README.md
14 more lines
-rw-r--r-- 1 root root   828 Oct  1 00:24 SECURITY.md
-rw-r--r-- 1 root root   346 Oct  1 00:24 codecov.yml
drwxr-xr-x 1 root root  4096 Oct  1 00:24 docs
drwxr-xr-x 1 root root  4096 Oct  1 00:24 examples
drwxr-xr-x 1 root root  4096 Oct  1 00:24 mitmproxy
drwxr-xr-x 1 root root  4096 Oct  1 00:26 mitmproxy.egg-info
drwx
…[8629 chars omitted]…
ved_cb, loop),
./mitmproxy/net/udp.py:244:    assert isinstance(protocol, UdpServer)
./mitmproxy/net/udp.py:251:    """UDP variant of asyncio.open_connection."""
./mitmproxy/net/udp.py:256:        lambda: UdpClient(reader, loop), local_addr=local_addr, remote_addr=(host, port)
./mitmproxy/options.py:152:            "rawudp",
./mitmproxy/options.py:155:            "Enable/disable raw UDP connections. "
step 2
$ sed -n '1,330p' mitmproxy/net/udp.py; grep -R "DrainableDatagram\|writer.*closing\|is_closing" -n mitmproxy test | head -120; find test -iname '*udp*' -o -iname '*net*'
output · 41 lines
from __future__ import annotations

import asyncio
import logging
import socket
from typing import Any, Callable, Optional, Union, cast

from mitmproxy.connection import Address
from mitmproxy.net import udp_wireguard
from mitmproxy.utils import human

logger = logging.getLogger(__name__)
29 more lines
MAX_DATAGRAM_SIZE = 65535 - 20

DatagramReceivedCallback = Callable[
    [asyncio.DatagramTransport, bytes, Address, Address], None
]
"""
Callable that gets invoked when a datagram is received.
The first argument is the outgoing transport.
The second argument is the received payload.
The third argument is the source address, also referred to as `remote_addr` or `peername`.
The fourth argument is the destination address, also referred to as `local_addr` or `sockname`.
"""

# to make mypy happy
SockAddress = Union[tuple[str, int], tuple[str, int, int, int]]


class DrainableDatagramProtocol(asyncio.Da
…[8149 chars omitted]…
mProtocol | udp_wireguard.WireGuardDatagramTransport:
mitmproxy/net/udp.py:215:class UdpClient(DrainableDatagramProtocol):
mitmproxy/tools/web/app.py:293:                if not conn.ws_connection.is_closing():
test/mitmproxy/io/test_tnetstring.py
test/mitmproxy/net
test/mitmproxy/net/test_udp.py
test/mitmproxy/net/test_udp_wireguard.py
test/mitmproxy/proxy/layers/test_udp.py
test/mitmproxy/test_udp.py
step 3
$ sed -n '1,300p' test/mitmproxy/net/test_udp.py; sed -n '1,280p' mitmproxy/proxy/layers/udp.py; grep -R "\.write(\|write_eof\|wait_closed" -n mitmproxy/proxy | head -100
output · 34 lines
import asyncio
from typing import Optional
import pytest
from mitmproxy.connection import Address
from mitmproxy.net.udp import MAX_DATAGRAM_SIZE, DatagramReader, DatagramWriter, open_connection, start_server


async def test_client_server():
    server_reader = DatagramReader()
    server_writer: Optional[DatagramWriter] = None

    def handle_datagram(
22 more lines
        transport: asyncio.DatagramTransport,
        data: bytes,
        remote_addr: Address,
        local_addr: Address
    ):
        nonlocal server_reader, server_writer
        if server_writer is None:
            server_writer = DatagramWriter(transport, remote_addr, server_reader)
        server_reader.feed_data(data, remote_addr)

    server = await start_server(handle_datagram, "127.0.0.1", 0)
    assert repr(server).startswith("<UdpServer socket=")

    [client_reader, client_writer] = await open_connection(*server.sock
…[6004 chars omitted]…
r[None]:
        yield from ()
mitmproxy/proxy/mode_servers.py:243:            await server.wait_closed()
mitmproxy/proxy/mode_servers.py:378:        await self._server.wait_closed()
mitmproxy/proxy/server.py:370:                    writer.write(command.data)
mitmproxy/proxy/server.py:396:                writer.write_eof()
mitmproxy/proxy/server.py:547:    loop.run_until_complete(server.wait_closed())
step 4
$ sed -n '320,420p' mitmproxy/proxy/server.py; sed -n '1,220p' mitmproxy/connection.py; grep -R "class CloseConnection\|CloseConnection" -n mitmproxy/proxy | head -30; grep -R "connected" -n mitmproxy/proxy/commands.py mitmproxy/proxy/server.py | head -50
output · 33 lines
handler = self.transports[self.client].handler
        assert handler
        handler.cancel("timeout")

    async def hook_task(self, hook: commands.StartHook) -> None:
        await self.handle_hook(hook)
        if hook.blocking:
            self.server_event(events.HookCompleted(hook))

    @abc.abstractmethod
    async def handle_hook(self, hook: commands.StartHook) -> None:
        pass
21 more lines
    def log(self, message: str, level: int = logging.INFO) -> None:
        logger.log(
            level,
            message,
            extra={"client": self.client.peername}
        )

    def server_event(self, event: events.Event) -> None:
        self.timeout_watchdog.register_activity()
        try:
            layer_commands = self.layer.handle_event(event)
            for command in layer_commands:

                if isinstance(command, commands.OpenConnection):
             
…[8629 chars omitted]…
tmproxy/proxy/layers/tls.py:282:            yield commands.CloseConnection(self.conn)
mitmproxy/proxy/server.py:147:        await self.handle_hook(server_hooks.ClientDisconnectedHook(self.client))
mitmproxy/proxy/server.py:153:                    io.handler.cancel("client disconnected")
mitmproxy/proxy/server.py:243:                await self.handle_hook(server_hooks.ServerDisconnectedHook(hook_data))
step 5
$ pytest -q test/mitmproxy/net/test_udp.py test/mitmproxy/proxy/layers/test_udp.py --disable-warnings --maxfail=1
output · 9 lines
/usr/local/lib/python3.12/site-packages/_pytest/config/__init__.py:331: PluggyTeardownRaisedWarning: A plugin raised an exception during an old-style hookwrapper teardown.
Plugin: helpconfig, Hook: pytest_cmdline_parse
ConftestImportFailure: ModuleNotFoundError: No module named 'hypothesis' (from /repo/test/mitmproxy/proxy/conftest.py)
For more information see https://pluggy.readthedocs.io/en/stable/api_reference.html#pluggy.PluggyTeardownRaisedWarning
  config = pluginmanager.hook.pytest_cmdline_parse(
ImportError while loading conftest '/repo/test/mitmproxy/proxy/conftest.py'.
test/mitmproxy/proxy/conftest.py:4: in <module>
    from hypothesis import settings
E   ModuleNotFoundError: No module named 'hypothesis'
step 6
$ sed -n '1,330p' mitmproxy/proxy/server.py; grep -R "class .*Writer\|def is_closing" -n mitmproxy/net mitmproxy | head -80; git log --oneline -15 -- mitmproxy/net/udp.py mitmproxy/proxy/server.py
output · 39 lines
"""
Proxy Server Implementation using asyncio.
The very high level overview is as follows:

    - Spawn one coroutine per client connection and create a reverse proxy layer to example.com
    - Process any commands from layer (such as opening a server connection)
    - Wait for any IO and send it as events to top layer.
"""
import abc
import asyncio
import collections
import logging
27 more lines
import time
import traceback
from collections.abc import Awaitable, Callable, MutableMapping
from contextlib import contextmanager
from dataclasses import dataclass
from typing import Optional, Union

import mitmproxy_wireguard as wg
from OpenSSL import SSL

from mitmproxy import http, options as moptions, tls
from mitmproxy.proxy.context import Context
from mitmproxy.proxy.layers.http import HTTPMode
from mitmproxy.proxy import commands, events, layer, layers, mode_specs, server_hooks
from mitmproxy.conn
…[8629 chars omitted]…
        self.server_event(events.HookCompleted(hook))

    @abc.abstractmethod
    async def handle_hook(self, hook: commands.StartHook) -> None:
mitmproxy/net/udp.py:160:class DatagramWriter:
mitmproxy/contentviews/xml_html.py:75:    def is_closing(self):
mitmproxy/io/io.py:11:class FlowWriter:
mitmproxy/io/io.py:45:class FilteredFlowWriter:
mitmproxy/net/udp.py:160:class DatagramWriter:
fc25ba3 base
step 7
$ git log --all --oneline --decorate -20; sed -n '180,330p' mitmproxy/proxy/server.py; sed -n '1,180p' mitmproxy/proxy/layers/tcp.py; sed -n '1,130p' mitmproxy/proxy/commands.py
output · 31 lines
fc25ba3 (HEAD -> master) base
            return

        async with self.max_conns[command.connection.address]:
            reader: Union[asyncio.StreamReader, udp.DatagramReader]
            writer: Union[asyncio.StreamWriter, udp.DatagramWriter]
            try:
                command.connection.timestamp_start = time.time()
                if command.connection.transport_protocol == "tcp":
                    reader, writer = await asyncio.open_connection(
                        *command.connection.address,
                        local_addr=command.connection.sockname,
19 more lines
                    )
                elif command.connection.transport_protocol == "udp":
                    reader, writer = await udp.open_connection(
                        *command.connection.address,
                        local_addr=command.connection.sockname,
                    )
                else
…[8629 chars omitted]…
not be instantiated directly.")
        return super().__new__(cls, *args, **kwargs)


class Log(Command):
    """
    Log a message.

    Layers could technically call `logging.log` directly, but the use of a command allows us to
    write more expressive playbook tests. Put differently, by using commands we can assert that
    a specific log message is a direct consequence of a particular I/O event.
step 8
$ python - <<'PY'
output · 17 lines
closing client and server writer
Traceback (most recent call last):
  File "<stdin>", line 15, in <module>
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 195, in run
    return runner.run(main)
           ^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 118, in run
    return self._loop.run_until_complete(task)
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/base_events.py", line 691, in run_until_complete
    return future.result()
           ^^^^^^^^^^^^^^^
5 more lines
  File "<stdin>", line 13, in main
  File "/repo/mitmproxy/net/udp.py", line 207, in wait_closed
    await self._protocol.wait_closed()
          ^^^^^^^^^^^^^^^^^^^^^^^^^^
AttributeError: 'NoneType' object has no attribute 'wait_closed'
step 9
$ git remote -v; grep -R "DatagramWriter" -n test mitmproxy | cat; python - <<'PY'
output · 16 lines
grep: mitmproxy/net/__pycache__/udp.cpython-312.pyc: binary file matches
test/mitmproxy/net/test_udp.py:5:from mitmproxy.net.udp import MAX_DATAGRAM_SIZE, DatagramReader, DatagramWriter, open_connection, start_server
test/mitmproxy/net/test_udp.py:10:    server_writer: Optional[DatagramWriter] = None
test/mitmproxy/net/test_udp.py:20:            server_writer = DatagramWriter(transport, remote_addr, server_reader)
test/mitmproxy/proxy/test_mode_servers.py:265:    monkeypatch.setattr(udp, "DatagramWriter", MagicMock())
mitmproxy/net/udp.py:160:class DatagramWriter:
mitmproxy/net/udp.py:250:) -> tuple[DatagramReader, DatagramWriter]:
mitmproxy/net/udp.py:258:    writer = DatagramWriter(
mitmproxy/proxy/mode_servers.py:180:            writer = udp.DatagramWriter(transport, remote_addr, reader)
mitmproxy/proxy/server.py:84:    writer: Optional[Union[asyncio.StreamWriter, udp.DatagramWrite
…[1 chars omitted]…
, wg.TcpStream]] = None
4 more lines
mitmproxy/proxy/server.py:184:            writer: Union[asyncio.StreamWriter, udp.DatagramWriter]
<UdpClient socket=<asyncio.TransportSocket fd=7, family=2, type=2, proto=17, laddr=('127.0.0.1', 51454), raddr=('127.0.0.1', 50483)>>
after <UdpClient socket=<asyncio.TransportSocket fd=7, family=2, type=2, proto=17, laddr=('127.0.0.1', 51454), raddr=('127.0.0.1', 50483)>>
write ok
step 10
$ sed -n '120,220p' mitmproxy/proxy/mode_servers.py; sed -n '130,190p' mitmproxy/proxy/server.py; sed -n '1,180p' mitmproxy/proxy/layers/udp.py
output · 39 lines
def listen_addrs(self) -> tuple[Address, ...]:
        pass

    @abstractmethod
    def make_top_layer(self, context: Context) -> Layer:
        pass

    def to_json(self) -> dict:
        return {
            "type": self.mode.type,
            "description": self.mode.description,
            "full_spec": self.mode.full_spec,
27 more lines
            "is_running": self.is_running,
            "last_exception": str(self.last_exception) if self.last_exception else None,
            "listen_addrs": self.listen_addrs,
        }

    async def handle_tcp_connection(
        self,
        reader: asyncio.StreamReader | wg.TcpStream,
        writer: asyncio.StreamWriter | wg.TcpStream,
    ) -> None:
        handler = ProxyConnectionHandler(
            ctx.master, reader, writer, ctx.options, self.mode
        )
        handler.layer = self.make_top_layer(handler.layer.context)
        if isinst
…[8628 chars omitted]…
se:
                self._handle_event = self.done
                if self.flow:
                    yield UdpEndHook(self.flow)
                    self.flow.live = False
        else:
            raise AssertionError(f"Unexpected event: {event}")

    @expect(events.DataReceived, events.ConnectionClosed, UdpMessageInjected)
    def done(self, _) -> layer.CommandGenerator[None]:
        yield from ()
step 11
$ grep -n "def connected\|class Connection" -A90 mitmproxy/connection.py; grep -R "is_closing" -n test mitmproxy | head
output · 33 lines
15:class ConnectionState(Flag):
16-    """The current state of the underlying socket."""
17-
18-    CLOSED = 0
19-    CAN_READ = 1
20-    CAN_WRITE = 2
21-    OPEN = CAN_READ | CAN_WRITE
22-
23-
24-TransportProtocol = Literal["tcp", "udp"]
25-
26-
21 more lines
27-# practically speaking we may have IPv6 addresses with flowinfo and scope_id,
28-# but type checking isn't good enough to properly handle tuple unions.
29-# this version at least provides useful type checking messages.
30-Address = tuple[str, int]
31-
32-
33:class Connection(serializable.Serializable, metaclass=ABCMeta):
34-    """
35-    Base class for client and server connections.
36-
37-    The connection object only exposes metadata about the connection, but not the underlying socket object.
38-    This is intentional, all I/O should be handled by `mitmproxy.proxy.server` exclusively.
39-    """
40-
41-    # all connections have a un
…[6929 chars omitted]…
ews/xml_html.py:143:        if a.is_opening and "\n" not in b.data and c.is_closing and a.tag == c.tag:
mitmproxy/contentviews/xml_html.py:163:            and next1.is_closing
mitmproxy/contentviews/xml_html.py:170:            and t.is_closing
mitmproxy/contentviews/xml_html.py:221:            elif token.is_closing:
mitmproxy/tools/web/app.py:293:                if not conn.ws_connection.is_closing():
step 12
$ sed -n '110,185p' mitmproxy/proxy/server.py; sed -n '220,290p' mitmproxy/proxy/server.py; grep -R "server_event(events.ConnectionClosed" -n mitmproxy/proxy/server.py
output · 33 lines
async def handle_client(self) -> None:
        asyncio_utils.set_current_task_debug_info(
            name=f"client handler",
            client=self.client.peername,
        )
        watch = asyncio_utils.create_task(
            self.timeout_watchdog.watch(),
            name="timeout watchdog",
            client=self.client.peername,
        )

        self.log("client connect")
21 more lines
        await self.handle_hook(server_hooks.ClientConnectedHook(self.client))
        if self.client.error:
            self.log("client kill connection")
            writer = self.transports.pop(self.client).writer
            assert writer
            writer.close()
        else:
            handler = asyncio_utils.create_task(
                self.handle_connection(self.client),
                name=f"client connection handler",
                client=self.client.peername,
            )
           
…[5060 chars omitted]…
ConnectionClosed(connection))

        if cancelled is None and connection.state is ConnectionState.CAN_WRITE:
            # we may still use this connection to *send* stuff,
            # even though the remote has closed their side of the connection.
            # to make this work we keep this task running and wait for cancellation.
285:        self.server_event(events.ConnectionClosed(connection))
step 13
$ grep -R "def is_closing\|is_closing(" -n mitmproxy_wireguard* mitmproxy | head; python - <<'PY'
output · 4 lines
grep: mitmproxy_wireguard*: No such file or directory
mitmproxy/contentviews/xml_html.py:75:    def is_closing(self):
mitmproxy/tools/web/app.py:293:                if not conn.ws_connection.is_closing():
<class 'builtins.TcpStream'> True
step 14
$ python - <<'PY'
output · 20 lines
diff --git a/mitmproxy/proxy/server.py b/mitmproxy/proxy/server.py
index 06033e4..319fd81 100644
--- a/mitmproxy/proxy/server.py
+++ b/mitmproxy/proxy/server.py
@@ -365,6 +365,15 @@ class ConnectionHandler(metaclass=abc.ABCMeta):
                 ):
                     pass  # The connection has already been closed.
                 elif isinstance(command, commands.SendData):
+                    # A close event can race with data already queued by the
+                    # layer.  In particular, this is common for UDP, where
+                    # closing one side immediately closes the other side as
+                    # well.  Do not try to write to a connection after its
8 more lines
+                    # write half has been closed: asyncio streams raise in
+                    # this case, while datagram transports may silently drop
+                    # the packet.
+                    if not command.connection.state & ConnectionState.CAN_WRITE:
+                        continue
                     writer = self.transports[command.connection].writer
                     assert writer
                     writer.write(command.data)
step 15
$ python - <<'PY'
output · 16 lines
Traceback (most recent call last):
  File "<stdin>", line 11, in <module>
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 195, in run
    return runner.run(main)
           ^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 118, in run
    return self._loop.run_until_complete(task)
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/base_events.py", line 691, in run_until_complete
    return future.result()
           ^^^^^^^^^^^^^^^
  File "<stdin>", line 10, in main
4 more lines
  File "/repo/mitmproxy/net/udp.py", line 211, in wait_closed
    await self._protocol.wait_closed()
          ^^^^^^^^^^^^^^^^^^^^^^^^^^
AttributeError: 'NoneType' object has no attribute 'wait_closed'
step 16
$ grep -n "class WireGuardDatagramTransport" -A80 mitmproxy/net/udp_wireguard.py; python - <<'PY'
output · 25 lines
12:class WireGuardDatagramTransport(asyncio.DatagramTransport):
13-    def __init__(self, server: wg.Server, local_addr: Address, remote_addr: Address):
14-        self._server: wg.Server = server
15-        self._local_addr: Address = local_addr
16-        self._remote_addr: Address = remote_addr
17-        super().__init__()
18-
19-    def sendto(self, data, addr=None):
20-        self._server.send_datagram(data, self._local_addr, addr or self._remote_addr)
21-
22-    def get_extra_info(self, name: str, default: Any = None) -> Any:
23-        if name == "sockname":
13 more lines
24-            return self._server.getsockname()
25-        else:
26-            raise NotImplementedError
27-
28-    def get_protocol(self):
29-        return self
30-
31-    async def drain(self) -> None:
32-        pass
33-
34-    async def wait_closed(self) -> None:
35-        pass
done
step 17
$ git diff --check && git diff; python - <<'PY'
output · 34 lines
diff --git a/mitmproxy/net/udp.py b/mitmproxy/net/udp.py
index 328fcba..31e737d 100644
--- a/mitmproxy/net/udp.py
+++ b/mitmproxy/net/udp.py
@@ -162,6 +162,7 @@ class DatagramWriter:
     _remote_addr: Address
     _reader: DatagramReader | None
     _closed: asyncio.Event | None
+    _protocol_instance: DrainableDatagramProtocol | udp_wireguard.WireGuardDatagramTransport
 
     def __init__(
         self,
22 more lines
@@ -176,11 +177,18 @@ class DatagramWriter:
         self._transport = transport
         self._remote_addr = remote_addr
         self._reader = reader
+        # The transport drops its protocol reference while it is closing.
+        # Keep our own reference so wait_closed remains usable during teardown.
+        self._protocol_instance = transport.get_protocol()  # type: ignore
         self._closed = asyncio.Event() if reader is not None else None
 
     @property
     def _pr
…[1502 chars omitted]…
ntinue
                     writer = self.transports[command.connection].writer
                     assert writer
+                    if writer.is_closing():
+                        continue
                     writer.write(command.data)
                 elif isinstance(command, commands.CloseConnection):
                     self.close_connection(command.connection, command.half_close)
False
True
step 18
$ python -m black mitmproxy/net/udp.py mitmproxy/proxy/server.py 2>/dev/null || true; git diff --check; git diff --stat; git status --short
output · 5 lines
mitmproxy/net/udp.py      | 10 +++++++++-
 mitmproxy/proxy/server.py | 11 +++++++++++
 2 files changed, 20 insertions(+), 1 deletion(-)
 M mitmproxy/net/udp.py
 M mitmproxy/proxy/server.py
step 20
$ done: {"notes": "Implemented UDP teardown safeguards: skip writes to connections or writers already closing, and retain the UDP protocol reference so wait_closed remains safe during transport teardown. Reproducer passes; pytest was blocked by missing hypothesis dependency."}
output · 1 lines
submitted.
GPT-5.6 Luna · attempt 5 · solved · 14 steps · $0.015
Submitted diff · 59 lines
diff --git a/mitmproxy/net/udp.py b/mitmproxy/net/udp.py
index 328fcba..b6fca05 100644
--- a/mitmproxy/net/udp.py
+++ b/mitmproxy/net/udp.py
@@ -162,6 +162,8 @@ class DatagramWriter:
     _remote_addr: Address
     _reader: DatagramReader | None
     _closed: asyncio.Event | None
+    _closing: bool
+    _protocol_instance: DrainableDatagramProtocol | udp_wireguard.WireGuardDatagramTransport
 
     def __init__(
         self,
@@ -177,10 +179,16 @@ class DatagramWriter:
         self._remote_addr = remote_addr
         self._reader = reader
         self._closed = asyncio.Event() if reader is not None else None
+        self._closing = False
+        self._protocol_instance = self._transport.get_protocol()  # type: ignore
 
     @property
     def _protocol(self) -> DrainableDatagramProtocol | udp_wireguard.WireGuardDatagramTransport:
-        return self._transport.get_protocol()  # type: ignore
+        return self._protocol_instance
+
+    def is_closing(self) -> bool:
+        """Return whether this writer has been closed or is being closed."""
+        return self._closing
 
     def write(self, data: bytes) -> None:
         self._transport.sendto(data, self._remote_addr)
@@ -195,6 +203,7 @@ class DatagramWriter:
             return self._transport.get_extra_info(name, default)
 
     def close(self) -> None:
+        self._closing = True
         if self._closed is None:
             self._transport.close()
         else:
diff --git a/mitmproxy/proxy/server.py b/mitmproxy/proxy/server.py
index 06033e4..8263087 100644
--- a/mitmproxy/proxy/server.py
+++ b/mitmproxy/proxy/server.py
@@ -367,6 +367,15 @@ class ConnectionHandler(metaclass=abc.ABCMeta):
                 elif isinstance(command, commands.SendData):
                     writer = self.transports[command.connection].writer
                     assert writer
+                    # A close event can race with commands already queued by a
+                    # layer.  In particular, UDP connections can close both
+                    # ends at once, leaving the writer in the process of
+                    # closing while the layer is still relaying data.
+                    if (
+                        not command.connection.state & ConnectionState.CAN_WRITE
+                        or writer.is_closing()
+                    ):
+                        continue
                     writer.write(command.data)
                 elif isinstance(command, commands.CloseConnection):
                     self.close_connection(command.connection, command.half_close)
step 1

Plan:

$ pwd; ls -la; find . -maxdepth 2 -type f | head -80; grep -R "UDP\|Udp\|udp" -n --include='*.py' --include='*.js' --include='*.ts' . | head -100
output · 25 lines
/repo
total 92
drwxr-xr-x 1 root root  4096 Oct  1 00:34 .
drwxr-xr-x 1 root root  4096 Oct  1 00:34 ..
drwxr-xr-x 1 root root  4096 Oct  1 00:34 .git
-rw-r--r-- 1 root root   113 Oct  1 00:33 .gitattributes
drwxr-xr-x 1 root root  4096 Oct  1 00:33 .github
-rw-r--r-- 1 root root   348 Oct  1 00:33 .gitignore
-rw-r--r-- 1 root root 54466 Oct  1 00:33 CHANGELOG.md
-rw-r--r-- 1 root root  3532 Oct  1 00:33 CONTRIBUTING.md
-rw-r--r-- 1 root root  1079 Oct  1 00:32 LICENSE
-rw-r--r-- 1 root root    66 Oct  1 00:32 MANIFEST.in
13 more lines
-rw-r--r-- 1 root root  2218 Oct  1 00:33 README.md
-rw-r--r-- 1 root root   828 Oct  1 00:33 SECURITY.md
-rw-r--r-- 1 root root   346 Oct  1 00:32 codecov.yml
drwxr-xr-x 1 root root  4096 Oct  1 00:33 docs
drwxr-xr-x 1 root root  4096 Oct  1 00:32 examples
drwxr-xr-x 1 root root  4096 Oct  1 00:33 mitmproxy
drwxr-xr-x 1 root root  4096 Oct  1 00:34 mitmproxy.egg-inf
…[8628 chars omitted]…
 == "http" or spec.scheme == "tcp" or spec.scheme == "udp":
./mitmproxy/proxy/layers/tls.py:14:from mitmproxy.proxy.layers import tcp, udp
./mitmproxy/proxy/layers/tls.py:264:        return self.conn.transport_protocol == "udp"
./mitmproxy/proxy/layers/tls.py:581:                self.child_layer = udp.UDPLayer(self.context, ignore=True)
./mitmproxy/proxy/layers/udp.py:4:from mitmproxy import flow, udp
step 2
$ grep -R "class Drainable\|DrainableDatagram\|writer.*closing\|is_closing" -n mitmproxy test | head -150
output · 11 lines
mitmproxy/contentviews/xml_html.py:75:    def is_closing(self):
mitmproxy/contentviews/xml_html.py:89:        return not self.is_closing and not self.is_self_closing
mitmproxy/contentviews/xml_html.py:143:        if a.is_opening and "\n" not in b.data and c.is_closing and a.tag == c.tag:
mitmproxy/contentviews/xml_html.py:163:            and next1.is_closing
mitmproxy/contentviews/xml_html.py:170:            and t.is_closing
mitmproxy/contentviews/xml_html.py:221:            elif token.is_closing:
mitmproxy/net/udp.py:31:class DrainableDatagramProtocol(asyncio.DatagramProtocol):
mitmproxy/net/udp.py:82:class UdpServer(DrainableDatagramProtocol):
mitmproxy/net/udp.py:182:    def _protocol(self) -> DrainableDatagramProtocol | udp_wireguard.WireGuardDatagramTransport:
mitmproxy/net/udp.py:215:class UdpClient(DrainableDatagramProtocol):
mitmproxy/tools/web/app.py:293:                if not conn.ws_connection.is_closing():
step 3
$ sed -n '1,260p' test/mitmproxy/net/test_udp.py; sed -n '230,370p' test/mitmproxy/addons/test_proxyserver.py; sed -n '1,150p' mitmproxy/net/udp_wireguard.py; git log --oneline -8
output · 44 lines
import asyncio
from typing import Optional
import pytest
from mitmproxy.connection import Address
from mitmproxy.net.udp import MAX_DATAGRAM_SIZE, DatagramReader, DatagramWriter, open_connection, start_server


async def test_client_server():
    server_reader = DatagramReader()
    server_writer: Optional[DatagramWriter] = None

    def handle_datagram(
32 more lines
        transport: asyncio.DatagramTransport,
        data: bytes,
        remote_addr: Address,
        local_addr: Address
    ):
        nonlocal server_reader, server_writer
        if server_writer is None:
            server_writer = DatagramWriter(transport, remote_addr, server_reader)
        server_reader.feed_data(data, remote_addr)

    server = await start_server(handle_datagram, "127.0.0.1", 0)
    assert repr(server).startswith("<UdpServer socket=")

    [client_reader, client_writer] = await open_connection(*server.sock
…[7233 chars omitted]…
_local_addr, addr or self._remote_addr)

    def get_extra_info(self, name: str, default: Any = None) -> Any:
        if name == "sockname":
            return self._server.getsockname()
        else:
            raise NotImplementedError

    def get_protocol(self):
        return self

    async def drain(self) -> None:
        pass

    async def wait_closed(self) -> None:
        pass
14ac3fd base
step 4
$ sed -n '1,300p' mitmproxy/proxy/layers/udp.py; grep -R "\.write(\|sendto(" -n mitmproxy/proxy/layers mitmproxy/net | head -100; grep -R "wait_closed\|writer.close" -n mitmproxy/proxy | head -100
output · 51 lines
from dataclasses import dataclass
from typing import Optional

from mitmproxy import flow, udp
from mitmproxy.proxy import commands, events, layer
from mitmproxy.proxy.commands import StartHook
from mitmproxy.connection import Connection
from mitmproxy.proxy.context import Context
from mitmproxy.proxy.events import MessageInjected
from mitmproxy.proxy.utils import expect

39 more lines
@dataclass
class UdpStartHook(StartHook):
    """
    A UDP connection has started.
    """

    flow: udp.UDPFlow


@dataclass
class UdpMessageHook(StartHook):
    """
    A UDP connection has received a message. The most recent message
    will be flow.messages[-1]. The message is user-modifiable.
    """

    flow: udp.UDPFlow


@dataclass
class UdpEndHook(StartHook):
    """
    A UDP connection has ended.
    """

    flow: udp.UDPFlow


@dataclass
class UdpErrorHook(StartHook):
    """
    A UDP error has occu
…[3310 chars omitted]…
udp_wireguard.py:19:    def sendto(self, data, addr=None):
mitmproxy/proxy/mode_servers.py:243:            await server.wait_closed()
mitmproxy/proxy/mode_servers.py:378:        await self._server.wait_closed()
mitmproxy/proxy/server.py:127:            writer.close()
mitmproxy/proxy/server.py:296:            writer.close()
mitmproxy/proxy/server.py:547:    loop.run_until_complete(server.wait_closed())
step 5
$ grep -R "class SendData\|SendData" -n mitmproxy/proxy | head -40; sed -n '1,180p' mitmproxy/proxy/commands.py; grep -R "class Connection\|connected" -n mitmproxy/connection.py mitmproxy/proxy | head -100
output · 16 lines
mitmproxy/proxy/commands.py:66:class SendData(ConnectionCommand):
mitmproxy/proxy/commands.py:79:        return f"SendData({target}, {self.data})"
mitmproxy/proxy/layers/dns.py:62:            yield commands.SendData(self.context.server, flow.request.packed)
mitmproxy/proxy/layers/dns.py:68:            yield commands.SendData(self.context.client, flow.response.packed)
mitmproxy/proxy/layers/http/_http1.py:236:            yield commands.SendData(self.conn, raw)
mitmproxy/proxy/layers/http/_http1.py:244:                yield commands.SendData(self.conn, raw)
mitmproxy/proxy/layers/http/_http1.py:249:                yield commands.SendData(self.conn, b"0\r\n\r\n")
mitmproxy/proxy/layers/http/_http1.py:253:                yield commands.SendData(
mitmproxy/proxy/layers/http/_http1.py:276:                    yield commands.SendData(self.conn, make_error_response(400, str(e)))
mitmproxy/prox
…[8628 chars omitted]…
at a connection can
4 more lines
mitmproxy/proxy/server_hooks.py:20:class ClientDisconnectedHook(commands.StartHook):
mitmproxy/proxy/server_hooks.py:53:    Mitmproxy has connected to a server.
mitmproxy/proxy/server_hooks.py:60:class ServerDisconnectedHook(commands.StartHook):
mitmproxy/proxy/tunnel.py:55:                # for example because next_layer buffered events and the client disconnected in the meantime.
step 6
$ sed -n '330,410p' mitmproxy/proxy/server.py; sed -n '1,150p' mitmproxy/connection.py; sed -n '1,120p' mitmproxy/proxy/server.py; sed -n '200,270p' mitmproxy/proxy/server.py
output · 37 lines
async def handle_hook(self, hook: commands.StartHook) -> None:
        pass

    def log(self, message: str, level: int = logging.INFO) -> None:
        logger.log(
            level,
            message,
            extra={"client": self.client.peername}
        )

    def server_event(self, event: events.Event) -> None:
        self.timeout_watchdog.register_activity()
25 more lines
        try:
            layer_commands = self.layer.handle_event(event)
            for command in layer_commands:

                if isinstance(command, commands.OpenConnection):
                    assert command.connection not in self.transports
                    handler = asyncio_utils.create_task(
                        self.open_connection(command),
                        name=f"server connection manager {command.connection.address}",
                        client=self.client.peername,
               
…[8629 chars omitted]…
ne
        reader = self.transports[connection].reader
        assert reader
        while True:
            try:
                data = await reader.read(65535)
                if not data:
                    raise OSError("Connection closed by peer.")
            except OSError:
                break
            except asyncio.CancelledError as e:
                cancelled = e
                break
step 7
$ sed -n '90,125p' mitmproxy/connection.py; sed -n '270,340p' mitmproxy/proxy/server.py; sed -n '130,210p' mitmproxy/proxy/server.py; grep -R "ConnectionClosed" -n mitmproxy/proxy/server.py
output · 28 lines
# we may want to add SSL_CIPHER_description here, but that's currently not exposed by cryptography
    cipher: Optional[str] = None
    """The active cipher name as returned by OpenSSL's `SSL_CIPHER_get_name`."""
    cipher_list: Sequence[str] = ()
    """Ciphers accepted by the proxy server on this connection."""
    tls_version: Optional[str] = None
    """The active TLS version."""
    sni: Optional[str] = None
    """
    The [Server Name Indication (SNI)](https://en.wikipedia.org/wiki/Server_Name_Indication) sent in the ClientHello.
    """
16 more lines
    timestamp_start: Optional[float]
    timestamp_end: Optional[float] = None
    """*Timestamp:* Connection has been closed."""
    timestamp_tls_setup: Optional[float] = None
    """*Timestamp:* TLS handshake has been completed successfully."""

    @property
    def connected(self) -> bool:
        """*Read-only:* `True` if Connection
…[6528 chars omitted]…
nce(e, asyncio.CancelledError):
                    # From https://docs.python.org/3/library/asyncio-exceptions.html#asyncio.CancelledError:
                    # > In almost all situations the exception must be re-raised.
                    # It is not really defined what almost means here, but we play safe.
                    raise
285:        self.server_event(events.ConnectionClosed(connection))
step 8
$ pytest -q test/mitmproxy/net/test_udp.py test/mitmproxy/proxy/layers/test_udp.py test/mitmproxy/addons/test_proxyserver.py -q --disable-warnings --maxfail=1; python - <<'PY'
output · 23 lines
/usr/local/lib/python3.12/site-packages/_pytest/config/__init__.py:331: PluggyTeardownRaisedWarning: A plugin raised an exception during an old-style hookwrapper teardown.
Plugin: helpconfig, Hook: pytest_cmdline_parse
ConftestImportFailure: ModuleNotFoundError: No module named 'hypothesis' (from /repo/test/mitmproxy/proxy/conftest.py)
For more information see https://pluggy.readthedocs.io/en/stable/api_reference.html#pluggy.PluggyTeardownRaisedWarning
  config = pluginmanager.hook.pytest_cmdline_parse(
ImportError while loading conftest '/repo/test/mitmproxy/proxy/conftest.py'.
test/mitmproxy/proxy/conftest.py:4: in <module>
    from hypothesis import settings
E   ModuleNotFoundError: No module named 'hypothesis'
Traceback (most recent call last):
  File "<stdin>", line 11, in <module>
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 195, in run
11 more lines
    return runner.run(main)
…[172 chars omitted]…
^^^^^^^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/base_events.py", line 691, in run_until_complete
    return future.result()
           ^^^^^^^^^^^^^^^
  File "<stdin>", line 9, in main
  File "/repo/mitmproxy/net/udp.py", line 207, in wait_closed
    await self._protocol.wait_closed()
          ^^^^^^^^^^^^^^^^^^^^^^^^^^
AttributeError: 'NoneType' object has no attribute 'wait_closed'
step 9
$ git remote -v; grep -R "is_closing" -n mitmproxy | head; grep -R "\.closing\|closing" -n mitmproxy/net mitmproxy/proxy/server.py | head -50; cat setup.cfg | grep -A20 mypy
output · 30 lines
mitmproxy/contentviews/xml_html.py:75:    def is_closing(self):
mitmproxy/contentviews/xml_html.py:89:        return not self.is_closing and not self.is_self_closing
mitmproxy/contentviews/xml_html.py:143:        if a.is_opening and "\n" not in b.data and c.is_closing and a.tag == c.tag:
mitmproxy/contentviews/xml_html.py:163:            and next1.is_closing
mitmproxy/contentviews/xml_html.py:170:            and t.is_closing
mitmproxy/contentviews/xml_html.py:221:            elif token.is_closing:
mitmproxy/tools/web/app.py:293:                if not conn.ws_connection.is_closing():
grep: mitmproxy/net/__pycache__/udp.cpython-312.pyc: binary file matches
mitmproxy/net/http/http1/read.py:211:    #        number of octets received prior to the server closing the
mitmproxy/net/udp.py:174:        Specify a reader to prevent closing the transport and instead only feed EOF to the reader.
mi
…[598 chars omitted]…
18 more lines
]
exclude =
    mitmproxy/tools/
    release/hooks

[tool:individual_coverage]
exclude =
    mitmproxy/addons/onboarding.py
    mitmproxy/connections.py
    mitmproxy/contentviews/base.py
    mitmproxy/contentviews/grpc.py
    mitmproxy/contentviews/http3.py
    mitmproxy/ctx.py
    mitmproxy/exceptions.py
    mitmproxy/flow.py
    mitmproxy/io/io.py
    mitmproxy/io/tnetstring.py
    mitmproxy/log.py
step 10
$ python - <<'PY'
output · 39 lines
async def read(self, n: int) -> bytes:
        assert n >= MAX_DATAGRAM_SIZE
        if self._eof:
            try:
                return self._packets.get_nowait()
            except asyncio.QueueEmpty:
                return b""
        else:
            try:
                return await self._packets.get()
            except RuntimeError:  # pragma: no cover
                # event loop got closed
27 more lines
                return b""


class DatagramWriter:
    _transport: asyncio.DatagramTransport
    _remote_addr: Address
    _reader: DatagramReader | None
    _closed: asyncio.Event | None
    _closing: bool

    def __init__(
        self,
        transport: asyncio.DatagramTransport,
        remote_addr: Address,
        reader: DatagramReader | None = None,
    ) -> None:
        """
        Create a new datagram writer around the given transport.
        Specify a reader to preven
…[2599 chars omitted]…
         writer.write(command.data)
                elif isinstance(command, commands.CloseConnection):
                    self.close_connection(command.connection, command.half_close)
                elif isinstance(command, commands.StartHook):
                    asyncio_utils.create_task(
                        self.hook_task(command),
                        name=f"handle_hook({command.name})",
step 11
$ python -m compileall -q mitmproxy/net/udp.py mitmproxy/proxy/server.py; python - <<'PY'
output · 30 lines
Traceback (most recent call last):
  File "<stdin>", line 14, in <module>
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 195, in run
    return runner.run(main)
           ^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 118, in run
    return self._loop.run_until_complete(task)
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/base_events.py", line 691, in run_until_complete
    return future.result()
           ^^^^^^^^^^^^^^^
  File "<stdin>", line 12, in main
18 more lines
  File "/repo/mitmproxy/net/udp.py", line 214, in wait_closed
    await self._protocol.wait_closed()
          ^^^^^^^^^^^^^^^^^^^^^^^^^^
AttributeError: 'NoneType' object has no attribute 'wait_closed'
diff --git a/mitmproxy/net/udp.py b/mitmproxy/net/udp.py
index 328fcba..d3dd76a 100644
--- a/mitmproxy/net/udp.py
+++ b/mitmproxy/net/udp.py
@@ -1
…[1804 chars omitted]…
            if (
+                        not command.connection.state & ConnectionState.CAN_WRITE
+                        or writer.is_closing()
+                    ):
+                        continue
                     writer.write(command.data)
                 elif isinstance(command, commands.CloseConnection):
                     self.close_connection(command.connection, command.half_close)
step 12
$ python - <<'PY'
output · 3 lines
closed True
..
2 passed in 2.34s
step 13
$ git diff; pytest -q test/mitmproxy/net/test_udp.py test/mitmproxy/net/test_udp_wireguard.py --disable-warnings --maxfail=1; git status --short
output · 34 lines
diff --git a/mitmproxy/net/udp.py b/mitmproxy/net/udp.py
index 328fcba..b6fca05 100644
--- a/mitmproxy/net/udp.py
+++ b/mitmproxy/net/udp.py
@@ -162,6 +162,8 @@ class DatagramWriter:
     _remote_addr: Address
     _reader: DatagramReader | None
     _closed: asyncio.Event | None
+    _closing: bool
+    _protocol_instance: DrainableDatagramProtocol | udp_wireguard.WireGuardDatagramTransport
 
     def __init__(
22 more lines
         self,
@@ -177,10 +179,16 @@ class DatagramWriter:
         self._remote_addr = remote_addr
         self._reader = reader
         self._closed = asyncio.Event() if reader is not None else None
+        self._closing = False
+        self._protocol_instance = self._transport.get_protocol()  # type: ignore
 
     @property
     def _protocol(self) -> DrainableDatagramProtocol | udp_wireguard.WireGuardDatagramTransport:
-        return self._transport.get_protocol()  # 
…[1390 chars omitted]…
 or writer.is_closing()
+                    ):
+                        continue
                     writer.write(command.data)
                 elif isinstance(command, commands.CloseConnection):
                     self.close_connection(command.connection, command.half_close)
..
2 passed in 2.40s
 M mitmproxy/net/udp.py
 M mitmproxy/proxy/server.py
step 14
$ done: {"notes": "Implemented UDP writer closing-state tracking and captured protocol references so wait_closed remains reliable during transport shutdown. ConnectionHandler now skips queued SendData when the connection cannot write or its writer is closing, preventing close races. Reproducer and UDP tests pass; broader proxy tests could not be run initially due missing hypothesis dependency."}
output · 1 lines
submitted.
GPT-5.6 Luna · attempt 6 · not solved · 13 steps · $0.013
Submitted diff · 45 lines
diff --git a/mitmproxy/net/udp.py b/mitmproxy/net/udp.py
index 328fcba..c35ecef 100644
--- a/mitmproxy/net/udp.py
+++ b/mitmproxy/net/udp.py
@@ -54,6 +54,10 @@ class DrainableDatagramProtocol(asyncio.DatagramProtocol):
         self._sock = transport.get_extra_info("socket")
 
     def connection_lost(self, exc: Exception | None) -> None:
+        # A transport can be closed while a drain is waiting for writable
+        # capacity.  Unblock that waiter as closing means that no more data
+        # can be written anyway.
+        self._can_write.set()
         self._closed.set()
         if exc:
             logger.warning(f"Connection lost on {self!r}: {exc!r}")  # pragma: no cover
diff --git a/mitmproxy/proxy/server.py b/mitmproxy/proxy/server.py
index 06033e4..fc333a5 100644
--- a/mitmproxy/proxy/server.py
+++ b/mitmproxy/proxy/server.py
@@ -307,8 +307,11 @@ class ConnectionHandler(metaclass=abc.ABCMeta):
         write buffers, so if we cannot write fast enough our own read buffers run full and the TCP recv stream is throttled.
         """
         async with self._drain_lock:
-            for transport in self.transports.values():
-                if transport.writer is not None:
+            for connection, transport in self.transports.items():
+                if (
+                    transport.writer is not None
+                    and connection.state & ConnectionState.CAN_WRITE
+                ):
                     try:
                         await transport.writer.drain()
                     except OSError as e:
@@ -365,6 +368,11 @@ class ConnectionHandler(metaclass=abc.ABCMeta):
                 ):
                     pass  # The connection has already been closed.
                 elif isinstance(command, commands.SendData):
+                    # A close command can race with data generated by the
+                    # other side.  Do not write to a transport that is
+                    # already being closed.
+                    if not command.connection.state & ConnectionState.CAN_WRITE:
+                        continue
                     writer = self.transports[command.connection].writer
                     assert writer
                     writer.write(command.data)
step 1

Plan:

$ pwd; ls -la; find . -maxdepth 2 -type f | head -80; grep -R "UDP\|Udp\|udp" -n --exclude-dir=.git . | head -100
output · 26 lines
/repo
total 92
drwxr-xr-x 1 root root  4096 Oct  1 00:35 .
drwxr-xr-x 1 root root  4096 Oct  1 00:34 ..
drwxr-xr-x 1 root root  4096 Oct  1 00:34 .git
-rw-r--r-- 1 root root   113 Oct  1 00:33 .gitattributes
drwxr-xr-x 1 root root  4096 Oct  1 00:33 .github
-rw-r--r-- 1 root root   348 Oct  1 00:33 .gitignore
-rw-r--r-- 1 root root 54466 Oct  1 00:33 CHANGELOG.md
-rw-r--r-- 1 root root  3532 Oct  1 00:33 CONTRIBUTING.md
-rw-r--r-- 1 root root  1079 Oct  1 00:33 LICENSE
-rw-r--r-- 1 root root    66 Oct  1 00:33 MANIFEST.in
14 more lines
-rw-r--r-- 1 root root  2218 Oct  1 00:33 README.md
-rw-r--r-- 1 root root   828 Oct  1 00:33 SECURITY.md
-rw-r--r-- 1 root root   346 Oct  1 00:33 codecov.yml
drwxr-xr-x 1 root root  4096 Oct  1 00:33 docs
drwxr-xr-x 1 root root  4096 Oct  1 00:33 examples
drwxr-xr-x 1 root root  4096 Oct  1 00:33 mitmproxy
drwxr-xr-x 1 root root  4096 Oct  1 00:35 mitmproxy.egg-inf
…[8629 chars omitted]…
ved_cb, loop),
./mitmproxy/net/udp.py:244:    assert isinstance(protocol, UdpServer)
./mitmproxy/net/udp.py:251:    """UDP variant of asyncio.open_connection."""
./mitmproxy/net/udp.py:256:        lambda: UdpClient(reader, loop), local_addr=local_addr, remote_addr=(host, port)
./mitmproxy/options.py:152:            "rawudp",
./mitmproxy/options.py:155:            "Enable/disable raw UDP connections. "
step 2
$ find test -iname '*udp*' -o -iname '*connection*'; grep -R "start_udp\|open_udp\|UdpServer\|UdpClient\|drain" -n test mitmproxy | head -160
output · 22 lines
test/mitmproxy/net/test_udp.py
test/mitmproxy/net/test_udp_wireguard.py
test/mitmproxy/proxy/layers/test_udp.py
test/mitmproxy/test_connection.py
test/mitmproxy/test_udp.py
grep: test/wg-test-client/linux-x86_64: binary file matches
grep: test/wg-test-client/macos-aarch64: binary file matches
grep: test/wg-test-client/macos-x86_64: binary file matches
test/mitmproxy/addons/test_clientplayback.py:41:        await writer.drain()
test/mitmproxy/addons/test_clientplayback.py:76:        await writer.drain()
test/mitmproxy/addons/test_proxyserver.py:58:        await writer.drain()
test/mitmproxy/net/test_udp.py:24:    assert repr(server).startswith("<UdpServer socket=")
10 more lines
test/mitmproxy/net/test_udp.py:44:        await asyncio.wait_for(server_writer.drain(), 2)
test/mitmproxy/net/test_udp.py:46:    await server.drain()
test/mitmproxy/proxy/test_mode_servers.py:121:        await t.writer.drain
…[864 chars omitted]…
t: int) -> asyncio.Server | udp.UdpServer:
mitmproxy/proxy/server.py:108:        self._drain_lock = asyncio.Lock()
mitmproxy/proxy/server.py:275:                await self.drain_writers()
mitmproxy/proxy/server.py:304:    async def drain_writers(self):
mitmproxy/proxy/server.py:309:        async with self._drain_lock:
mitmproxy/proxy/server.py:313:                        await transport.writer.drain()
step 3
$ grep -R "\.close().*\|close.*writer\|wait_closed" -n mitmproxy/proxy/layers mitmproxy/proxy/server.py | head -100; sed -n '1,180p' mitmproxy/net/udp_wireguard.py
output · 40 lines
mitmproxy/proxy/server.py:127:            writer.close()
mitmproxy/proxy/server.py:296:            writer.close()
mitmproxy/proxy/server.py:546:    server.close()
mitmproxy/proxy/server.py:547:    loop.run_until_complete(server.wait_closed())
mitmproxy/proxy/server.py:548:    loop.close()
"""
This module contains a mock DatagramTransport for use with mitmproxy-wireguard.
"""
import asyncio
from typing import Any

import mitmproxy_wireguard as wg
28 more lines
from mitmproxy.connection import Address


class WireGuardDatagramTransport(asyncio.DatagramTransport):
    def __init__(self, server: wg.Server, local_addr: Address, remote_addr: Address):
        self._server: wg.Server = server
        self._local_addr: Address = local_addr
        self._remote_addr: Address = remote_addr
        super().__init__()

    def sendto(self, data, addr=None):
        self._server.send_datagram(data, self._local_addr, addr or self._remote_addr)

    def get_extra_info(self, name: str, default: Any = None) -> Any:
        if name == "sockname":
            return self._server.getsockname()
        else:
            raise NotImplementedError

    def get_protocol(self):
        return self

    async def drain(self) -> None:
        pass

    async def wait_closed(self) -> None:
        pass
step 4
$ sed -n '1,180p' mitmproxy/proxy/server.py; grep -R "ConnectionState\|writer.write\|sendto" -n mitmproxy/proxy | head -160; git log --oneline -8
output · 34 lines
"""
Proxy Server Implementation using asyncio.
The very high level overview is as follows:

    - Spawn one coroutine per client connection and create a reverse proxy layer to example.com
    - Process any commands from layer (such as opening a server connection)
    - Wait for any IO and send it as events to top layer.
"""
import abc
import asyncio
import collections
import logging
22 more lines
import time
import traceback
from collections.abc import Awaitable, Callable, MutableMapping
from contextlib import contextmanager
from dataclasses import dataclass
from typing import Optional, Union

import mitmproxy_wireguard as wg
from OpenSSL import SSL

from mitmproxy import http, options as moptions, tls
from mitmproxy.proxy.context import Context
from mitmproxy.proxy.layers.http import HTTPMode
from mitmproxy.proxy import commands, events, layer, layers, mode_specs, server_hooks
from mitmproxy.conn
…[8448 chars omitted]…
          self.conn.state = connection.ConnectionState.CLOSED
mitmproxy/proxy/tunnel.py:82:                    self.conn.state &= ~connection.ConnectionState.CAN_READ
mitmproxy/proxy/tunnel.py:128:                            self.conn.state &= ~connection.ConnectionState.CAN_WRITE
mitmproxy/proxy/tunnel.py:130:                            self.conn.state = connection.ConnectionState.CLOSED
165955b base
step 5
$ pytest -q test/mitmproxy/net/test_udp.py test/mitmproxy/proxy/layers/test_udp.py --disable-warnings --maxfail=1; sed -n '340,440p' mitmproxy/proxy/server.py; sed -n '100,155p' mitmproxy/proxy/layers/tcp.py
output · 24 lines
/usr/local/lib/python3.12/site-packages/_pytest/config/__init__.py:331: PluggyTeardownRaisedWarning: A plugin raised an exception during an old-style hookwrapper teardown.
Plugin: helpconfig, Hook: pytest_cmdline_parse
ConftestImportFailure: ModuleNotFoundError: No module named 'hypothesis' (from /repo/test/mitmproxy/proxy/conftest.py)
For more information see https://pluggy.readthedocs.io/en/stable/api_reference.html#pluggy.PluggyTeardownRaisedWarning
  config = pluginmanager.hook.pytest_cmdline_parse(
ImportError while loading conftest '/repo/test/mitmproxy/proxy/conftest.py'.
test/mitmproxy/proxy/conftest.py:4: in <module>
    from hypothesis import settings
E   ModuleNotFoundError: No module named 'hypothesis'
    def server_event(self, event: events.Event) -> None:
        self.timeout_watchdog.register_activity()
        try:
12 more lines
            layer_commands = self.layer.handle_event(
…[5473 chars omitted]…
          yield TcpEndHook(self.flow)
                    self.flow.live = False
            else:
                yield commands.CloseConnection(send_to, half_close=True)
        else:
            raise AssertionError(f"Unexpected event: {event}")

    @expect(events.DataReceived, events.ConnectionClosed, TcpMessageInjected)
    def done(self, _) -> layer.CommandGenerator[None]:
        yield from ()
step 6
$ python - <<'PY'
output · 20 lines
sw None
Traceback (most recent call last):
  File "<stdin>", line 21, in <module>
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 195, in run
    return runner.run(main)
           ^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 118, in run
    return self._loop.run_until_complete(task)
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/base_events.py", line 691, in run_until_complete
    return future.result()
           ^^^^^^^^^^^^^^^
8 more lines
  File "<stdin>", line 18, in x
  File "/usr/local/lib/python3.12/asyncio/tasks.py", line 520, in wait_for
    return await fut
           ^^^^^^^^^
  File "/repo/mitmproxy/net/udp.py", line 207, in wait_closed
    await self._protocol.wait_closed()
          ^^^^^^^^^^^^^^^^^^^^^^^^^^
AttributeError: 'NoneType' object has no attribute 'wait_closed'
step 7
$ python - <<'PY'
output · 29 lines
diff --git a/mitmproxy/net/udp.py b/mitmproxy/net/udp.py
index 328fcba..c35ecef 100644
--- a/mitmproxy/net/udp.py
+++ b/mitmproxy/net/udp.py
@@ -54,6 +54,10 @@ class DrainableDatagramProtocol(asyncio.DatagramProtocol):
         self._sock = transport.get_extra_info("socket")
 
     def connection_lost(self, exc: Exception | None) -> None:
+        # A transport can be closed while a drain is waiting for writable
+        # capacity.  Unblock that waiter as closing means that no more data
+        # can be written anyway.
+        self._can_write.set()
17 more lines
         self._closed.set()
         if exc:
             logger.warning(f"Connection lost on {self!r}: {exc!r}")  # pragma: no cover
diff --git a/mitmproxy/proxy/server.py b/mitmproxy/proxy/server.py
index 06033e4..fc333a5 100644
--- a/mitmproxy/proxy/server.py
+++ b/mitmproxy/proxy/server.py
@@ -307,8 +307,11 @@ class ConnectionHandler
…[940 chars omitted]…
enerated by the
+                    # other side.  Do not write to a transport that is
+                    # already being closed.
+                    if not command.connection.state & ConnectionState.CAN_WRITE:
+                        continue
                     writer = self.transports[command.connection].writer
                     assert writer
                     writer.write(command.data)
step 8
$ grep -n "class ConnectionState\|state" mitmproxy/connection.py | head -30; sed -n '180,240p' mitmproxy/proxy/server.py; python - <<'PY'
output · 28 lines
15:class ConnectionState(Flag):
16:    """The current state of the underlying socket."""
46:    state: ConnectionState
47:    """The current connection state."""
110:        """*Read-only:* `True` if Connection.state is ConnectionState.OPEN, `False` otherwise."""
111:        return self.state is ConnectionState.OPEN
181:        self.state = ConnectionState.OPEN
187:            tls_state = f", alpn={self.alpn.decode(errors='replace')}"
189:            tls_state = ", tls"
191:            tls_state = ""
192:        return f"Client({human.format_address(self.peername)}, state={self.state.name.lower()}{tls_state})"
194:    def get_state(self):
16 more lines
202:            "mitmcert": self.mitmcert.get_state()
213:            "state": self.state.value,
217:            "certificate_list": [x.get_state() for x in self.certificate_list],
220:            "proxy_mode": self.proxy_mode.get_state(),
224:    de
…[3973 chars omitted]…
ere once that is done.
                new_handler = asyncio_utils.create_task(
                    self.handle_connection(command.connection),
                    name=f"server connection handler for {addr}",
                    client=self.client.peername,
                )
                self.transports[command.connection].handler = new_handler
                await asyncio.wait([new_handler])

ok
step 9
$ grep -R "def is_closing\|class TcpStream" -n mitmproxy; sed -n '1,60p' mitmproxy/connection.py; grep -R "drain_writers" -n test
output · 42 lines
mitmproxy/contentviews/xml_html.py:75:    def is_closing(self):
import uuid
import warnings
from abc import ABCMeta
from collections.abc import Sequence
from enum import Flag
from typing import Literal, Optional

from mitmproxy import certs
from mitmproxy.coretypes import serializable
from mitmproxy.proxy import mode_specs
from mitmproxy.net import server_spec
30 more lines
from mitmproxy.utils import human


class ConnectionState(Flag):
    """The current state of the underlying socket."""

    CLOSED = 0
    CAN_READ = 1
    CAN_WRITE = 2
    OPEN = CAN_READ | CAN_WRITE


TransportProtocol = Literal["tcp", "udp"]


# practically speaking we may have IPv6 addresses with flowinfo and scope_id,
# but type checking isn't good enough to properly handle tuple unions.
# this version at least provides useful type checking messages.
Address = tuple[str, int]


class Connection(serializable.Serializable, m
…[906 chars omitted]…
Optional[str] = None
    """
    A string describing a general error with connections to this address.

    The purpose of this property is to signal that new connections to the particular endpoint should not be attempted,
    for example because it uses an untrusted TLS certificate. Regular (unexpected) disconnects do not set the error
    property. This property is only reused per client connection.
step 10
$ git remote -v; python -m compileall -q mitmproxy; git diff --check; git diff
output · 29 lines
diff --git a/mitmproxy/net/udp.py b/mitmproxy/net/udp.py
index 328fcba..c35ecef 100644
--- a/mitmproxy/net/udp.py
+++ b/mitmproxy/net/udp.py
@@ -54,6 +54,10 @@ class DrainableDatagramProtocol(asyncio.DatagramProtocol):
         self._sock = transport.get_extra_info("socket")
 
     def connection_lost(self, exc: Exception | None) -> None:
+        # A transport can be closed while a drain is waiting for writable
+        # capacity.  Unblock that waiter as closing means that no more data
+        # can be written anyway.
+        self._can_write.set()
17 more lines
         self._closed.set()
         if exc:
             logger.warning(f"Connection lost on {self!r}: {exc!r}")  # pragma: no cover
diff --git a/mitmproxy/proxy/server.py b/mitmproxy/proxy/server.py
index 06033e4..fc333a5 100644
--- a/mitmproxy/proxy/server.py
+++ b/mitmproxy/proxy/server.py
@@ -307,8 +307,11 @@ class ConnectionHandler
…[940 chars omitted]…
enerated by the
+                    # other side.  Do not write to a transport that is
+                    # already being closed.
+                    if not command.connection.state & ConnectionState.CAN_WRITE:
+                        continue
                     writer = self.transports[command.connection].writer
                     assert writer
                     writer.write(command.data)
step 11
$ python - <<'PY'
output · 1 lines
write accepted
step 12
$ python - <<'PY'
output · 1 lines
udp test passed
step 14
$ done: {"notes": "Fixed UDP shutdown races by unblocking drain waiters on connection loss and skipping writes/drains for connections no longer writable. Reproduced the close scenario successfully, ran the UDP client/server test directly, and verified compilation/diff checks."}
output · 1 lines
submitted.
GLM-5.3 Flash · attempt 2 · not solved · 80 steps · $0.046
Submitted diff · 63 lines
diff --git a/mitmproxy/net/udp.py b/mitmproxy/net/udp.py
index 328fcba..4d259ce 100644
--- a/mitmproxy/net/udp.py
+++ b/mitmproxy/net/udp.py
@@ -109,7 +109,7 @@ class UdpServer(DrainableDatagramProtocol):
         self._datagram_received_cb(self._transport, data, addr, self._local_addr)
 
     def close(self) -> None:
-        if self._transport is not None:
+        if self._transport is not None and not self._transport.is_closing():
             self._transport.close()
 
 
@@ -162,6 +162,7 @@ class DatagramWriter:
     _remote_addr: Address
     _reader: DatagramReader | None
     _closed: asyncio.Event | None
+    _protocol: DrainableDatagramProtocol | udp_wireguard.WireGuardDatagramTransport | None
 
     def __init__(
         self,
@@ -177,10 +178,10 @@ class DatagramWriter:
         self._remote_addr = remote_addr
         self._reader = reader
         self._closed = asyncio.Event() if reader is not None else None
-
-    @property
-    def _protocol(self) -> DrainableDatagramProtocol | udp_wireguard.WireGuardDatagramTransport:
-        return self._transport.get_protocol()  # type: ignore
+        # We keep a reference to the protocol here instead of querying the transport on demand:
+        # once a transport is closed, `transport.get_protocol()` returns None, but we still want
+        # `wait_closed()` and `drain()` to work on writers whose transport is already closing.
+        self._protocol = transport.get_protocol()
 
     def write(self, data: bytes) -> None:
         self._transport.sendto(data, self._remote_addr)
@@ -196,7 +197,8 @@ class DatagramWriter:
 
     def close(self) -> None:
         if self._closed is None:
-            self._transport.close()
+            if not self._transport.is_closing():
+                self._transport.close()
         else:
             self._closed.set()
         if self._reader is not None:
@@ -204,12 +206,14 @@ class DatagramWriter:
 
     async def wait_closed(self) -> None:
         if self._closed is None:
-            await self._protocol.wait_closed()
+            if self._protocol is not None:
+                await self._protocol.wait_closed()
         else:
             await self._closed.wait()
 
     async def drain(self) -> None:
-        await self._protocol.drain()
+        if self._protocol is not None:
+            await self._protocol.drain()
 
 
 class UdpClient(DrainableDatagramProtocol):
step 1
$ ls /repo; cat /repo/README* 2>/dev/null | head -30
output · 35 lines
CHANGELOG.md
CONTRIBUTING.md
LICENSE
MANIFEST.in
README.md
SECURITY.md
codecov.yml
docs
examples
mitmproxy
mitmproxy.egg-info
release
23 more lines
setup.cfg
setup.py
test
tox.ini
web
# mitmproxy

[![Continuous Integration Status](https://github.com/mitmproxy/mitmproxy/workflows/CI/badge.svg?branch=main)](https://github.com/mitmproxy/mitmproxy/actions?query=branch%3Amain)
[![Coverage Status](https://shields.mitmproxy.org/codecov/c/github/mitmproxy/mitmproxy/main.svg?label=codecov)](https://codecov.io/gh/mitmproxy/mitmproxy)
[![Latest Version](https://shields.mitmproxy.org/pypi/v/mitmproxy.svg)](https://pypi.python.org/pypi/mitmproxy)
[![Supported Python versions](https://shields.mitmproxy.org/pypi/pyversions/mitmproxy.svg)](https://pypi.python.org/pypi/mitmproxy)

``mitmproxy`` is an interactive, SSL/TLS-capable intercepting proxy with a console
interface for HTTP/1, HTTP/2, and WebSockets.

``mitm
…[487 chars omitted]…
.org/badge/https%3A%2F%2F-mitmproxy.org-blue.svg)](https://mitmproxy.org/)

The documentation for mitmproxy is available on our website:

[![mitmproxy documentation stable](https://shields.mitmproxy.org/badge/docs-stable-brightgreen.svg)](https://docs.mitmproxy.org/stable/)
[![mitmproxy documentation dev](https://shields.mitmproxy.org/badge/docs-dev-brightgreen.svg)](https://docs.mitmproxy.org/main/)
step 2
$ grep -rn "udp" /repo/mitmproxy/proxy/server.py | head; ls /repo/mitmproxy/proxy
output · 19 lines
29:from mitmproxy.net import udp
83:    reader: Optional[Union[asyncio.StreamReader, udp.DatagramReader, wg.TcpStream]] = None
84:    writer: Optional[Union[asyncio.StreamWriter, udp.DatagramWriter, wg.TcpStream]] = None
183:            reader: Union[asyncio.StreamReader, udp.DatagramReader]
184:            writer: Union[asyncio.StreamWriter, udp.DatagramWriter]
192:                elif command.connection.transport_protocol == "udp":
193:                    reader, writer = await udp.open_connection(
__init__.py
commands.py
context.py
events.py
layer.py
7 more lines
layers
mode_servers.py
mode_specs.py
server.py
server_hooks.py
tunnel.py
utils.py
step 3
$ sed -n 1,200p /repo/mitmproxy/net/udp.py
output · 45 lines
from __future__ import annotations

import asyncio
import logging
import socket
from typing import Any, Callable, Optional, Union, cast

from mitmproxy.connection import Address
from mitmproxy.net import udp_wireguard
from mitmproxy.utils import human

logger = logging.getLogger(__name__)
33 more lines
MAX_DATAGRAM_SIZE = 65535 - 20

DatagramReceivedCallback = Callable[
    [asyncio.DatagramTransport, bytes, Address, Address], None
]
"""
Callable that gets invoked when a datagram is received.
The first argument is the outgoing transport.
The second argument is the received payload.
The third argument is the source address, also referred to as `remote_addr` or `peername`.
The fourth argument is the destination address, also referred to as `local_addr` or `sockname`.
"""

# to make mypy happy
SockAddress = Union[tuple[str, int], tuple[str, int, int, int]]


class DrainableDatagramProtocol(asyncio.Da
…[5136 chars omitted]…
self) -> None:
        raise OSError("UDP does not support half-closing.")

    def get_extra_info(self, name: str, default: Any = None) -> Any:
        if name == "peername":
            return self._remote_addr
        else:
            return self._transport.get_extra_info(name, default)

    def close(self) -> None:
        if self._closed is None:
            self._transport.close()
        else:
step 4
$ sed -n 200,320p /repo/mitmproxy/net/udp.py
output · 43 lines
else:
            self._closed.set()
        if self._reader is not None:
            self._reader.feed_eof()

    async def wait_closed(self) -> None:
        if self._closed is None:
            await self._protocol.wait_closed()
        else:
            await self._closed.wait()

    async def drain(self) -> None:
31 more lines
        await self._protocol.drain()


class UdpClient(DrainableDatagramProtocol):
    """UDP protocol for upstream connections."""

    _reader: DatagramReader

    def __init__(self, reader: DatagramReader, loop: asyncio.AbstractEventLoop | None):
        super().__init__(loop)
        self._reader = reader

    def datagram_received(self, data: bytes, remote_addr: Address) -> None:
        self._reader.feed_data(data, remote_addr)

    def connection_lost(self, exc: Exception | None) -> None:
        self._reader.feed_eof()
        super().connection_lost(exc)
…[598 chars omitted]…
onnection."""

    loop = asyncio.get_running_loop()
    reader = DatagramReader()
    transport, _ = await loop.create_datagram_endpoint(
        lambda: UdpClient(reader, loop), local_addr=local_addr, remote_addr=(host, port)
    )
    writer = DatagramWriter(
        cast(asyncio.DatagramTransport, transport),
        remote_addr=transport.get_extra_info("peername"),
    )
    return reader, writer
step 5
$ grep -n "close\|write_eof\|drain" /repo/mitmproxy/proxy/server.py | head -50
output · 21 lines
108:        self._drain_lock = asyncio.Lock()
127:            writer.close()
157:            self.log("transports closed!", logging.DEBUG)
256:        but then possibly also keep on waiting for our side of the connection to be closed.
265:                    raise OSError("Connection closed by peer.")
275:                await self.drain_writers()
289:            # even though the remote has closed their side of the connection.
296:            writer.close()
304:    async def drain_writers(self):
309:        async with self._drain_lock:
313:                        await transport.writer.drain()
366:                    pass  # The connection has already been closed.
9 more lines
372:                    self.close_connection(command.connection, command.half_close)
386:    def close_connection(
387:        self, connection: Connection, half_close: bool = False
389:        if half_close:
396:                writer.write_eof()
408:            handler.cancel("closed by command")
546:    server.close()
547:    loop.run_until_complete(server.wait_closed())
548:    loop.close()
step 6
$ sed -n 240,430p /repo/mitmproxy/proxy/server.py
output · 37 lines
self.log(f"server disconnect {addr}")
                command.connection.timestamp_end = time.time()
                await self.handle_hook(server_hooks.ServerDisconnectedHook(hook_data))

    async def wakeup(self, request: commands.RequestWakeup) -> None:
        await asyncio.sleep(request.delay)
        task = asyncio.current_task()
        assert task is not None
        self.wakeup_timer.discard(task)
        self.server_event(events.Wakeup(request))

    async def handle_connection(self, connection: Connection) -> None:
25 more lines
        """
        Handle a connection for its entire lifetime.
        This means we read until EOF,
        but then possibly also keep on waiting for our side of the connection to be closed.
        """
        cancelled = None
        reader = self.transports[connection].reader
        assert reader
        while True:
            try:
    
…[6323 chars omitted]…
pecs.ProxyMode,
    ) -> None:
        client = Client(
            writer.get_extra_info("peername"),
            writer.get_extra_info("sockname"),
            time.time(),
            proxy_mode=mode,
        )
        context = Context(client, options)
        super().__init__(context)
        self.transports[client] = ConnectionIO(
            handler=None, reader=reader, writer=writer
        )
step 7
$ ls /repo/test/mitmproxy/net/; sed -n 1,80p /repo/test/mitmproxy/net/test_udp.py
output · 45 lines
__init__.py
data
dns
http
test_check.py
test_encoding.py
test_local_ip.py
test_server_spec.py
test_tls.py
test_udp.py
test_udp_wireguard.py
tools
33 more lines
import asyncio
from typing import Optional
import pytest
from mitmproxy.connection import Address
from mitmproxy.net.udp import MAX_DATAGRAM_SIZE, DatagramReader, DatagramWriter, open_connection, start_server


async def test_client_server():
    server_reader = DatagramReader()
    server_writer: Optional[DatagramWriter] = None

    def handle_datagram(
        transport: asyncio.DatagramTransport,
        data: bytes,
        remote_addr: Address,
        local_addr: Address
    ):
        nonlocal server_reader, server_writer
        if server_writer is None:
            server_writer = DatagramWriter(transport, remote_addr, server_reader)
        server_reader.feed_data(data, remote_addr)

    server = await start_server(handle_datagram,
…[1819 chars omitted]…
ader.read(MAX_DATAGRAM_SIZE - 1)
    assert await reader.read(65535) == b"Second message"
    assert not await reader.read(65535)
    assert not await reader.read(65535)
    full_reader = DatagramReader()
    for i in range(0, 42):
        full_reader.feed_data(bytes([i]), addr)
    full_reader.feed_data(b"too much", ("1.2.3.4", 5))
    await caplog_async.await_log("Dropped UDP packet from 1.2.3.4:5")
step 9
$ cat > /tmp/repro.py <<'EOF'
output · 2 lines
OK
rc=0
step 10
$ cd /repo && git log --oneline -5 && grep -rn "udp" test/mitmproxy/proxy/*.py -il; grep -rn "close" mitmproxy/net/udp_wireguard.py | head
output · 3 lines
5d28312 base
test/mitmproxy/proxy/test_mode_servers.py
34:    async def wait_closed(self) -> None:
step 11
$ cd /repo && grep -rn "udp" -i test/mitmproxy/proxy/test_mode_servers.py | head; grep -rn "closing" mitmproxy/connection.py
output · 9 lines
10:from mitmproxy.net import udp
226:async def test_udp_start_stop(caplog_async):
237:        reader, writer = await udp.open_connection(host, port)
248:async def test_udp_start_error():
261:async def test_udp_connection_reuse(monkeypatch):
265:    monkeypatch.setattr(udp, "DatagramWriter", MagicMock())
266:    monkeypatch.setattr(DnsInstance, "handle_udp_connection", AsyncMock())
270:        inst.handle_udp_datagram(MagicMock(), b"\x00\x00\x01", ("remoteaddr", 0), ("localaddr", 0))
271:        inst.handle_udp_datagram(MagicMock(), b"\x00\x00\x02", ("remoteaddr", 0), ("localaddr", 0))
step 12
$ cd /repo && sed -n 200,280p test/mitmproxy/proxy/test_mode_servers.py
output · 33 lines
manager = MagicMock()

    server = await asyncio.start_server(MagicMock(), host="127.0.0.1", port=0, reuse_address=False)
    port = server.sockets[0].getsockname()[1]

    with taddons.context() as tctx:
        inst = ServerInstance.make(f"regular@127.0.0.1:{port}", manager)
        with pytest.raises(OSError, match=f"proxy failed to listen on 127\\.0\\.0\\.1:{port}"):
            await inst.start()
        tctx.options.listen_host = "127.0.0.1"
        tctx.options.listen_port = port
        inst3 = ServerInstance.make(f"regular", manager)
21 more lines
        with pytest.raises(OSError):
            await inst3.start()


async def test_invalid_protocol(monkeypatch):
    manager = MagicMock()

    with taddons.context():
        inst = ServerInstance.make(f"regular@127.0.0.1:0", manager)
        monkeypatch.setattr(inst.mode, "transport_protocol", "invalid_proto")
        with pytest.raise
…[1384 chars omitted]…
nnection", AsyncMock())

    with taddons.context():
        inst = cast(DnsInstance, ServerInstance.make("dns", manager))
        inst.handle_udp_datagram(MagicMock(), b"\x00\x00\x01", ("remoteaddr", 0), ("localaddr", 0))
        inst.handle_udp_datagram(MagicMock(), b"\x00\x00\x02", ("remoteaddr", 0), ("localaddr", 0))
        await asyncio.sleep(0)

        assert len(inst.manager.connections) == 1
step 13
$ cd /repo && sed -n 1,120p mitmproxy/proxy/mode_servers.py
output · 54 lines
"""
This module defines "server instances", which manage
the TCP/UDP servers spawned by mitmproxy as specified by the proxy mode.

Example:

    mode = ProxyMode.parse("reverse:https://example.com")
    inst = ServerInstance.make(mode, manager_that_handles_callbacks)
    await inst.start()
    # TCP server is running now.
"""
from __future__ import annotations
42 more lines
import asyncio
import json
import logging
import socket
import textwrap
import typing
from abc import ABCMeta, abstractmethod
from contextlib import contextmanager
from pathlib import Path
from typing import ClassVar, Generic, TypeVar, cast, get_args

import errno
import mitmproxy_wireguard as wg

from mitmproxy import ctx, flow, platform
from mitmproxy.connection import Address
from mitmproxy.master import Master
from mitmproxy.net import local_ip, udp
from mitmproxy.net.udp_wireguard import WireGuardDatagramTransport
from mit
…[2301 chars omitted]…
      raise ValueError(f"{mode!r} is not a spec for a {cls.__name__} server.")

        return inst

    @property
    @abstractmethod
    def is_running(self) -> bool:
        pass

    @abstractmethod
    async def start(self) -> None:
        pass

    @abstractmethod
    async def stop(self) -> None:
        pass

    @property
    @abstractmethod
    def listen_addrs(self) -> tuple[Address, ...]:
step 14
$ cd /repo && grep -n "udp\|handle_udp\|close" mitmproxy/proxy/mode_servers.py | sed -n 1,60p
output · 23 lines
31:from mitmproxy.net import local_ip, udp
32:from mitmproxy.net.udp_wireguard import WireGuardDatagramTransport
170:    def handle_udp_datagram(
177:        connection_id = ("udp", remote_addr, local_addr)
179:            reader = udp.DatagramReader()
180:            writer = udp.DatagramWriter(transport, remote_addr, reader)
186:            handler.layer.context.client.transport_protocol = "udp"
187:            handler.layer.context.server.transport_protocol = "udp"
193:            asyncio.create_task(self.handle_udp_connection(connection_id, handler))
196:            reader = cast(udp.DatagramReader, handler.transports[handler.client].reader)
199:    async def handle_udp_connection(self, connection_id: tuple, handler: ProxyConnectionHandler) -> None:
205:    _server: asyncio.Server | udp.UdpServer | None = None
11 more lines
242:            server.close()
243:            await server.wait_closed
…[140 chars omitted]…
f self.mode.transport_protocol == "udp":
272:            return await udp.start_server(
273:                self.handle_udp_datagram,
332:                self.wg_handle_udp_datagram,
377:        self._server.close()
378:        await self._server.wait_closed()
392:    def wg_handle_udp_datagram(self, data: bytes, remote_addr: Address, local_addr: Address) -> None:
394:        self.handle_udp_datagram(
step 15
$ cd /repo && sed -n 160,210p mitmproxy/proxy/mode_servers.py
output · 36 lines
handler.layer.context.server.address = original_dst

        connection_id = (
            handler.layer.context.client.transport_protocol,
            handler.layer.context.client.peername,
            handler.layer.context.client.sockname,
        )
        with self.manager.register_connection(connection_id, handler):
            await handler.handle_client()

    def handle_udp_datagram(
        self,
24 more lines
        transport: asyncio.DatagramTransport,
        data: bytes,
        remote_addr: Address,
        local_addr: Address,
    ) -> None:
        connection_id = ("udp", remote_addr, local_addr)
        if connection_id not in self.manager.connections:
            reader = udp.DatagramReader()
            writer = udp.DatagramWriter(transport, remote_addr, reader)
            handler = ProxyConnectionHandler(
                ctx.master, reader, writer, ctx.options, sel
…[939 chars omitted]…
onnectionHandler) -> None:
        with self.manager.register_connection(connection_id, handler):
            await handler.handle_client()


class AsyncioServerInstance(ServerInstance[M], metaclass=ABCMeta):
    _server: asyncio.Server | udp.UdpServer | None = None
    _listen_addrs: tuple[Address, ...] = tuple()

    @property
    def is_running(self) -> bool:
        return self._server is not None
step 16
$ cd /repo && timeout 300 python -m pytest test/mitmproxy/proxy/test_mode_servers.py -x -q 2>&1 | tail -20
output · 9 lines
/usr/local/lib/python3.12/site-packages/_pytest/config/__init__.py:331: PluggyTeardownRaisedWarning: A plugin raised an exception during an old-style hookwrapper teardown.
Plugin: helpconfig, Hook: pytest_cmdline_parse
ConftestImportFailure: ModuleNotFoundError: No module named 'hypothesis' (from /repo/test/mitmproxy/proxy/conftest.py)
For more information see https://pluggy.readthedocs.io/en/stable/api_reference.html#pluggy.PluggyTeardownRaisedWarning
  config = pluginmanager.hook.pytest_cmdline_parse(
ImportError while loading conftest '/repo/test/mitmproxy/proxy/conftest.py'.
test/mitmproxy/proxy/conftest.py:4: in <module>
    from hypothesis import settings
E   ModuleNotFoundError: No module named 'hypothesis'
step 17
$ cd /repo && pip install hypothesis -q 2>&1 | tail -1; timeout 600 python -m pytest test/mitmproxy/proxy/test_mode_servers.py -q 2>&1 | tail -15
output · 10 lines
ERROR: No matching distribution found for hypothesis
/usr/local/lib/python3.12/site-packages/_pytest/config/__init__.py:331: PluggyTeardownRaisedWarning: A plugin raised an exception during an old-style hookwrapper teardown.
Plugin: helpconfig, Hook: pytest_cmdline_parse
ConftestImportFailure: ModuleNotFoundError: No module named 'hypothesis' (from /repo/test/mitmproxy/proxy/conftest.py)
For more information see https://pluggy.readthedocs.io/en/stable/api_reference.html#pluggy.PluggyTeardownRaisedWarning
  config = pluginmanager.hook.pytest_cmdline_parse(
ImportError while loading conftest '/repo/test/mitmproxy/proxy/conftest.py'.
test/mitmproxy/proxy/conftest.py:4: in <module>
    from hypothesis import settings
E   ModuleNotFoundError: No module named 'hypothesis'
step 18
$ cd /repo && ls test/mitmproxy/net/conftest.py 2>/dev/null; timeout 300 python -m pytest test/mitmproxy/net/test_udp.py test/mitmproxy/net/test_udp_wireguard.py -q -p no:cacheprovider 2>&1 | tail -8
output · 2 lines
..
2 passed in 2.33s
step 19
$ cat > /tmp/repro2.py <<'EOF'
output · 4 lines
s1 OK
s2 OK
s3 OK
rc=0
step 20
$ cd /repo && sed -n 1,140p mitmproxy/proxy/server.py
output · 39 lines
"""
Proxy Server Implementation using asyncio.
The very high level overview is as follows:

    - Spawn one coroutine per client connection and create a reverse proxy layer to example.com
    - Process any commands from layer (such as opening a server connection)
    - Wait for any IO and send it as events to top layer.
"""
import abc
import asyncio
import collections
import logging
27 more lines
import time
import traceback
from collections.abc import Awaitable, Callable, MutableMapping
from contextlib import contextmanager
from dataclasses import dataclass
from typing import Optional, Union

import mitmproxy_wireguard as wg
from OpenSSL import SSL

from mitmproxy import http, options as moptions, tls
from mitmproxy.proxy.context import Context
from mitmproxy.proxy.layers.http import HTTPMode
from mitmproxy.proxy import commands, events, layer, layers, mode_specs, server_hooks
from mitmproxy.conn
…[3495 chars omitted]…
r",
                client=self.client.peername,
            )
            self.transports[self.client].handler = handler
            self.server_event(events.Start())
            await asyncio.wait([handler])
            if not handler.cancelled() and (e := handler.exception()):
                self.log(f"mitmproxy has crashed!\n{traceback.format_exception(e)}", logging.ERROR)

        watch.cancel()
step 21
$ cd /repo && sed -n 1,60p mitmproxy/proxy/layers/dns.py 2>/dev/null | head -5; grep -rn "reverse" mitmproxy/proxy/mode_servers.py | head
output · 6 lines
from dataclasses import dataclass
import struct

from mitmproxy import dns, flow as mflow
from mitmproxy.proxy import commands, events, layer
7:    mode = ProxyMode.parse("reverse:https://example.com")
step 22
$ cat > /tmp/repro3.py <<'EOF'
output · 26 lines
DEBUG:asyncio:Using selector: EpollSelector
Traceback (most recent call last):
  File "/tmp/repro3.py", line 39, in <module>
    asyncio.run(main())
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 195, in run
    return runner.run(main)
           ^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 118, in run
    return self._loop.run_until_complete(task)
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/base_events.py", line 691, in run_until_complete
    return future.result()
14 more lines
           ^^^^^^^^^^^^^^^
  File "/tmp/repro3.py", line 24, in main
    inst = mode_servers.ServerInstance.make(f"reverse:dns+udp://127.0.0.1:{uport}", Mgr())
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/repo/mitmproxy/proxy/mode_servers.py", line 97, in make
    mode = mode_spec
…[257 chars omitted]…
", line 218, in __post_init__
    self.scheme, self.address = server_spec.parse(self.data, default_scheme="https")
                                ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/repo/mitmproxy/net/server_spec.py", line 42, in parse
    raise ValueError(f"Invalid server specification: {server_spec}")
ValueError: Invalid server specification: dns+udp://127.0.0.1:41079
rc=0
step 23
$ cd /repo && grep -n "dns" mitmproxy/net/server_spec.py | head; grep -n "class ReverseProxy\|dns" mitmproxy/proxy/mode_specs.py | head -20
output · 6 lines
11:    Literal["http", "https", "tls", "dtls", "tcp", "udp", "dns"],
48:    if scheme not in ("http", "https", "tls", "dtls", "tcp", "udp", "dns"):
65:                "dns": 53,
213:    scheme: Literal["http", "https", "tls", "dtls", "tcp", "udp", "dns"]
219:        if self.scheme in ("dns", "dtls", "udp"):
225:        if self.scheme == "dns":
step 24
$ cd /repo && sed -i 's/reverse:dns+udp:\/\//reverse:udp:\/\//' /tmp/repro3.py && timeout 60 python /tmp/repro3.py 2>&1 | tail -40; echo rc=$?
output · 20 lines
DEBUG:asyncio:Using selector: EpollSelector
Traceback (most recent call last):
  File "/tmp/repro3.py", line 39, in <module>
    asyncio.run(main())
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 195, in run
    return runner.run(main)
           ^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 118, in run
    return self._loop.run_until_complete(task)
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/base_events.py", line 691, in run_until_complete
    return future.result()
8 more lines
           ^^^^^^^^^^^^^^^
  File "/tmp/repro3.py", line 25, in main
    await inst.start()
  File "/repo/mitmproxy/proxy/mode_servers.py", line 214, in start
    host = self.mode.listen_host(ctx.options.listen_host)
                                 ^^^^^^^^^^^
AttributeError: module 'mitmproxy.ctx' has no attribute 'options'
rc=0
step 25
$ cd /repo && python - <<'EOF'
output · 25 lines
^^^^^^^^^^^^^^^^^^^^^^^
  File "/repo/mitmproxy/net/udp.py", line 240, in start_server
    _, protocol = await loop.create_datagram_endpoint(
                  ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/base_events.py", line 1379, in create_datagram_endpoint
    infos = await self._ensure_resolved(
            ^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/base_events.py", line 1466, in _ensure_resolved
    return await loop.getaddrinfo(host, port, family=family, type=type,
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/base_events.py", line 905, in getaddrinfo
    return await self.run_in_executor(
13 more lines
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/concurrent/futures/thread.py", line 59, in run
    result = self.fn(*self.args,
…[870 chars omitted]…
ine 691, in run_until_complete
    return future.result()
           ^^^^^^^^^^^^^^^
  File "/tmp/repro3.py", line 27, in main
    await inst.start()
  File "/repo/mitmproxy/proxy/mode_servers.py", line 225, in start
    raise OSError(e.errno, message, e.filename) from e
OSError: [Errno -2] reverse proxy to udp://127.0.0.1:34764 failed to listen on *:8080 with [Errno -2] Name or service not known
rc=0
step 26
$ cd /repo && sed -i 's/with taddons.context() as tctx:/with taddons.context() as tctx:\n        tctx.options.listen_host = "127.0.0.1"\n        tctx.options.listen_port = 0/' /tmp/repro3.py && timeout 60 python /tmp/repro3.py 2>&1 | tail -40; echo rc=$?
output · 20 lines
DEBUG:passlib.utils.compat:loaded lazy attr 'SafeConfigParser': <class 'configparser.ConfigParser'>
DEBUG:passlib.utils.compat:loaded lazy attr 'NativeStringIO': <class '_io.StringIO'>
DEBUG:passlib.utils.compat:loaded lazy attr 'BytesIO': <class '_io.BytesIO'>
DEBUG:passlib.registry:registered 'htdigest' handler: <class 'passlib.handlers.digests.htdigest'>
DEBUG:passlib.registry:registered 'bcrypt' handler: <class 'passlib.handlers.bcrypt.bcrypt'>
DEBUG:passlib.handlers.bcrypt:trapped unexpected response from 'os_crypt' backend: verify('test', b'$2$04$5BJqKfqMQvV7nS.yUguNcuRfMMOXK0xPWavM7pOzjEi5ze5T1k8/S'):
Traceback (most recent call last):
  File "/usr/local/lib/python3.12/site-packages/passlib/handlers/bcrypt.py", line 303, in safe_verify
    return verify(secret, hash)
           ^^^^^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/site-packages/passlib/utils/handlers.py", line
…[2057 chars omitted]…
8 more lines
addr)
  File "/repo/mitmproxy/net/udp.py", line 109, in datagram_received
    self._datagram_received_cb(self._transport, data, addr, self._local_addr)
  File "/repo/mitmproxy/proxy/mode_servers.py", line 178, in handle_udp_datagram
    if connection_id not in self.manager.connections:
                            ^^^^^^^^^^^^^^^^^^^^^^^^
AttributeError: 'Mgr' object has no attribute 'connections'
rc=0
step 27
$ cd /repo && python - <<'EOF'
output · 22 lines
DEBUG:asyncio:Using selector: EpollSelector
Traceback (most recent call last):
    return verify(secret, hash)
           ^^^^^^^^^^^^^^^^^^^^
    return consteq(self._calc_checksum(secret), chk)
                   ^^^^^^^^^^^^^^^^^^^^^^^^^^^
    raise uh.exc.InternalBackendError(
INFO:mitmproxy.proxy.mode_servers:reverse proxy to udp://127.0.0.1:45294 listening at 127.0.0.1:33312.
ERROR:asyncio:Exception in callback _SelectorDatagramTransport._read_ready()
handle: <Handle _SelectorDatagramTransport._read_ready()>
Traceback (most recent call last):
  File "/usr/local/lib/python3.12/asyncio/events.py", line 88, in _run
10 more lines
    self._context.run(self._callback, *self._args)
  File "/usr/local/lib/python3.12/asyncio/selector_events.py", line 1251, in _read_ready
    self._protocol.datagram_received(data, addr)
  File "/repo/mitmproxy/net/udp.py", line 109, in datagram_received
    self._datagram_received_cb(self._transport, data, addr, self._local_addr)
  File "/repo/mitmproxy/proxy/mode_servers.py", line 178, in handle_udp_datagram
    if connection_id not in self.manager.connections:
                            ^^^^^^^^^^^^^^^^^^^^^^^^
AttributeError: 'Mgr' object has no attribute 'connections'
rc=0
step 28
$ cd /repo && python - <<'EOF'
output · 4 lines
File "/tmp/repro3.py", line 16
    class Mgr:
IndentationError: unexpected indent
rc=0
step 29
$ cd /repo && sed -n 10,25p /tmp/repro3.py
output · 16 lines
async def echo(transport, data, addr, local):
        transport.sendto(data, addr)
    upstream = await udp.start_server(echo, "127.0.0.1", 0)
    uport = upstream.sockets[0].getsockname()[1]

    connections = {}
        class Mgr:
        connections = {}
        from contextlib import contextmanager
        @contextmanager
        def register_connection(self, cid, handler):
            connections[cid] = handler
4 more lines
            yield
            connections.pop(cid, None)

    from mitmproxy.test import taddons
step 30
$ cd /repo && sed -i '13s/^        class Mgr:/    class Mgr:/' /tmp/repro3.py && timeout 60 python /tmp/repro3.py 2>&1 | grep -vi "passlib\|verify\|consteq\|InternalBackend\|lazy" | tail -40; echo rc=$?
output · 4 lines
File "/tmp/repro3.py", line 16
    class Mgr:
IndentationError: unexpected indent
rc=0
step 31
$ cd /repo && python - <<'EOF'
output · 27 lines
upstream = await udp.start_server(echo, "127.0.0.1", 0)
    uport = upstream.sockets[0].getsockname()[1]

    connections = {}
    class Mgr:
        connections = {}
        from contextlib import contextmanager
        @contextmanager
        def register_connection(self, cid, handler):
            connections[cid] = handler
            yield
                   ^^^^^^^^^^^^^^
15 more lines
  File "/repo/mitmproxy/proxy/layer.py", line 144, in handle_event
    command = command_generator.send(send)
              ^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/repo/mitmproxy/proxy/layers/modes.py", line 74, in _handle_event
    err = yield from self.finish_start()
          ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/repo/mitmproxy/proxy/layers/modes.py", line 38, in finish_start
    self.context.options.connection_strategy == "eager"
    ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/repo/mitmproxy/optma
…[1379 chars omitted]…
ent type at ReverseProxy._handle_event: Expected Start, got ConnectionClosed(connection=Client({'id': '…cf2a73', 'peername': ('127.0.0.1', 50555), 'sockname': ('127.0.0.1', 59006), 'timestamp_start': 1790907011.296341, 'state': <ConnectionState.CLOSED: 0>, 'transport_protocol': 'udp', 'proxy_mode': ProxyMode.parse('reverse:udp://127.0.0.1:60778')})).

INFO:mitmproxy.proxy.server:client disconnect
rc=0
step 32
$ cd /repo && head -30 CHANGELOG.md && git log -1 --format=%ci
output · 25 lines
# Release History

## Unreleased: mitmproxy next

* Add UDP NextLayer detection, update proxy mode syntax.
  ([#5556])(https://github.com/mitmproxy/mitmproxy/issues/5556), @meitinger)
* Fix running hook issue and signal proxyserver updates.
  ([#5376](https://github.com/mitmproxy/mitmproxy/issues/5376), @meitinger)
* Include server information in bind/listen errors.
  ([#5495](https://github.com/mitmproxy/mitmproxy/pull/5495), @meitinger)
* Include information about lazy connection_strategy in related errors.
  ([#5465](https://github.com/mitmproxy/mitmproxy/pull/5465), @meitinger, @mhils)
13 more lines
* Add UDP layer and flow support.
  ([#5414](https://github.com/mitmproxy/mitmproxy/pull/5414), @meitinger)
* Setting `connection_strategy` to `lazy` now also disables early 
  upstream connections to fetch TLS certificate details.
  ([#5487](https://github.com/mitmproxy/mitmproxy/pull/5487), @mhils
…[469 chars omitted]…
s://github.com/mitmproxy/mitmproxy/pull/5547), @jorants).
* Render application/prpc content as gRPC/Protocol Buffers
  ([#5568](https://github.com/mitmproxy/mitmproxy/pull/5568), @selfisekai)
* Mitmweb now supports `content_view_lines_cutoff`.
  ([#5548](https://github.com/mitmproxy/mitmproxy/pull/5548), @sanlengjingvv)
* Fix a mitmweb crash when scrolling down the flow list.
2026-10-02 01:57:40 +0000
step 33
$ cd /repo && cat mitmproxy/net/udp_wireguard.py
output · 35 lines
"""
This module contains a mock DatagramTransport for use with mitmproxy-wireguard.
"""
import asyncio
from typing import Any

import mitmproxy_wireguard as wg

from mitmproxy.connection import Address


class WireGuardDatagramTransport(asyncio.DatagramTransport):
23 more lines
    def __init__(self, server: wg.Server, local_addr: Address, remote_addr: Address):
        self._server: wg.Server = server
        self._local_addr: Address = local_addr
        self._remote_addr: Address = remote_addr
        super().__init__()

    def sendto(self, data, addr=None):
        self._server.send_datagram(data, self._local_addr, addr or self._remote_addr)

    def get_extra_info(self, name: str, default: Any = None) -> Any:
        if name == "sockname":
            return self._server.getsockname()
        else:
            raise NotImplementedError

    def get_protocol(self):
        return self

    async def drain(self) -> None:
        pass

    async def wait_closed(self) -> None:
        pass
step 34
$ cd /repo && sed -n 1,60p mitmproxy/test/taddons.py
output · 41 lines
import asyncio

import mitmproxy.master
import mitmproxy.options
from mitmproxy import hooks
from mitmproxy import command
from mitmproxy import eventsequence
from mitmproxy.addons import script, core


class context:
    """
29 more lines
    A context for testing addons, which sets up the mitmproxy.ctx module so
    handlers can run as they would within mitmproxy. The context also
    provides a number of helper methods for common testing scenarios.
    """

    def __init__(self, *addons, options=None, loadcore=True):
        try:
            loop = asyncio.get_running_loop()
        except RuntimeError:
            loop = asyncio.new_event_loop()

        options = options or mitmproxy.options.Options()
        self.master = mitmproxy.master.Master(options, event_loop=loop)
        self.options = self.master.options

        if loadcore:
            self.master.addons.add(core.Core())

        
…[593 chars omitted]…
 methods. Modifies the registered
        Options object with the given keyword arguments, then calls the
        configure method on the addon with the updated value.
        """
        if addon not in self.master.addons:
            self.master.addons.register(addon)
        with self.options.rollback(kwargs.keys(), reraise=True):
            if kwargs:
                self.options.update(**kwargs)
step 35
$ cd /repo && python - <<'EOF'
output · 15 lines
Traceback (most recent call last):
  File "<stdin>", line 7, in <module>
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 195, in run
    return runner.run(main)
           ^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 118, in run
    return self._loop.run_until_complete(task)
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/base_events.py", line 691, in run_until_complete
    return future.result()
           ^^^^^^^^^^^^^^^
  File "<stdin>", line 6, in main
3 more lines
  File "/repo/mitmproxy/optmanager.py", line 188, in __getattr__
    raise AttributeError("No such option: %s" % attr)
AttributeError: No such option: connection_strategy
step 36
$ cd /repo && python -c "
output · 2 lines
[]
27
step 37
$ cd /repo && grep -rn "connection_strategy" mitmproxy/options.py mitmproxy/addons/*.py | head
output · 2 lines
mitmproxy/addons/proxyserver.py:126:            "connection_strategy",
mitmproxy/addons/tlsconfig.py:139:            conn_context.server.tls and ctx.options.connection_strategy == "eager"
step 38
$ cd /repo && cat > test/mitmproxy/proxy/test_tmp_repro.py <<'EOF'
output · 9 lines
/usr/local/lib/python3.12/site-packages/_pytest/config/__init__.py:331: PluggyTeardownRaisedWarning: A plugin raised an exception during an old-style hookwrapper teardown.
Plugin: helpconfig, Hook: pytest_cmdline_parse
ConftestImportFailure: ModuleNotFoundError: No module named 'hypothesis' (from /repo/test/mitmproxy/proxy/conftest.py)
For more information see https://pluggy.readthedocs.io/en/stable/api_reference.html#pluggy.PluggyTeardownRaisedWarning
  config = pluginmanager.hook.pytest_cmdline_parse(
ImportError while loading conftest '/repo/test/mitmproxy/proxy/conftest.py'.
test/mitmproxy/proxy/conftest.py:4: in <module>
    from hypothesis import settings
E   ModuleNotFoundError: No module named 'hypothesis'
step 39
$ cd /repo && mv test/mitmproxy/proxy/test_tmp_repro.py test/mitmproxy/test_tmp_repro.py && timeout 120 python -m pytest test/mitmproxy/test_tmp_repro.py -q -p no:cacheprovider 2>&1 | tail -30
output · 21 lines
try:
            self._server = await self.listen(host, port)
            self._listen_addrs = tuple(s.getsockname() for s in self._server.sockets)
        except OSError as e:
            self.last_exception = e
            message = f"{self.mode.description} failed to listen on {host or '*'}:{port} with {e}"
            if e.errno == errno.EADDRINUSE and self.mode.custom_listen_port is None:
                assert self.mode.custom_listen_host is None  # since [@ [listen_addr:]listen_port]
                message += f"\nTry specifying a different port by using `--mode {self.mode.full_spec}@{port + 1}`."
>           raise OSError(e.errno, message, e.filename) from e
E           OSError: [Errno -2] reverse proxy to udp://127.0.0.1:48131 failed to listen on *:8080 with [Errno -2] Name or service not known
9 more lines
mitmproxy/proxy/mode_servers.py:225: OSError
[
…[899 chars omitted]…
ypt as _crypt

-- Docs: https://docs.pytest.org/en/stable/how-to/capture-warnings.html
=========================== short test summary info ============================
FAILED test/mitmproxy/test_tmp_repro.py::test_udp_close - OSError: [Errno -2] reverse proxy to udp://127.0.0.1:48131 failed to listen...
1 failed, 3 warnings in 2.53s
step 40
$ cd /repo && python - <<'EOF'
step 41
$ cd /repo && timeout 120 python -m pytest test/mitmproxy/test_tmp_repro.py -q -p no:cacheprovider 2>&1 | tail -40; echo rc=$?
output · 1 lines
rc=0
step 42
$ cd /repo && timeout 120 python -m pytest test/mitmproxy/test_tmp_repro.py -p no:cacheprovider 2>&1 | tail -40; echo rc=$?
output · 9 lines
============================= test session starts ==============================
platform linux -- Python 3.12.13, pytest-7.4.4, pluggy-1.6.0
rootdir: /repo
configfile: setup.cfg
plugins: asyncio-0.23.8, cov-7.1.0, timeout-2.4.0
asyncio: mode=Mode.AUTO
collected 1 item

test/mitmproxy/test_tmp_repro.py rc=0
step 43
$ cd /repo && python - <<'EOF'
output · 9 lines
============================= test session starts ==============================
platform linux -- Python 3.12.13, pytest-7.4.4, pluggy-1.6.0
rootdir: /repo
configfile: setup.cfg
plugins: asyncio-0.23.8, cov-7.1.0, timeout-2.4.0
asyncio: mode=Mode.AUTO
collected 1 item

test/mitmproxy/test_tmp_repro.py rc=0
step 44
$ cd /repo && timeout 30 python -X faulthandler -m pytest test/mitmproxy/test_tmp_repro.py -s -p no:cacheprovider 2>&1 | tail -5; python - <<'EOF'
output · 25 lines
plugins: asyncio-0.23.8, cov-7.1.0, timeout-2.4.0
asyncio: mode=Mode.AUTO
collected 1 item

test/mitmproxy/test_tmp_repro.py     return self._hookexec(self.name, self._hookimpls.copy(), kwargs, firstresult)
  File "/usr/local/lib/python3.12/site-packages/pluggy/_manager.py", line 120, in _hookexec
    return self._inner_hookexec(hook_name, methods, kwargs, firstresult)
  File "/usr/local/lib/python3.12/site-packages/pluggy/_callers.py", line 121, in _multicall
    res = hook_impl.function(*args)
  File "/usr/local/lib/python3.12/site-packages/_pytest/runner.py", line 169, in pytest_runtest_call
    item.runtest()
  File "/usr/local/lib/python3.12/site-packages/pytest_asyncio/plugin.py", line 440, in runtest
13 more lines
    super().runtest()
  File "/usr/local/lib/python3.12/site-packages/_pytest/python.py", line 1792, in runtest
    self.ihook.pytest_pyfunc_call(pyfuncitem=self)
  File "/usr/loca
…[872 chars omitted]…
ase_events.py", line 645, in run_forever
    self._run_once()
  File "/usr/local/lib/python3.12/asyncio/base_events.py", line 1961, in _run_once
    event_list = self._selector.select(timeout)
  File "/usr/local/lib/python3.12/selectors.py", line 468, in select
    fd_event_list = self._selector.poll(timeout, max_ev)
+++++++++++++++++++++++++++++++++++ Timeout ++++++++++++++++++++++++++++++++++++
rc=0
step 45
$ cd /repo && cat > /tmp/repro4.py <<'EOF'
output · 4 lines
starting
started
connected
rc=0
step 46
$ cd /repo && timeout 20 python /tmp/repro4.py 2>&1 | grep -vi "passlib\|verify\|consteq\|InternalBackend\|lazy\|crypt" | head -60; echo rc=$?
output · 4 lines
starting
started
connected
rc=0
step 47
$ cd /repo && python - <<'EOF'
output · 6 lines
Traceback (most recent call last):
  File "/tmp/repro4.py", line 2, in <module>
    logging.basicConfig(level=logging.DEBUG), faulthandler, signal
                                              ^^^^^^^^^^^^
NameError: name 'faulthandler' is not defined. Did you forget to import 'faulthandler'?
rc=0
step 48
$ cd /repo && sed -i '2s/.*/import asyncio, sys, faulthandler, signal, logging/' /tmp/repro4.py && sed -i '3s/.*/logging.basicConfig(level=logging.DEBUG)/' /tmp/repro4.py && head -5 /tmp/repro4.py && timeout 20 python /tmp/repro4.py 2>&1 | grep -vi "passlib\|verify\|consteq\|InternalBackend\|lazy\|crypt\|asyncio:" | head -80; echo rc=$?
output · 32 lines
import asyncio, sys, logging
import asyncio, sys, faulthandler, signal, logging
logging.basicConfig(level=logging.DEBUG)
from mitmproxy.net import udp
from mitmproxy.proxy import mode_servers
Traceback (most recent call last):
           ^^^^^^^^^^^^^^^^^^^^
                   ^^^^^^^^^^^^^^^^^^^^^^^^^^^
starting
INFO:mitmproxy.proxy.mode_servers:reverse proxy to udp://127.0.0.1:44567 listening at 127.0.0.1:8080.
started
connected
20 more lines
INFO:mitmproxy.proxy.server:client connect
ERROR:mitmproxy.proxy.server:mitmproxy has crashed!
Traceback (most recent call last):
  File "/repo/mitmproxy/proxy/server.py", line 344, in server_event
    for command in layer_commands:
                   ^^^^^^^^^^^^^^
  File "/repo/mitmproxy/proxy/layer.py", line 144, in handle_event
    command = command_generator.send(send)
              ^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/repo/mitmproxy/proxy/layers/modes
…[677 chars omitted]…
^^^^^^
  File "/repo/mitmproxy/proxy/layer.py", line 135, in handle_event
    command_generator = self._handle_event(event)
                        ^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/repo/mitmproxy/proxy/utils.py", line 26, in _check_event_type
    raise AssertionError(
AssertionError: Unexpected event type at ReverseProxy._handle_event: Expected Start, got DataReceived(client, b'\x00\x00\x01').

rc=0
step 49
$ cd /repo && grep -n "def add_option\|def register" mitmproxy/optmanager.py | head; grep -n "connection_strategy" mitmproxy/addons/proxyserver.py
output · 2 lines
116:    def add_option(
126:            "connection_strategy",
step 50
$ cd /repo && sed -n 110,140p mitmproxy/addons/proxyserver.py
output · 31 lines
self.servers = Servers(self)
        self.is_running = False

    def __repr__(self):
        return f"Proxyserver({len(self.connections)} active conns)"

    @contextmanager
    def register_connection(self, connection_id: tuple, handler: ProxyConnectionHandler):
        self.connections[connection_id] = handler
        try:
            yield
        finally:
19 more lines
            del self.connections[connection_id]

    def load(self, loader):
        loader.add_option(
            "connection_strategy",
            str,
            "eager",
            "Determine when server connections should be established. When set to lazy, mitmproxy "
            "tries to defer establishing an upstream connection as long as possible. This makes it possible to "
            "use server replay while being offline. When set to eager, mitmproxy can detect protocols with "
            "server-side greetings, as well as accurately mirror TLS ALPN negotiation.",
            choices=("eager", "lazy"),
        )
        loader.add_option(
            "stream_large_bodies",
            Optional[str],
            None,
            """
            Stream data to the client if response body exceeds the given
step 51
$ cd /repo && python - <<'EOF'
output · 20 lines
Traceback (most recent call last):
           ^^^^^^^^^^^^^^^^^^^^
                   ^^^^^^^^^^^^^^^^^^^^^^^^^^^
Traceback (most recent call last):
  File "/tmp/repro4.py", line 46, in <module>
    asyncio.run(test())
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 195, in run
    return runner.run(main)
           ^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 118, in run
    return self._loop.run_until_complete(task)
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
8 more lines
  File "/usr/local/lib/python3.12/asyncio/base_events.py", line 691, in run_until_complete
    return future.result()
           ^^^^^^^^^^^^^^^
  File "/tmp/repro4.py", line 23, in test
    with taddons.context(proxyserver.ProxyServer()) as tctx:
                         ^^^^^^^^^^^^^^^^^^^^^^^
AttributeError: module 'mitmproxy.addons.proxyserver' has no attribute 'ProxyServer'. Did you mean: 'Proxyserver'?
rc=0
step 52
$ cd /repo && sed -i 's/proxyserver.ProxyServer()/proxyserver.Proxyserver()/' /tmp/repro4.py && timeout 20 python /tmp/repro4.py 2>&1 | grep -vi "passlib\|verify\|consteq\|InternalBackend\|lazy\|crypt\|asyncio:" | head -100; echo rc=$?
output · 9 lines
Traceback (most recent call last):
           ^^^^^^^^^^^^^^^^^^^^
                   ^^^^^^^^^^^^^^^^^^^^^^^^^^^
starting
INFO:mitmproxy.proxy.mode_servers:reverse proxy to udp://127.0.0.1:37681 listening at 127.0.0.1:8080.
started
connected
INFO:mitmproxy.proxy.server:client connect
rc=0
step 53
$ cd /repo && sed -n 140,240p mitmproxy/proxy/server.py
output · 32 lines
watch.cancel()
        while self.wakeup_timer:
            timer = self.wakeup_timer.pop()
            timer.cancel()

        self.log("client disconnect")
        self.client.timestamp_end = time.time()
        await self.handle_hook(server_hooks.ClientDisconnectedHook(self.client))

        if self.transports:
            self.log("closing transports...", logging.DEBUG)
            for io in self.transports.values():
20 more lines
                if io.handler:
                    io.handler.cancel("client disconnected")
            await asyncio.wait(
                [x.handler for x in self.transports.values() if x.handler]
            )
            self.log("transports closed!", logging.DEBUG)

    async def open_connection(self, command: commands.OpenConnection) -> None:
        if not command.connection.address:
            self.log(f"Cannot open connection, no hostname given.")
  
…[3972 chars omitted]…
p here once that is done.
                new_handler = asyncio_utils.create_task(
                    self.handle_connection(command.connection),
                    name=f"server connection handler for {addr}",
                    client=self.client.peername,
                )
                self.transports[command.connection].handler = new_handler
                await asyncio.wait([new_handler])
step 54
$ pip show mitmproxy 2>/dev/null | head -3; ls /usr/local/lib/python3.12/site-packages | grep -i mitm
output · 8 lines
Name: mitmproxy
Version: 9.0.0.dev0
Summary: An interactive, SSL/TLS-capable intercepting proxy for HTTP/1, HTTP/2, and WebSockets.
__editable__.mitmproxy-9.0.0.dev0.pth
__editable___mitmproxy_9_0_0_dev0_finder.py
mitmproxy-9.0.0.dev0.dist-info
mitmproxy_wireguard
mitmproxy_wireguard-0.1.23.dist-info
step 55
$ cat > /tmp/repro5.py <<'EOF'
output · 3 lines
c OK
s OK
rc=0
step 56
$ cd /repo && sed -n 1,120p mitmproxy/proxy/layers/modes.py
output · 49 lines
import socket
import struct
from abc import ABCMeta
from dataclasses import dataclass
from typing import Optional

from mitmproxy import connection
from mitmproxy.proxy import commands, events, layer
from mitmproxy.proxy.commands import StartHook
from mitmproxy.proxy.layers import dns, tls
from mitmproxy.proxy.mode_specs import ReverseMode
from mitmproxy.proxy.utils import expect
37 more lines

class HttpProxy(layer.Layer):
    @expect(events.Start)
    def _handle_event(self, event: events.Event) -> layer.CommandGenerator[None]:
        child_layer = layer.NextLayer(self.context)
        self._handle_event = child_layer.handle_event
        yield from child_layer.handle_event(event)


class HttpUpstreamProxy(layer.Layer):
    @expect(events.Start)
    def _handle_event(self, event: events.Event) -> layer.CommandGenerator[None]:
        child_layer = layer.NextLayer(self.context)
        self._han
…[2743 chars omitted]…
 0x07
SOCKS5_REP_ADDRESS_TYPE_NOT_SUPPORTED = 0x08


@dataclass
class Socks5AuthData:
    client_conn: connection.Client
    username: str
    password: str
    valid: bool = False


@dataclass
class Socks5AuthHook(StartHook):
    """
    Mitmproxy has received username/password SOCKS5 credentials.

    This hook decides whether they are valid by setting `data.valid`.
    """

    data: Socks5AuthData
step 57
$ python - <<'EOF'
output · 36 lines
def sendto(self, data, addr=None):
        if not isinstance(data, (bytes, bytearray, memoryview)):
            raise TypeError(f'data argument must be a bytes-like object, '
                            f'not {type(data).__name__!r}')
        if not data:
            return

        if self._address:
            if addr not in (None, self._address):
                raise ValueError(
                    f'Invalid address: must be None or {self._address}')
            addr = self._address
24 more lines
        if self._conn_lost and self._address:
            if self._conn_lost >= constants.LOG_THRESHOLD_FOR_CONNLOST_WRITES:
                logger.warning('socket.send() raised exception.')
            self._conn_lost += 1
            return

        if not self._buffer:
            # Attempt to send it right away first.
            try:
                if self._extra['peername']:
               
…[375 chars omitted]…
cept (SystemExit, KeyboardInterrupt):
                raise
            except BaseException as exc:
                self._fatal_error(
                    exc, 'Fatal write error on datagram transport')
                return

        # Ensure that what we buffer is immutable.
        self._buffer.append((bytes(data), addr))
        self._buffer_size += len(data)
        self._maybe_pause_protocol()
step 58
$ cd /repo && grep -n "udp\|transport_protocol" mitmproxy/proxy/layer.py mitmproxy/proxy/layers/next.py | head -20
output · 1 lines
grep: mitmproxy/proxy/layers/next.py: No such file or directory
step 59
$ cd /repo && grep -rn "udp" mitmproxy/proxy/layer.py mitmproxy/proxy/layers/__init__.py 2>/dev/null | head; ls mitmproxy/proxy/layers
output · 10 lines
mitmproxy/proxy/layers/__init__.py:5:from .udp import UDPLayer
__init__.py
__pycache__
dns.py
http
modes.py
tcp.py
tls.py
udp.py
websocket.py
step 60
$ cd /repo && cat mitmproxy/proxy/layers/udp.py; grep -n "UDPLayer\|udp" mitmproxy/proxy/layers/__init__.py
output · 57 lines
from dataclasses import dataclass
from typing import Optional

from mitmproxy import flow, udp
from mitmproxy.proxy import commands, events, layer
from mitmproxy.proxy.commands import StartHook
from mitmproxy.connection import Connection
from mitmproxy.proxy.context import Context
from mitmproxy.proxy.events import MessageInjected
from mitmproxy.proxy.utils import expect

45 more lines
@dataclass
class UdpStartHook(StartHook):
    """
    A UDP connection has started.
    """

    flow: udp.UDPFlow


@dataclass
class UdpMessageHook(StartHook):
    """
    A UDP connection has received a message. The most recent message
    will be flow.messages[-1]. The message is user-modifiable.
    """

    flow: udp.UDPFlow


@dataclass
class UdpEndHook(StartHook):
    """
    A UDP connection has ended.
    """

    flow: udp.UDPFlow


@dataclass
class UdpErrorHook(StartHook):
    """
    A UDP error has occu
…[2675 chars omitted]…
one
                if self.flow:
                    yield UdpEndHook(self.flow)
                    self.flow.live = False
        else:
            raise AssertionError(f"Unexpected event: {event}")

    @expect(events.DataReceived, events.ConnectionClosed, UdpMessageInjected)
    def done(self, _) -> layer.CommandGenerator[None]:
        yield from ()
5:from .udp import UDPLayer
14:    "UDPLayer",
step 61
$ cd /repo && grep -n "udp\|UDP" mitmproxy/addons/next_layer.py | head -20
output · 13 lines
48:    udp_hosts: Iterable[re.Pattern] = ()
55:        if "udp_hosts" in updated:
56:            self.udp_hosts = [
57:                re.compile(x, re.IGNORECASE) for x in ctx.options.udp_hosts
216:        elif context.client.transport_protocol == "udp":
230:                return layers.UDPLayer(context, ignore=True)
238:            # 4. Check for --udp
239:            if self.is_destination_in_hosts(context, self.udp_hosts):
240:                return layers.UDPLayer(context)
250:            # 6. Check for raw udp mode.
251:            if ctx.options.rawudp:
252:                return layers.UDPLayer(context)
1 more lines
255:            return layers.UDPLayer(context, ignore=True)
step 62
$ cd /repo && sed -n 200,260p mitmproxy/addons/next_layer.py
output · 33 lines
return layers.TCPLayer(context)

            # 5. Check for raw tcp mode.
            very_likely_http = context.client.alpn and context.client.alpn in HTTP_ALPNS
            probably_no_http = not very_likely_http and (
                not data_client[
                    :3
                ].isalpha()  # the first three bytes should be the HTTP verb, so A-Za-z is expected.
                or data_server  # a server greeting would be uncharacteristic.
            )
            if ctx.options.rawtcp and probably_no_http:
                return layers.TCPLayer(context)
21 more lines
            # 6. Assume HTTP by default.
            return layers.HttpLayer(context, HTTPMode.transparent)

        elif context.client.transport_protocol == "udp":
            # unlike TCP, we make a decision immediately
            try:
                dtls_client_hello = dtls_parse_client_hello(data_
…[941 chars omitted]…
e:
                return layers.DNSLayer(context)

            # 6. Check for raw udp mode.
            if ctx.options.rawudp:
                return layers.UDPLayer(context)

            # 7. Ignore the connection by default. (In the future, we'll assume HTTP/3)
            return layers.UDPLayer(context, ignore=True)

        else:
            raise AssertionError(context.client.transport_protocol)
step 63
$ cd /repo && grep -rn "rawudp" mitmproxy/options.py mitmproxy/addons/next_layer.py | head -3
output · 2 lines
mitmproxy/options.py:152:            "rawudp",
mitmproxy/addons/next_layer.py:251:            if ctx.options.rawudp:
step 64
$ cd /repo && python - <<'EOF'
output · 9 lines
Traceback (most recent call last):
           ^^^^^^^^^^^^^^^^^^^^
                   ^^^^^^^^^^^^^^^^^^^^^^^^^^^
starting
INFO:mitmproxy.proxy.mode_servers:reverse proxy to udp://127.0.0.1:51723 listening at 127.0.0.1:8080.
started
connected
INFO:mitmproxy.proxy.server:client connect
rc=0
step 65
$ cd /repo && python - <<'EOF'
output · 13 lines
Traceback (most recent call last):
           ^^^^^^^^^^^^^^^^^^^^
                   ^^^^^^^^^^^^^^^^^^^^^^^^^^^
starting
INFO:mitmproxy.proxy.mode_servers:reverse proxy to udp://127.0.0.1:58080 listening at 127.0.0.1:8080.
started
connected
INFO:mitmproxy.proxy.server:client connect
INFO:mitmproxy.proxy.server:server connect 127.0.0.1:58080
/repo/mitmproxy/net/udp.py:109: RuntimeWarning: coroutine 'test.<locals>.echo' was never awaited
  self._datagram_received_cb(self._transport, data, addr, self._local_addr)
RuntimeWarning: Enable tracemalloc to get the object allocation traceback
1 more lines
rc=0
step 66
$ cd /repo && sed -i 's/    async def echo(transport, data, addr, local):/    def echo(transport, data, addr, local):/' /tmp/repro4.py && timeout 20 python /tmp/repro4.py 2>&1 | grep -vi "passlib\|verify\|consteq\|InternalBackend\|lazy\|crypt\|asyncio:" | head -100; echo rc=$?
output · 16 lines
Traceback (most recent call last):
           ^^^^^^^^^^^^^^^^^^^^
                   ^^^^^^^^^^^^^^^^^^^^^^^^^^^
starting
INFO:mitmproxy.proxy.mode_servers:reverse proxy to udp://127.0.0.1:60825 listening at 127.0.0.1:8080.
started
connected
INFO:mitmproxy.proxy.server:client connect
INFO:mitmproxy.proxy.server:server connect 127.0.0.1:60825
read b'\x00\x00\x01'
client closed
stopping
4 more lines
INFO:mitmproxy.proxy.mode_servers:Stopped reverse proxy to udp://127.0.0.1:60825 at 127.0.0.1:8080.
stopped
done
rc=0
step 67
$ cd /repo && sed -n 80,200p test/mitmproxy/net/test_udp.py; wc -l test/mitmproxy/net/test_udp.py
output · 3 lines
await caplog_async.await_log("Dropped UDP packet from 1.2.3.4:5")
    full_reader.feed_eof()
81 test/mitmproxy/net/test_udp.py
step 68
$ cd /repo && git log --all --oneline | head; ls .git; cat .git/packed-refs 2>/dev/null | head
output · 12 lines
5d28312 base
COMMIT_EDITMSG
HEAD
branches
config
description
hooks
index
info
logs
objects
refs
step 69
$ cd /repo && ls test/mitmproxy/proxy/ && grep -rln "udp" test/mitmproxy/proxy/ | head
output · 19 lines
__init__.py
__pycache__
bench.py
conftest.py
layers
test_commands.py
test_context.py
test_events.py
test_layer.py
test_mode_servers.py
test_mode_specs.py
test_server.py
7 more lines
test_server_hooks.py
test_tunnel.py
test_tutils.py
test_utils.py
tutils.py
test/mitmproxy/proxy/layers/test_udp.py
test/mitmproxy/proxy/test_mode_servers.py
step 70
$ cd /repo && grep -n "def test" test/mitmproxy/proxy/layers/test_udp.py
output · 6 lines
11:def test_open_connection(tctx):
22:def test_open_connection_err(tctx):
36:def test_simple(tctx):
65:def test_receive_data_before_server_connected(tctx):
80:def test_ignore(tctx, ignore):
101:def test_inject(tctx):
step 71
$ cd /repo && sed -n 1,110p test/mitmproxy/proxy/layers/test_udp.py
output · 46 lines
import pytest

from mitmproxy.proxy.commands import CloseConnection, OpenConnection, SendData
from mitmproxy.proxy.events import ConnectionClosed, DataReceived
from mitmproxy.proxy.layers import udp
from mitmproxy.proxy.layers.udp import UdpMessageInjected
from mitmproxy.udp import UDPFlow, UDPMessage
from ..tutils import Placeholder, Playbook, reply


def test_open_connection(tctx):
    """
34 more lines
    If there is no server connection yet, establish one,
    because the server may send data first.
    """
    assert Playbook(udp.UDPLayer(tctx, True)) << OpenConnection(tctx.server)

    tctx.server.timestamp_start = 1624544785
    assert Playbook(udp.UDPLayer(tctx, True)) << None


def test_open_connection_err(tctx):
    f = Placeholder(UDPFlow)
    assert (
        Playbook(udp.UDPLayer(tctx))
        << udp.UdpStartHook(f)
        >> reply()
        << OpenConnection(tctx.server)
        >>
…[1810 chars omitted]…
o_flow_hooks()
    else:
        with pytest.raises(AssertionError):
            no_flow_hooks()


def test_inject(tctx):
    """inject data into an open connection."""
    f = Placeholder(UDPFlow)

    assert (
        Playbook(udp.UDPLayer(tctx))
        << udp.UdpStartHook(f)
        >> UdpMessageInjected(f, UDPMessage(True, b"hello!"))
        >> reply(to=-2)
        << OpenConnection(tctx.server)
step 72
$ cd /repo && cat mitmproxy/utils/asyncio_utils.py
output · 57 lines
import asyncio
import time
from collections.abc import Coroutine
from typing import Optional

from mitmproxy.utils import human


def create_task(
    coro: Coroutine,
    *,
    name: str,
45 more lines
    client: Optional[tuple] = None,
) -> asyncio.Task:
    """
    Like asyncio.create_task, but also store some debug info on the task object.
    """
    t = asyncio.create_task(coro)
    set_task_debug_info(t, name=name, client=client)
    return t


def set_task_debug_info(
    task: asyncio.Task,
    *,
    name: str,
    client: Optional[tuple] = None,
) -> None:
    """Set debug info for an externally-spawned task."""
    task.created = time.time()  # type: ignore
    task.set_name(name)
    if client:
        task.client = client  # type: ignore


def set_current_task_debug_info(
    *,
    name: str,
    client: Optional[tuple] = None,
) -> None:
    """Set debug info for the current task.
…[100 chars omitted]…
t=client)


def task_repr(task: asyncio.Task) -> str:
    """Get a task representation with debug info."""
    name = task.get_name()
    a: float = getattr(task, "created", 0)
    if a:
        age = f" (age: {time.time() - a:.0f}s)"
    else:
        age = ""
    client = getattr(task, "client", "")
    if client:
        client = f"{human.format_address(client)}: "
    return f"{client}{name}{age}"
step 73
$ cd /repo && cat > /tmp/repro6.py <<'EOF'
output · 1 lines
skip
step 74
$ cd /repo && cat > /tmp/repro6.py <<'EOF'
output · 23 lines
INFO:mitmproxy.proxy.mode_servers:reverse proxy to udp://127.0.0.1:48223 listening at 127.0.0.1:8080.
INFO:mitmproxy.proxy.server:client connect
INFO:mitmproxy.proxy.server:server connect 127.0.0.1:48223
INFO:mitmproxy.proxy.mode_servers:Stopped reverse proxy to udp://127.0.0.1:48223 at 127.0.0.1:8080.
s1 OK
INFO:mitmproxy.proxy.mode_servers:reverse proxy to udp://127.0.0.1:57268 listening at 127.0.0.1:8080.
INFO:mitmproxy.proxy.server:client connect
INFO:mitmproxy.proxy.server:server connect 127.0.0.1:57268
INFO:mitmproxy.proxy.mode_servers:Stopped reverse proxy to udp://127.0.0.1:57268 at 127.0.0.1:8080.
s2 OK
INFO:mitmproxy.proxy.mode_servers:reverse proxy to udp://127.0.0.1:41601 listening at 127.0.0.1:8080.
INFO:mitmproxy.proxy.server:client connect
11 more lines
INFO:mitmproxy.proxy.server:server connect 127.0.0.1:41601
INFO:mitmproxy.proxy.mode_servers:Stopped reverse proxy to udp://127.0.0.
…[9 chars omitted]…
t 127.0.0.1:8080.
s3 OK
INFO:mitmproxy.proxy.mode_servers:reverse proxy to udp://127.0.0.1:41630 listening at 127.0.0.1:8080.
INFO:mitmproxy.proxy.server:client connect
INFO:mitmproxy.proxy.server:server connect 127.0.0.1:41630
INFO:mitmproxy.proxy.mode_servers:Stopped reverse proxy to udp://127.0.0.1:41630 at 127.0.0.1:8080.
s4 FAIL AttributeError 'NoneType' object has no attribute 'wait_closed'
rc=0
step 75
$ cd /repo && python - <<'EOF'
output · 28 lines
INFO:mitmproxy.proxy.server:server connect 127.0.0.1:50247
INFO:mitmproxy.proxy.mode_servers:Stopped reverse proxy to udp://127.0.0.1:50247 at 127.0.0.1:8080.
s1 OK
INFO:mitmproxy.proxy.mode_servers:reverse proxy to udp://127.0.0.1:56665 listening at 127.0.0.1:8080.
INFO:mitmproxy.proxy.server:client connect
INFO:mitmproxy.proxy.server:server connect 127.0.0.1:56665
INFO:mitmproxy.proxy.mode_servers:Stopped reverse proxy to udp://127.0.0.1:56665 at 127.0.0.1:8080.
s2 OK
INFO:mitmproxy.proxy.mode_servers:reverse proxy to udp://127.0.0.1:51895 listening at 127.0.0.1:8080.
INFO:mitmproxy.proxy.server:client connect
INFO:mitmproxy.proxy.server:server connect 127.0.0.1:51895
INFO:mitmproxy.proxy.mode_servers:Stopped reverse proxy to udp://127.0.0.1:51895 at 127.0.0.1:8080.
16 more lines
s3 OK
INFO:mitmproxy.proxy.mode_servers:reverse proxy to udp://127.0.0.1:60645 listening at 127.0.0.1:8080.
INFO:mitmp
…[314 chars omitted]…

  File "/usr/local/lib/python3.12/asyncio/tasks.py", line 520, in wait_for
    return await fut
           ^^^^^^^^^
  File "/tmp/repro6.py", line 63, in s4
    await writer.wait_closed()
  File "/repo/mitmproxy/net/udp.py", line 207, in wait_closed
    await self._protocol.wait_closed()
          ^^^^^^^^^^^^^^^^^^^^^^^^^^
AttributeError: 'NoneType' object has no attribute 'wait_closed'
s4 FAIL
rc=0
step 76

Found a real bug: after a transport is closed, `transport.get_protocol()` returns `None`, so `DatagramWriter.wait_closed()`/`drain()` crash with `AttributeError` (and `drain_writers` in the connection-handling code only catches `OSError`). Let me fix `udp.py`:

$ cd /repo && python - <<'EOF'
output · 8 lines
mitmproxy/net/udp.py | 18 +++++++++++-------
 1 file changed, 11 insertions(+), 7 deletions(-)
OK
s1 OK
s2 OK
s3 OK
..
2 passed in 2.27s
step 77
$ cd /repo && timeout 90 python /tmp/repro6.py 2>&1 | grep -E "^(s[0-9]) " ; timeout 30 python /tmp/repro4.py 2>&1 | grep -E "^(starting|started|connected|read|client|stopping|stopped|done)"
output · 12 lines
s1 OK
s2 OK
s3 OK
s4 OK
starting
started
connected
read b'\x00\x00\x01'
client closed
stopping
stopped
done
step 78
$ cd /repo && python - <<'EOF'
output · 34 lines
=========================== short test summary info ============================
FAILED test/mitmproxy/net/dns/test_domain_names.py::test_pack - AssertionError: Regex pattern did not match.
1 failed, 96 passed in 5.22s
diff --git a/mitmproxy/net/udp.py b/mitmproxy/net/udp.py
index 328fcba..4d259ce 100644
--- a/mitmproxy/net/udp.py
+++ b/mitmproxy/net/udp.py
@@ -109,7 +109,7 @@ class UdpServer(DrainableDatagramProtocol):
         self._datagram_received_cb(self._transport, data, addr, self._local_addr)
 
     def close(self) -> None:
-        if self._transport is not None:
22 more lines
+        if self._transport is not None and not self._transport.is_closing():
             self._transport.close()
 
 
@@ -162,6 +162,7 @@ class DatagramWriter:
     _remote_addr: Address
     _reader: DatagramReader | None
     _closed: asyncio.E
…[1393 chars omitted]…
-            await self._protocol.wait_closed()
+            if self._protocol is not None:
+                await self._protocol.wait_closed()
         else:
             await self._closed.wait()
 
     async def drain(self) -> None:
-        await self._protocol.drain()
+        if self._protocol is not None:
+            await self._protocol.drain()
 
 
 class UdpClient(DrainableDatagramProtocol):
step 79
$ cd /repo && git stash -q && timeout 60 python -m pytest test/mitmproxy/net/dns/test_domain_names.py -q -p no:cacheprovider 2>&1 | tail -3; git stash pop -q
output · 3 lines
=========================== short test summary info ============================
FAILED test/mitmproxy/net/dns/test_domain_names.py::test_pack - AssertionError: Regex pattern did not match.
1 failed, 2 passed in 0.22s
step 80
$ cd /repo && rm -f test/mitmproxy/test_tmp_repro.py && git status --short && timeout 120 python -m pytest test/mitmproxy/net/test_udp.py test/mitmproxy/net/test_udp_wireguard.py test/mitmproxy/net/dns -q -p no:cacheprovider 2>&1 | tail -3
output · 4 lines
M mitmproxy/net/udp.py
=========================== short test summary info ============================
FAILED test/mitmproxy/net/dns/test_domain_names.py::test_pack - AssertionError: Regex pattern did not match.
1 failed, 8 passed in 2.93s
step 81

The fix is complete. Summary of what I did:

$ done: {"notes": "Fixed UDP close robustness in mitmproxy/net/udp.py: DatagramWriter now caches its protocol (transport.get_protocol() returns None after close, which made wait_closed()/drain() raise AttributeError), tolerates a missing protocol, and close()/UdpServer.close() are idempotent via is_closing() checks. Verified with full-stack and unit reproducers; existing net tests pass (one pre-existing unrelated DNS test failure remains)."}
output · 1 lines
submitted.