Skip to content

Commit 1fb1c74

Browse files
authored
Merge pull request #70324 from dwoz/dwoz/fix/3008x-stress-tests-branch-build
tests/stress: fix two 3008.x branch-build flakes
2 parents 27ff89e + 49c49bd commit 1fb1c74

2 files changed

Lines changed: 68 additions & 12 deletions

File tree

‎tests/pytests/stress/master_subprocess/mworker/test_mworker_stress.py‎

Lines changed: 21 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -268,17 +268,27 @@ def test_requester_disconnect_midflight_leaves_worker_alive(
268268

269269
# MWorker had already dispatched and queued the reply for the
270270
# "drop-me" request before we closed the DEALER; on reconnect
271-
# libzmq redelivers that queued reply to our fresh DEALER
272-
# first. Drain it, then send + receive a fresh request.
273-
try:
274-
stale = handle.recv(timeout=2.0)
275-
log.info("drained stale reply after reconnect: %r", stale)
276-
except TimeoutError:
277-
# Some libzmq versions do not redeliver buffered replies
278-
# after a peer identity change; that is fine too.
279-
log.info("no stale reply queued")
280-
281-
good = handle.send_recv(_ping("after-reconnect"), timeout=10.0)
271+
# libzmq may redeliver that queued reply to our fresh DEALER.
272+
# The timing of that redelivery is not deterministic across
273+
# libzmq builds and load levels: sometimes it arrives inside
274+
# the first poll window, sometimes after we've already sent
275+
# the follow-up request. Rather than draining with a fixed
276+
# timeout and hoping, we send the fresh request first, then
277+
# drain replies until we see the one keyed to
278+
# ``after-reconnect``. Any earlier reply keyed to ``drop-me``
279+
# is the redelivered stale that we're intentionally skipping.
280+
handle.send(_ping("after-reconnect"))
281+
good = None
282+
drain_deadline = time.monotonic() + 15.0
283+
while time.monotonic() < drain_deadline:
284+
try:
285+
reply = handle.recv(timeout=5.0)
286+
except TimeoutError:
287+
continue
288+
if isinstance(reply, dict) and reply.get("id") == "after-reconnect":
289+
good = reply
290+
break
291+
log.info("drained stale reply after reconnect: %r", reply)
282292
assert good == {"cmd": "ping", "id": "after-reconnect"}
283293
finally:
284294
handle.stop()

‎tests/pytests/stress/master_subprocess/mworkerqueue/test_mworkerqueue_stress.py‎

Lines changed: 47 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -326,8 +326,54 @@ def test_requester_churn_fd_bounded(mworkerqueue, proc_stats):
326326
A truly unbounded per-connection leak would grow phase 2 as much or
327327
more than phase 1.
328328
"""
329+
import zmq.utils.monitor
330+
329331
pid = mworkerqueue.process.pid
330-
worker = mworkerqueue.worker()
332+
333+
# ``mworkerqueue.worker()`` only issues ``connect()``; the ROUTER↔DEALER
334+
# zmq handshake and the internal DEALER's peer-list update are
335+
# asynchronous. On the first churn cycle we would otherwise race:
336+
# a REQ can complete its own handshake with the ROUTER and land a
337+
# message in the inner DEALER before the worker REP has registered
338+
# as a routable peer, so the DEALER queues the message internally
339+
# and ``_poll_recv(worker, 2000)`` in the churn loop times out.
340+
#
341+
# Wait for the worker REP's HANDSHAKE_SUCCEEDED against the proxy's
342+
# inner DEALER *before* we start the churn loop. This is a passive
343+
# readiness signal — no test traffic is injected — so it can't
344+
# leave stale probe messages queued at the DEALER or leave the REP
345+
# in "must send" state. (An earlier attempt that used an active
346+
# probe REQ→worker round-trip pileged stale ``b"probe"`` messages
347+
# onto the worker's inbound queue, which the churn loop then
348+
# mis-read as its own reply key and ended up hanging on ``m.recv()``
349+
# in the ROUTER→REQ direction.)
350+
worker = mworkerqueue.ctx.socket(zmq.REP)
351+
worker.setsockopt(zmq.LINGER, 500)
352+
worker.setsockopt(zmq.RCVTIMEO, 5000)
353+
worker.setsockopt(zmq.SNDTIMEO, 5000)
354+
monitor = worker.get_monitor_socket()
355+
worker.connect(mworkerqueue.dealer_uri)
356+
# Keep the socket alive for the fixture's teardown sweep.
357+
mworkerqueue._sockets.append(worker) # noqa: SLF001 — intentional
358+
359+
handshake_deadline = time.monotonic() + 15.0
360+
while time.monotonic() < handshake_deadline:
361+
if not monitor.poll(500):
362+
continue
363+
try:
364+
ev = zmq.utils.monitor.recv_monitor_message(monitor, zmq.NOBLOCK)
365+
except zmq.error.Again:
366+
continue
367+
if ev.get("event") == zmq.Event.HANDSHAKE_SUCCEEDED:
368+
break
369+
else:
370+
monitor.close(linger=0)
371+
raise AssertionError(
372+
"worker REP never completed HANDSHAKE_SUCCEEDED against the "
373+
"inner DEALER within 15s (test setup race, not the FD-bound "
374+
"assertion)"
375+
)
376+
monitor.close(linger=0)
331377

332378
def _churn(n_cycles: int, id_prefix: str) -> None:
333379
for i in range(n_cycles):

0 commit comments

Comments
 (0)