From 810a3106e1d518eaf9791591fc794330d85ea651 Mon Sep 17 00:00:00 2001 From: "devin-ai-integration[bot]" <158243242+devin-ai-integration[bot]@users.noreply.github.com> Date: Wed, 7 Oct 2026 12:01:12 -0700 Subject: [PATCH] test(integration): hold the wire barrier until the test releases it or the wire closes (#45124) Co-authored-by: mateo-berri <277851410+mateo-berri@users.noreply.github.com> --- tests/integration/_support/wire.py | 17 +++++++++++++++-- 1 file changed, 15 insertions(+), 2 deletions(-) diff --git a/tests/integration/_support/wire.py b/tests/integration/_support/wire.py index 10bb7787cc9..d118da03daf 100644 --- a/tests/integration/_support/wire.py +++ b/tests/integration/_support/wire.py @@ -47,6 +47,13 @@ class Wire: return self.connected.qsize() +def _await_release(gate: threading.Event, closing: threading.Event) -> bool: + while not closing.is_set(): + if gate.wait(timeout=0.05): + return True + return gate.is_set() + + @contextmanager def wire_server( respond: Callable[[Request], Reply], @@ -60,6 +67,7 @@ def wire_server( errors: Final[SimpleQueue[Exception]] = SimpleQueue() disconnected: Final[SimpleQueue[str]] = SimpleQueue() connected: Final[SimpleQueue[str]] = SimpleQueue() + closing: Final = threading.Event() class Handler(BaseHTTPRequestHandler): protocol_version = "HTTP/1.1" @@ -108,8 +116,12 @@ def wire_server( break self.wfile.write(b"%x\r\n%s\r\n" % (len(chunk), chunk)) self.wfile.flush() - if index == 0 and reply.gate_after_first is not None: - assert reply.gate_after_first.wait(timeout=5), "Stream barrier was never released" + if ( + index == 0 + and reply.gate_after_first is not None + and not _await_release(reply.gate_after_first, closing) + ): + break if reply.pause_between_chunks and index + 1 < len(reply.chunks): time.sleep(reply.pause_between_chunks) else: @@ -151,6 +163,7 @@ def wire_server( connected, ) finally: + closing.set() server.shutdown() thread.join(timeout=6) assert not thread.is_alive(), "Owned HTTP server survived cleanup"