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", [