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>
This commit is contained in:
devin-ai-integration[bot] 2026-10-07 12:01:12 -07:00 • committed by GitHub
parent 1c6f714187
commit 810a3106e1
No known key found for this signature in database
GPG key ID: B5690EEEBB952194

View file

@ -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"