From b4569632e5036179ce7b54465467c76aa91931b8 Mon Sep 17 00:00:00 2001 From: Avinash Kumar Deepak Date: Wed, 9 Sep 2026 15:18:49 +0530 Subject: [PATCH] fix(java): recover REQ socket after missing reply --- ConcoreJavaRuntimeCore.java | 6 +++++- TestConcoredockerApi.java | 30 ++++++++++++++++++++++++++++++ 2 files changed, 35 insertions(+), 1 deletion(-) diff --git a/ConcoreJavaRuntimeCore.java b/ConcoreJavaRuntimeCore.java index d371080..3120e1e 100644 --- a/ConcoreJavaRuntimeCore.java +++ b/ConcoreJavaRuntimeCore.java @@ -504,6 +504,10 @@ private class ZeroMQPort { ZeroMQPort(String portType, String address, int socketType) { ZMQ.Context ctx = getZmqContext(); this.socket = ctx.socket(socketType); + if (socketType == ZMQ.REQ) { + this.socket.setReqRelaxed(true); + this.socket.setReqCorrelate(true); + } this.socket.setReceiveTimeOut(2000); this.socket.setSendTimeOut(2000); this.socket.setLinger(0); @@ -771,4 +775,4 @@ Object parseKeyword() { } } } -} \ No newline at end of file +} diff --git a/TestConcoredockerApi.java b/TestConcoredockerApi.java index d62c8de..90e2085 100644 --- a/TestConcoredockerApi.java +++ b/TestConcoredockerApi.java @@ -2,6 +2,8 @@ import java.nio.file.Files; import java.nio.file.Path; import java.util.*; +import org.zeromq.ZMQ; +import org.zeromq.ZMQException; /** * Tests for concoredocker read(), write(), unchanged(), initVal() @@ -29,6 +31,7 @@ public static void main(String[] args) { testReadParseError(); testReadTraversalBlocked(); testWriteTraversalBlocked(); + testReqCanSendAfterMissingReply(); System.out.println("\n=== Results: " + passed + " passed, " + failed + " failed out of " + (passed + failed) + " tests ==="); if (failed > 0) { @@ -251,4 +254,31 @@ static void testWriteTraversalBlocked() { concoredocker.write(1, "../escape", Collections.singletonList((Object) 1.0), 0); check("write traversal blocked: no escaped file", false, Files.exists(tmp.resolve("escape"))); } + + static void testReqCanSendAfterMissingReply() { + ZMQ.Context context = ZMQ.context(1); + ZMQ.Socket peer = context.socket(ZMQ.REP); + peer.setReceiveTimeOut(2000); + peer.setLinger(0); + int port = peer.bindToRandomPort("tcp://127.0.0.1"); + + concoredocker.terminateZmq(); + try { + concoredocker.initZmqPort("request", "connect", "tcp://127.0.0.1:" + port, "REQ"); + concoredocker.write("request", "signal", Collections.singletonList((Object) 1.0), 0); + check("REQ peer receives first request", true, peer.recvStr() != null); + + boolean secondWriteSucceeded = true; + try { + concoredocker.write("request", "signal", Collections.singletonList((Object) 2.0), 0); + } catch (ZMQException e) { + secondWriteSucceeded = false; + } + check("REQ can send again when reply is missing", true, secondWriteSucceeded); + } finally { + concoredocker.terminateZmq(); + peer.close(); + context.term(); + } + } }