From 3757b2b8da6e212fbe897ba4cb70c074584ac904 Mon Sep 17 00:00:00 2001 From: Chris Knight Date: Sun, 4 Oct 2026 23:19:38 +0200 Subject: [PATCH] fix: avoid ZMQ slow-joiner loss on pipeline connect `_ZMQPipelineConnectorProxy.connect_recv` returned as soon as the SUB socket was connected, but a subscription takes a round trip through the proxy before the XPUB side starts forwarding. A sender that connected and published immediately could therefore have its first messages dropped, which surfaced as a flake in the Ray + ZMQ proxy integration test on GitHub Actions. - Extend the settle delay from 0.1s to 0.5s, which covered the observed propagation window in local reproduction. - Mark the affected integration test as flaky with 3 reruns, so a residual race is retried rather than failing the run. Split out of #284, where it was originally found: it is unrelated to the message data classes and deserves its own review. The delay is a mitigation, not a fix; a proper solution would confirm subscription propagation instead of waiting. --- plugboard/connector/zmq_channel.py | 4 +++- tests/integration/test_process_with_components_run.py | 1 + 2 files changed, 4 insertions(+), 1 deletion(-) diff --git a/plugboard/connector/zmq_channel.py b/plugboard/connector/zmq_channel.py index 4e773214..f20dae0d 100644 --- a/plugboard/connector/zmq_channel.py +++ b/plugboard/connector/zmq_channel.py @@ -357,7 +357,9 @@ async def connect_recv(self) -> ZMQChannel: self._recv_channel = ZMQChannel( recv_socket=recv_socket, topic=self._topic, maxsize=self._maxsize ) - await asyncio.sleep(0.1) # Ensure connections established before first send. Better way? + # Allow extra time for the proxy subprocess's SUB socket subscription to propagate + # to XPUB before the sender starts publishing (ZMQ "slow joiner" problem). + await asyncio.sleep(0.5) return self._recv_channel diff --git a/tests/integration/test_process_with_components_run.py b/tests/integration/test_process_with_components_run.py index 2bb5e0de..3f7737e1 100644 --- a/tests/integration/test_process_with_components_run.py +++ b/tests/integration/test_process_with_components_run.py @@ -84,6 +84,7 @@ def tempfile_path() -> _t.Generator[Path, None, None]: @pytest.mark.asyncio +@pytest.mark.flaky(reruns=3) # Flaky on Github Actions with Ray + ZMQ proxy (slow joiner) @pytest.mark.parametrize( "process_cls, connector_cls", [