From d2f185bea3f2b9dc4507fc95c63f988e40277432 Mon Sep 17 00:00:00 2001 From: Miguel Prieto Date: Sat, 3 Oct 2026 13:35:26 -0300 Subject: [PATCH 1/5] fix(client): make request bodies one-shot Stops OkHttp re-sending a request that was already delivered. Adds retransmitRequestBodies(true) to opt back in. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01JvEmz7CHxsyLEmhWwPWm7F --- .../client/http/ConductorClient.java | 47 +++++- .../http/ConnectionRetryBehaviourTest.java | 138 ++++++++++++++++++ 2 files changed, 183 insertions(+), 2 deletions(-) create mode 100644 conductor-client/src/test/java/com/netflix/conductor/client/http/ConnectionRetryBehaviourTest.java diff --git a/conductor-client/src/main/java/com/netflix/conductor/client/http/ConductorClient.java b/conductor-client/src/main/java/com/netflix/conductor/client/http/ConductorClient.java index aaa3073d3..205db5001 100644 --- a/conductor-client/src/main/java/com/netflix/conductor/client/http/ConductorClient.java +++ b/conductor-client/src/main/java/com/netflix/conductor/client/http/ConductorClient.java @@ -68,6 +68,7 @@ import okhttp3.RequestBody; import okhttp3.Response; import okhttp3.internal.http.HttpMethod; +import okio.BufferedSink; public class ConductorClient { private static final Logger LOGGER = LoggerFactory.getLogger(ConductorClient.class); @@ -79,6 +80,7 @@ public class ConductorClient { private final KeyManager[] keyManagers; private final List headerSuppliers; private final MetricsCollector metricsCollector; + private final boolean retransmitRequestBodies; public static Builder builder() { return new Builder<>(); @@ -95,6 +97,7 @@ protected ConductorClient(Builder builder) { this.keyManagers = builder.keyManagers; this.headerSuppliers = builder.headerSupplier(); this.metricsCollector = builder.metricsCollector; + this.retransmitRequestBodies = builder.retransmitRequestBodies; if (this.metricsCollector != null) { ApiClientMetrics apiClientMetrics = this.metricsCollector.getApiClientMetrics(); @@ -491,13 +494,42 @@ private RequestBody requestBody(String method, String contentType, Object body) return null; } + RequestBody requestBody; if (body == null && "DELETE".equals(method)) { return null; } else if (body == null) { - return RequestBody.create("", MediaType.parse(contentType)); + requestBody = RequestBody.create("", MediaType.parse(contentType)); + } else { + requestBody = serialize(contentType, body); } - return serialize(contentType, body); + return retransmitRequestBodies ? requestBody : oneShot(requestBody); + } + + // Wraps a request body so OkHttp will not retransmit it on a retried connection. + private static RequestBody oneShot(RequestBody delegate) { + return new RequestBody() { + @Override + public MediaType contentType() { + return delegate.contentType(); + } + + @Override + public long contentLength() throws IOException { + return delegate.contentLength(); + } + + @Override + public void writeTo(@NotNull BufferedSink sink) throws IOException { + delegate.writeTo(sink); + } + + // isOneShot() == true means: never re-send a body that was already transmitted. + @Override + public boolean isOneShot() { + return true; + } + }; } private HttpUrl buildUrl(String path, List queryParams) { @@ -603,6 +635,7 @@ public static class Builder> { private Supplier objectMapperSupplier = () -> new ObjectMapperProvider().getObjectMapper(); private final List headerSuppliers = new ArrayList<>(); MetricsCollector metricsCollector; + private boolean retransmitRequestBodies = false; private boolean useEnvVariables = false; @@ -656,6 +689,16 @@ public T proxy(Proxy proxy) { return self(); } + /** + * Pass {@code true} to restore OkHttp's stock behaviour of retransmitting a request + * body on a retried connection. The default ({@code false}) marks bodies one-shot, + * because a request that was already delivered to the server should not be sent again. + */ + public T retransmitRequestBodies(boolean retransmitRequestBodies) { + this.retransmitRequestBodies = retransmitRequestBodies; + return self(); + } + public T connectionPoolConfig(ConnectionPoolConfig config) { this.connectionPoolConfig = config; return self(); diff --git a/conductor-client/src/test/java/com/netflix/conductor/client/http/ConnectionRetryBehaviourTest.java b/conductor-client/src/test/java/com/netflix/conductor/client/http/ConnectionRetryBehaviourTest.java new file mode 100644 index 000000000..f802611b3 --- /dev/null +++ b/conductor-client/src/test/java/com/netflix/conductor/client/http/ConnectionRetryBehaviourTest.java @@ -0,0 +1,138 @@ +/* + * Copyright 2026 Conductor Authors. + *

+ * Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + *

+ * http://www.apache.org/licenses/LICENSE-2.0 + *

+ * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on + * an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the + * specific language governing permissions and limitations under the License. + */ +package com.netflix.conductor.client.http; + +import java.io.IOException; +import java.net.ConnectException; +import java.net.InetAddress; +import java.util.List; + +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; + +import com.netflix.conductor.client.exception.ConductorClientException; + +import okhttp3.Dns; +import okhttp3.mockwebserver.MockResponse; +import okhttp3.mockwebserver.MockWebServer; +import okhttp3.mockwebserver.SocketPolicy; + +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * {@code localhost} resolves to both an IPv4 and an IPv6 route on this machine, but MockWebServer + * only listens on one of them. A request that is fully delivered and then has its socket severed + * is, with {@code retransmitRequestBodies(true)}, retried on the other (unlistened) route and + * reported as a connect failure - even though the server already received it. With the SDK + * default (one-shot bodies), the caller instead sees the raw failure from the single, + * already-delivered attempt. + */ +class ConnectionRetryBehaviourTest { + + private MockWebServer server; + + @BeforeEach + void setUp() throws IOException { + server = new MockWebServer(); + server.start(); + } + + @AfterEach + void tearDown() throws IOException { + server.shutdown(); + Thread.interrupted(); + } + + @Test + @DisplayName("retransmitRequestBodies(true) delivers the request but reports it as a connect failure") + void retransmitEnabled_requestIsDeliveredThenRetried_reportsConnectFailure() { + server.enqueue(new MockResponse().setSocketPolicy(SocketPolicy.DISCONNECT_AFTER_REQUEST)); + + var client = new ConductorClient( + ConductorClient.builder() + .basePath(server.url("/api").toString()) + .retransmitRequestBodies(true)); + + var request = ConductorClientRequest.builder() + .method(ConductorClientRequest.Method.POST) + .path("/workflow") + .body("{\"name\":\"test\"}") + .build(); + + var e = assertThrows(ConductorClientException.class, () -> client.execute(request)); + + assertEquals(1, server.getRequestCount(), "request must have been delivered to the server"); + assertTrue( + e.getCause() instanceof ConnectException, + "expected a ConnectException from the retried (unreachable) route, got: " + e.getCause()); + assertTrue( + e.getCause().getMessage().contains("Failed to connect"), + "message should read like a connect failure, was: " + e.getCause().getMessage()); + } + + @Test + @DisplayName("default one-shot bodies report the raw single-attempt failure, not a connect failure") + void oneShotByDefault_requestIsDeliveredOnce_reportsRawFailure() { + server.enqueue(new MockResponse().setSocketPolicy(SocketPolicy.DISCONNECT_AFTER_REQUEST)); + + var client = new ConductorClient( + ConductorClient.builder().basePath(server.url("/api").toString())); + + var request = ConductorClientRequest.builder() + .method(ConductorClientRequest.Method.POST) + .path("/workflow") + .body("{\"name\":\"test\"}") + .build(); + + var e = assertThrows(ConductorClientException.class, () -> client.execute(request)); + + assertEquals(1, server.getRequestCount(), "request must have been delivered to the server"); + assertTrue( + !(e.getCause() instanceof ConnectException), + "must not be the connect-failure seen under retransmission, got: " + e.getCause()); + assertTrue( + e.getCause() instanceof IOException + && e.getCause().getMessage().contains("unexpected end of stream"), + "expected the raw single-attempt stream failure, got: " + e.getCause()); + } + + @Test + @DisplayName("a pre-send connection failure still falls back to another route") + void preSendConnectFailure_fallsBackToAnotherRoute() throws Exception { + // Fake Dns: dead loopback route first, the real MockWebServer route second. + server.enqueue(new MockResponse().setBody("{}")); + + Dns twoRouteDns = hostname -> List.of( + InetAddress.getByName("127.0.0.2"), InetAddress.getByName("127.0.0.1")); + + var client = new ConductorClient( + ConductorClient.builder() + .basePath("http://multi-route.invalid:" + server.getPort() + "/api") + .connectTimeout(1000) + .configureOkHttp(b -> b.dns(twoRouteDns))); + + var request = ConductorClientRequest.builder() + .method(ConductorClientRequest.Method.GET) + .path("/workflow") + .build(); + + assertDoesNotThrow(() -> client.execute(request), + "the dead first route must not fail the call: OkHttp should fall back to the live second route"); + assertEquals(1, server.getRequestCount(), "request must have reached the server via the fallback route"); + } +} From 4a0e1d28303773dbb51bf1395b3d300cf7ca130c Mon Sep 17 00:00:00 2001 From: Miguel Prieto Date: Sat, 3 Oct 2026 14:56:57 -0300 Subject: [PATCH 2/5] test: reproduce same-route REFUSED_STREAM retransmission against a single address RESET_STREAM_AT_START resets the stream before MockWebServer reads it, so it can't prove delivery. Instead, build a tiny HTTP/2 peer directly from OkHttp's internal Http2Connection that fully reads each request before refusing the first stream it sees. OkHttp's same-route retry then reopens a second stream on the very same pooled connection (not even a new TCP connection), resending the body, and the caller ends up seeing the second stream's own local CANCEL reset instead of the first attempt's REFUSED_STREAM - confirmed suppressed when the body is one-shot. Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_01JvEmz7CHxsyLEmhWwPWm7F --- .../http/ConnectionRetryBehaviourTest.java | 219 ++++++++++++++++++ 1 file changed, 219 insertions(+) diff --git a/conductor-client/src/test/java/com/netflix/conductor/client/http/ConnectionRetryBehaviourTest.java b/conductor-client/src/test/java/com/netflix/conductor/client/http/ConnectionRetryBehaviourTest.java index f802611b3..ad3452259 100644 --- a/conductor-client/src/test/java/com/netflix/conductor/client/http/ConnectionRetryBehaviourTest.java +++ b/conductor-client/src/test/java/com/netflix/conductor/client/http/ConnectionRetryBehaviourTest.java @@ -12,10 +12,16 @@ */ package com.netflix.conductor.client.http; +import java.io.Closeable; import java.io.IOException; import java.net.ConnectException; import java.net.InetAddress; +import java.net.ServerSocket; +import java.net.Socket; import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; @@ -25,12 +31,25 @@ import com.netflix.conductor.client.exception.ConductorClientException; import okhttp3.Dns; +import okhttp3.MediaType; +import okhttp3.OkHttpClient; +import okhttp3.Protocol; +import okhttp3.Request; +import okhttp3.RequestBody; +import okhttp3.internal.concurrent.TaskRunner; +import okhttp3.internal.http2.ErrorCode; +import okhttp3.internal.http2.Http2Connection; +import okhttp3.internal.http2.Http2Stream; +import okhttp3.internal.http2.StreamResetException; import okhttp3.mockwebserver.MockResponse; import okhttp3.mockwebserver.MockWebServer; import okhttp3.mockwebserver.SocketPolicy; +import okio.BufferedSink; +import okio.Okio; import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -135,4 +154,204 @@ void preSendConnectFailure_fallsBackToAnotherRoute() throws Exception { "the dead first route must not fail the call: OkHttp should fall back to the live second route"); assertEquals(1, server.getRequestCount(), "request must have reached the server via the fallback route"); } + + // Same-route retransmission via a genuine mid-exchange HTTP/2 REFUSED_STREAM reset, using a + // hand-rolled Http2Connection peer since MockWebServer can't reset a stream after reading it. + + @Test + @DisplayName("non-one-shot body: a request already read by the server is retransmitted on " + + "the same connection and the caller sees the second attempt's failure") + void refusedStreamAfterFullRead_bodyNotOneShot_retransmitsAndReportsSecondAttempt() throws Exception { + try (var h2Server = new ScriptedH2Server()) { + var client = new OkHttpClient.Builder() + .protocols(List.of(Protocol.H2_PRIOR_KNOWLEDGE)) + .callTimeout(5, TimeUnit.SECONDS) + .build(); + + var request = new Request.Builder() + .url("http://127.0.0.1:" + h2Server.port() + "/api/workflow") + .post(jsonBody()) + .build(); + + var e = assertThrows(IOException.class, () -> client.newCall(request).execute()); + + assertEquals( + List.of(15L, 15L), + h2Server.receivedBodySizes(), + "the 15-byte body must have been fully read by the server on both the first " + + "and the retried attempt - proof of retransmission"); + assertEquals(1, h2Server.connectionsAccepted(), + "the retry must reuse the same TCP connection - same route, not a fallback"); + + // Empirically a severed connection surfaces client-side as a local CANCEL reset, not a plain IOException. + assertInstanceOf(StreamResetException.class, e, "got: " + e); + assertEquals( + ErrorCode.CANCEL, + ((StreamResetException) e).errorCode, + "final error must describe the second attempt's local teardown, not the " + + "first attempt's REFUSED_STREAM, got: " + e); + } + } + + @Test + @DisplayName("one-shot body (the SDK default): the request is delivered once and is not " + + "retransmitted after the same REFUSED_STREAM reset") + void refusedStreamAfterFullRead_oneShotBody_deliversOnceWithoutRetransmission() throws Exception { + try (var h2Server = new ScriptedH2Server()) { + var client = new OkHttpClient.Builder() + .protocols(List.of(Protocol.H2_PRIOR_KNOWLEDGE)) + .callTimeout(5, TimeUnit.SECONDS) + .build(); + + var request = new Request.Builder() + .url("http://127.0.0.1:" + h2Server.port() + "/api/workflow") + .post(oneShot(jsonBody())) + .build(); + + var e = assertThrows(IOException.class, () -> client.newCall(request).execute()); + + assertEquals( + List.of(15L), + h2Server.receivedBodySizes(), + "the body must have been read by the server exactly once - no retransmission"); + assertEquals(1, h2Server.connectionsAccepted(), "no second connection must be opened"); + + assertInstanceOf(StreamResetException.class, e, "got: " + e); + assertEquals(ErrorCode.REFUSED_STREAM, ((StreamResetException) e).errorCode, "got: " + e); + } + } + + @Test + @DisplayName("empirically: SocketPolicy.RESET_STREAM_AT_START resets before reading, so it " + + "cannot prove the server received the request") + void resetStreamAtStart_doesNotProveTheRequestWasRead() throws Exception { + server.shutdown(); + server = new MockWebServer(); + server.setProtocols(List.of(Protocol.H2_PRIOR_KNOWLEDGE)); + server.start(InetAddress.getByName("127.0.0.1"), 0); + server.enqueue(new MockResponse().setSocketPolicy(SocketPolicy.RESET_STREAM_AT_START)); + + var client = new OkHttpClient.Builder() + .protocols(List.of(Protocol.H2_PRIOR_KNOWLEDGE)) + .callTimeout(5, TimeUnit.SECONDS) + .build(); + + var request = new Request.Builder() + .url("http://127.0.0.1:" + server.getPort() + "/api/workflow") + .post(jsonBody()) + .build(); + + assertThrows(IOException.class, () -> client.newCall(request).execute()); + + assertEquals(1, server.getRequestCount(), + "getRequestCount() is incremented for bookkeeping even though nothing was read"); + + var recorded = server.takeRequest(); + assertEquals("", recorded.getRequestLine(), + "the recorded request is an empty placeholder - the real request line was never parsed"); + assertEquals(0L, recorded.getBodySize(), + "the recorded request carries no body - it was never read off the wire"); + } + + private static RequestBody jsonBody() { + return RequestBody.create("{\"name\":\"test\"}", MediaType.parse("application/json")); + } + + private static RequestBody oneShot(RequestBody delegate) { + return new RequestBody() { + @Override + public MediaType contentType() { + return delegate.contentType(); + } + + @Override + public long contentLength() throws IOException { + return delegate.contentLength(); + } + + @Override + public void writeTo(BufferedSink sink) throws IOException { + delegate.writeTo(sink); + } + + @Override + public boolean isOneShot() { + return true; + } + }; + } + + /** Tiny Http2Connection-based peer: refuses the first stream it reads, severs any further one. */ + private static final class ScriptedH2Server implements Closeable { + + private final ServerSocket serverSocket; + private final AtomicInteger connectionsAccepted = new AtomicInteger(); + private final AtomicInteger streamsHandled = new AtomicInteger(); + private final List receivedBodySizes = new CopyOnWriteArrayList<>(); + private final Thread acceptThread; + private volatile boolean stopped; + + ScriptedH2Server() throws IOException { + serverSocket = new ServerSocket(0, 50, InetAddress.getByName("127.0.0.1")); + acceptThread = new Thread(this::acceptLoop, "scripted-h2-server"); + acceptThread.setDaemon(true); + acceptThread.start(); + } + + int port() { + return serverSocket.getLocalPort(); + } + + int connectionsAccepted() { + return connectionsAccepted.get(); + } + + List receivedBodySizes() { + return List.copyOf(receivedBodySizes); + } + + private void acceptLoop() { + while (!stopped) { + Socket socket; + try { + socket = serverSocket.accept(); + } catch (IOException e) { + return; + } + connectionsAccepted.incrementAndGet(); + try { + handleConnection(socket); + } catch (IOException ignored) { + // The client side of the test observes and asserts on the resulting failure. + } + } + } + + private void handleConnection(Socket socket) throws IOException { + var listener = new Http2Connection.Listener() { + @Override + public void onStream(Http2Stream stream) throws IOException { + stream.takeHeaders(); + var body = Okio.buffer(stream.getSource()); + receivedBodySizes.add((long) body.readByteString().size()); + if (streamsHandled.incrementAndGet() == 1) { + stream.close(ErrorCode.REFUSED_STREAM, null); + } else { + socket.close(); + } + } + }; + new Http2Connection.Builder(false, TaskRunner.INSTANCE) + .socket(socket) + .listener(listener) + .build() + .start(); + } + + @Override + public void close() throws IOException { + stopped = true; + serverSocket.close(); + } + } } From 547d59822c7d823de9c37655d813029a5f22421b Mon Sep 17 00:00:00 2001 From: Miguel Prieto Date: Sat, 3 Oct 2026 15:12:14 -0300 Subject: [PATCH 3/5] test: update the class javadoc to match what it now covers Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01JvEmz7CHxsyLEmhWwPWm7F --- .../client/http/ConnectionRetryBehaviourTest.java | 15 +++++++++------ 1 file changed, 9 insertions(+), 6 deletions(-) diff --git a/conductor-client/src/test/java/com/netflix/conductor/client/http/ConnectionRetryBehaviourTest.java b/conductor-client/src/test/java/com/netflix/conductor/client/http/ConnectionRetryBehaviourTest.java index ad3452259..91c78eda8 100644 --- a/conductor-client/src/test/java/com/netflix/conductor/client/http/ConnectionRetryBehaviourTest.java +++ b/conductor-client/src/test/java/com/netflix/conductor/client/http/ConnectionRetryBehaviourTest.java @@ -54,12 +54,15 @@ import static org.junit.jupiter.api.Assertions.assertTrue; /** - * {@code localhost} resolves to both an IPv4 and an IPv6 route on this machine, but MockWebServer - * only listens on one of them. A request that is fully delivered and then has its socket severed - * is, with {@code retransmitRequestBodies(true)}, retried on the other (unlistened) route and - * reported as a connect failure - even though the server already received it. With the SDK - * default (one-shot bodies), the caller instead sees the raw failure from the single, - * already-delivered attempt. + * A request body that is not one-shot can be sent twice: OkHttp retransmits it after a recoverable + * connection failure, and the caller is then told about the second attempt rather than the first. + * For a non-idempotent call that means the work happened and the error says it did not. + * + *

The decisive pair binds to a single address over HTTP/2, so there is provably one route and no + * fallback involved: the server reads the body in full, resets the stream, and OkHttp retries on + * the same connection. The remaining tests cover the multi-route shape, and check that a failure + * before the request is sent still falls back to the next address - the behaviour that must not be + * lost in exchange. */ class ConnectionRetryBehaviourTest { From 717ba844a1c1ae435434867210410fb6c160732c Mon Sep 17 00:00:00 2001 From: Miguel Prieto Date: Sat, 3 Oct 2026 15:37:32 -0300 Subject: [PATCH 4/5] fix(client): default retransmitRequestBodies to true, fix test coverage Flip the default so the one-shot protection is opt-in (false), keeping today's OkHttp retransmit behaviour for existing callers; update the setter javadoc to cover the 307/308 redirect side effect. Delegate isDuplex() in the one-shot RequestBody wrapper so CallServerInterceptor doesn't hang on a duplex delegate. Replace ConnectionRetryBehaviourTest with RequestBodyRetransmissionTest: drop the dual-stack-dependent test, the HTTP/2 tests that never exercised ConductorClient, and the resetStreamAtStart test that passed only via an uncaught NPE and a 5s timeout. Add a deterministic pooled-connection retransmission test pair and a direct one-shot contract pin. Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_01JvEmz7CHxsyLEmhWwPWm7F --- .../client/http/ConductorClient.java | 15 +- .../http/ConnectionRetryBehaviourTest.java | 360 ------------------ .../http/RequestBodyRetransmissionTest.java | 155 ++++++++ 3 files changed, 166 insertions(+), 364 deletions(-) delete mode 100644 conductor-client/src/test/java/com/netflix/conductor/client/http/ConnectionRetryBehaviourTest.java create mode 100644 conductor-client/src/test/java/com/netflix/conductor/client/http/RequestBodyRetransmissionTest.java diff --git a/conductor-client/src/main/java/com/netflix/conductor/client/http/ConductorClient.java b/conductor-client/src/main/java/com/netflix/conductor/client/http/ConductorClient.java index 205db5001..2a89145a7 100644 --- a/conductor-client/src/main/java/com/netflix/conductor/client/http/ConductorClient.java +++ b/conductor-client/src/main/java/com/netflix/conductor/client/http/ConductorClient.java @@ -524,6 +524,11 @@ public void writeTo(@NotNull BufferedSink sink) throws IOException { delegate.writeTo(sink); } + @Override + public boolean isDuplex() { + return delegate.isDuplex(); + } + // isOneShot() == true means: never re-send a body that was already transmitted. @Override public boolean isOneShot() { @@ -635,7 +640,7 @@ public static class Builder> { private Supplier objectMapperSupplier = () -> new ObjectMapperProvider().getObjectMapper(); private final List headerSuppliers = new ArrayList<>(); MetricsCollector metricsCollector; - private boolean retransmitRequestBodies = false; + private boolean retransmitRequestBodies = true; private boolean useEnvVariables = false; @@ -690,9 +695,11 @@ public T proxy(Proxy proxy) { } /** - * Pass {@code true} to restore OkHttp's stock behaviour of retransmitting a request - * body on a retried connection. The default ({@code false}) marks bodies one-shot, - * because a request that was already delivered to the server should not be sent again. + * Pass {@code false} to mark request bodies one-shot, so a request already delivered to + * the server is never sent again on a retried connection. Doing so also stops OkHttp from + * following 307 and 308 redirects for requests with a body (301/302/303 are unaffected, + * since those convert to a bodyless GET). Defaults to {@code true}, OkHttp's stock + * retransmitting behaviour. */ public T retransmitRequestBodies(boolean retransmitRequestBodies) { this.retransmitRequestBodies = retransmitRequestBodies; diff --git a/conductor-client/src/test/java/com/netflix/conductor/client/http/ConnectionRetryBehaviourTest.java b/conductor-client/src/test/java/com/netflix/conductor/client/http/ConnectionRetryBehaviourTest.java deleted file mode 100644 index 91c78eda8..000000000 --- a/conductor-client/src/test/java/com/netflix/conductor/client/http/ConnectionRetryBehaviourTest.java +++ /dev/null @@ -1,360 +0,0 @@ -/* - * Copyright 2026 Conductor Authors. - *

- * Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - *

- * http://www.apache.org/licenses/LICENSE-2.0 - *

- * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on - * an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the - * specific language governing permissions and limitations under the License. - */ -package com.netflix.conductor.client.http; - -import java.io.Closeable; -import java.io.IOException; -import java.net.ConnectException; -import java.net.InetAddress; -import java.net.ServerSocket; -import java.net.Socket; -import java.util.List; -import java.util.concurrent.CopyOnWriteArrayList; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicInteger; - -import org.junit.jupiter.api.AfterEach; -import org.junit.jupiter.api.BeforeEach; -import org.junit.jupiter.api.DisplayName; -import org.junit.jupiter.api.Test; - -import com.netflix.conductor.client.exception.ConductorClientException; - -import okhttp3.Dns; -import okhttp3.MediaType; -import okhttp3.OkHttpClient; -import okhttp3.Protocol; -import okhttp3.Request; -import okhttp3.RequestBody; -import okhttp3.internal.concurrent.TaskRunner; -import okhttp3.internal.http2.ErrorCode; -import okhttp3.internal.http2.Http2Connection; -import okhttp3.internal.http2.Http2Stream; -import okhttp3.internal.http2.StreamResetException; -import okhttp3.mockwebserver.MockResponse; -import okhttp3.mockwebserver.MockWebServer; -import okhttp3.mockwebserver.SocketPolicy; -import okio.BufferedSink; -import okio.Okio; - -import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; -import static org.junit.jupiter.api.Assertions.assertEquals; -import static org.junit.jupiter.api.Assertions.assertInstanceOf; -import static org.junit.jupiter.api.Assertions.assertThrows; -import static org.junit.jupiter.api.Assertions.assertTrue; - -/** - * A request body that is not one-shot can be sent twice: OkHttp retransmits it after a recoverable - * connection failure, and the caller is then told about the second attempt rather than the first. - * For a non-idempotent call that means the work happened and the error says it did not. - * - *

The decisive pair binds to a single address over HTTP/2, so there is provably one route and no - * fallback involved: the server reads the body in full, resets the stream, and OkHttp retries on - * the same connection. The remaining tests cover the multi-route shape, and check that a failure - * before the request is sent still falls back to the next address - the behaviour that must not be - * lost in exchange. - */ -class ConnectionRetryBehaviourTest { - - private MockWebServer server; - - @BeforeEach - void setUp() throws IOException { - server = new MockWebServer(); - server.start(); - } - - @AfterEach - void tearDown() throws IOException { - server.shutdown(); - Thread.interrupted(); - } - - @Test - @DisplayName("retransmitRequestBodies(true) delivers the request but reports it as a connect failure") - void retransmitEnabled_requestIsDeliveredThenRetried_reportsConnectFailure() { - server.enqueue(new MockResponse().setSocketPolicy(SocketPolicy.DISCONNECT_AFTER_REQUEST)); - - var client = new ConductorClient( - ConductorClient.builder() - .basePath(server.url("/api").toString()) - .retransmitRequestBodies(true)); - - var request = ConductorClientRequest.builder() - .method(ConductorClientRequest.Method.POST) - .path("/workflow") - .body("{\"name\":\"test\"}") - .build(); - - var e = assertThrows(ConductorClientException.class, () -> client.execute(request)); - - assertEquals(1, server.getRequestCount(), "request must have been delivered to the server"); - assertTrue( - e.getCause() instanceof ConnectException, - "expected a ConnectException from the retried (unreachable) route, got: " + e.getCause()); - assertTrue( - e.getCause().getMessage().contains("Failed to connect"), - "message should read like a connect failure, was: " + e.getCause().getMessage()); - } - - @Test - @DisplayName("default one-shot bodies report the raw single-attempt failure, not a connect failure") - void oneShotByDefault_requestIsDeliveredOnce_reportsRawFailure() { - server.enqueue(new MockResponse().setSocketPolicy(SocketPolicy.DISCONNECT_AFTER_REQUEST)); - - var client = new ConductorClient( - ConductorClient.builder().basePath(server.url("/api").toString())); - - var request = ConductorClientRequest.builder() - .method(ConductorClientRequest.Method.POST) - .path("/workflow") - .body("{\"name\":\"test\"}") - .build(); - - var e = assertThrows(ConductorClientException.class, () -> client.execute(request)); - - assertEquals(1, server.getRequestCount(), "request must have been delivered to the server"); - assertTrue( - !(e.getCause() instanceof ConnectException), - "must not be the connect-failure seen under retransmission, got: " + e.getCause()); - assertTrue( - e.getCause() instanceof IOException - && e.getCause().getMessage().contains("unexpected end of stream"), - "expected the raw single-attempt stream failure, got: " + e.getCause()); - } - - @Test - @DisplayName("a pre-send connection failure still falls back to another route") - void preSendConnectFailure_fallsBackToAnotherRoute() throws Exception { - // Fake Dns: dead loopback route first, the real MockWebServer route second. - server.enqueue(new MockResponse().setBody("{}")); - - Dns twoRouteDns = hostname -> List.of( - InetAddress.getByName("127.0.0.2"), InetAddress.getByName("127.0.0.1")); - - var client = new ConductorClient( - ConductorClient.builder() - .basePath("http://multi-route.invalid:" + server.getPort() + "/api") - .connectTimeout(1000) - .configureOkHttp(b -> b.dns(twoRouteDns))); - - var request = ConductorClientRequest.builder() - .method(ConductorClientRequest.Method.GET) - .path("/workflow") - .build(); - - assertDoesNotThrow(() -> client.execute(request), - "the dead first route must not fail the call: OkHttp should fall back to the live second route"); - assertEquals(1, server.getRequestCount(), "request must have reached the server via the fallback route"); - } - - // Same-route retransmission via a genuine mid-exchange HTTP/2 REFUSED_STREAM reset, using a - // hand-rolled Http2Connection peer since MockWebServer can't reset a stream after reading it. - - @Test - @DisplayName("non-one-shot body: a request already read by the server is retransmitted on " - + "the same connection and the caller sees the second attempt's failure") - void refusedStreamAfterFullRead_bodyNotOneShot_retransmitsAndReportsSecondAttempt() throws Exception { - try (var h2Server = new ScriptedH2Server()) { - var client = new OkHttpClient.Builder() - .protocols(List.of(Protocol.H2_PRIOR_KNOWLEDGE)) - .callTimeout(5, TimeUnit.SECONDS) - .build(); - - var request = new Request.Builder() - .url("http://127.0.0.1:" + h2Server.port() + "/api/workflow") - .post(jsonBody()) - .build(); - - var e = assertThrows(IOException.class, () -> client.newCall(request).execute()); - - assertEquals( - List.of(15L, 15L), - h2Server.receivedBodySizes(), - "the 15-byte body must have been fully read by the server on both the first " - + "and the retried attempt - proof of retransmission"); - assertEquals(1, h2Server.connectionsAccepted(), - "the retry must reuse the same TCP connection - same route, not a fallback"); - - // Empirically a severed connection surfaces client-side as a local CANCEL reset, not a plain IOException. - assertInstanceOf(StreamResetException.class, e, "got: " + e); - assertEquals( - ErrorCode.CANCEL, - ((StreamResetException) e).errorCode, - "final error must describe the second attempt's local teardown, not the " - + "first attempt's REFUSED_STREAM, got: " + e); - } - } - - @Test - @DisplayName("one-shot body (the SDK default): the request is delivered once and is not " - + "retransmitted after the same REFUSED_STREAM reset") - void refusedStreamAfterFullRead_oneShotBody_deliversOnceWithoutRetransmission() throws Exception { - try (var h2Server = new ScriptedH2Server()) { - var client = new OkHttpClient.Builder() - .protocols(List.of(Protocol.H2_PRIOR_KNOWLEDGE)) - .callTimeout(5, TimeUnit.SECONDS) - .build(); - - var request = new Request.Builder() - .url("http://127.0.0.1:" + h2Server.port() + "/api/workflow") - .post(oneShot(jsonBody())) - .build(); - - var e = assertThrows(IOException.class, () -> client.newCall(request).execute()); - - assertEquals( - List.of(15L), - h2Server.receivedBodySizes(), - "the body must have been read by the server exactly once - no retransmission"); - assertEquals(1, h2Server.connectionsAccepted(), "no second connection must be opened"); - - assertInstanceOf(StreamResetException.class, e, "got: " + e); - assertEquals(ErrorCode.REFUSED_STREAM, ((StreamResetException) e).errorCode, "got: " + e); - } - } - - @Test - @DisplayName("empirically: SocketPolicy.RESET_STREAM_AT_START resets before reading, so it " - + "cannot prove the server received the request") - void resetStreamAtStart_doesNotProveTheRequestWasRead() throws Exception { - server.shutdown(); - server = new MockWebServer(); - server.setProtocols(List.of(Protocol.H2_PRIOR_KNOWLEDGE)); - server.start(InetAddress.getByName("127.0.0.1"), 0); - server.enqueue(new MockResponse().setSocketPolicy(SocketPolicy.RESET_STREAM_AT_START)); - - var client = new OkHttpClient.Builder() - .protocols(List.of(Protocol.H2_PRIOR_KNOWLEDGE)) - .callTimeout(5, TimeUnit.SECONDS) - .build(); - - var request = new Request.Builder() - .url("http://127.0.0.1:" + server.getPort() + "/api/workflow") - .post(jsonBody()) - .build(); - - assertThrows(IOException.class, () -> client.newCall(request).execute()); - - assertEquals(1, server.getRequestCount(), - "getRequestCount() is incremented for bookkeeping even though nothing was read"); - - var recorded = server.takeRequest(); - assertEquals("", recorded.getRequestLine(), - "the recorded request is an empty placeholder - the real request line was never parsed"); - assertEquals(0L, recorded.getBodySize(), - "the recorded request carries no body - it was never read off the wire"); - } - - private static RequestBody jsonBody() { - return RequestBody.create("{\"name\":\"test\"}", MediaType.parse("application/json")); - } - - private static RequestBody oneShot(RequestBody delegate) { - return new RequestBody() { - @Override - public MediaType contentType() { - return delegate.contentType(); - } - - @Override - public long contentLength() throws IOException { - return delegate.contentLength(); - } - - @Override - public void writeTo(BufferedSink sink) throws IOException { - delegate.writeTo(sink); - } - - @Override - public boolean isOneShot() { - return true; - } - }; - } - - /** Tiny Http2Connection-based peer: refuses the first stream it reads, severs any further one. */ - private static final class ScriptedH2Server implements Closeable { - - private final ServerSocket serverSocket; - private final AtomicInteger connectionsAccepted = new AtomicInteger(); - private final AtomicInteger streamsHandled = new AtomicInteger(); - private final List receivedBodySizes = new CopyOnWriteArrayList<>(); - private final Thread acceptThread; - private volatile boolean stopped; - - ScriptedH2Server() throws IOException { - serverSocket = new ServerSocket(0, 50, InetAddress.getByName("127.0.0.1")); - acceptThread = new Thread(this::acceptLoop, "scripted-h2-server"); - acceptThread.setDaemon(true); - acceptThread.start(); - } - - int port() { - return serverSocket.getLocalPort(); - } - - int connectionsAccepted() { - return connectionsAccepted.get(); - } - - List receivedBodySizes() { - return List.copyOf(receivedBodySizes); - } - - private void acceptLoop() { - while (!stopped) { - Socket socket; - try { - socket = serverSocket.accept(); - } catch (IOException e) { - return; - } - connectionsAccepted.incrementAndGet(); - try { - handleConnection(socket); - } catch (IOException ignored) { - // The client side of the test observes and asserts on the resulting failure. - } - } - } - - private void handleConnection(Socket socket) throws IOException { - var listener = new Http2Connection.Listener() { - @Override - public void onStream(Http2Stream stream) throws IOException { - stream.takeHeaders(); - var body = Okio.buffer(stream.getSource()); - receivedBodySizes.add((long) body.readByteString().size()); - if (streamsHandled.incrementAndGet() == 1) { - stream.close(ErrorCode.REFUSED_STREAM, null); - } else { - socket.close(); - } - } - }; - new Http2Connection.Builder(false, TaskRunner.INSTANCE) - .socket(socket) - .listener(listener) - .build() - .start(); - } - - @Override - public void close() throws IOException { - stopped = true; - serverSocket.close(); - } - } -} diff --git a/conductor-client/src/test/java/com/netflix/conductor/client/http/RequestBodyRetransmissionTest.java b/conductor-client/src/test/java/com/netflix/conductor/client/http/RequestBodyRetransmissionTest.java new file mode 100644 index 000000000..8332b6fa1 --- /dev/null +++ b/conductor-client/src/test/java/com/netflix/conductor/client/http/RequestBodyRetransmissionTest.java @@ -0,0 +1,155 @@ +/* + * Copyright 2026 Conductor Authors. + *

+ * Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + *

+ * http://www.apache.org/licenses/LICENSE-2.0 + *

+ * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on + * an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the + * specific language governing permissions and limitations under the License. + */ +package com.netflix.conductor.client.http; + +import java.io.IOException; +import java.net.InetAddress; +import java.util.List; +import java.util.Map; + +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; + +import com.netflix.conductor.client.exception.ConductorClientException; + +import okhttp3.Dns; +import okhttp3.RequestBody; +import okhttp3.mockwebserver.MockResponse; +import okhttp3.mockwebserver.MockWebServer; +import okhttp3.mockwebserver.SocketPolicy; + +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * A request body that is not one-shot can be sent twice: OkHttp retransmits it after a + * recoverable connection failure, and the caller is then told about the second attempt rather + * than the first. For a non-idempotent call that means the work happened and the error says it + * did not. + * + *

The decisive pair below makes one successful call first, so the connection is pooled, then + * severs the next one: a connection taken from the pool retries on its own address, which is the + * real-world shape this SDK must handle. The remaining tests pin the one-shot contract directly + * and check that a failure before the request is sent still falls back to another route - the + * behaviour that must not be lost in exchange. + */ +class RequestBodyRetransmissionTest { + + private MockWebServer server; + + @BeforeEach + void setUp() throws IOException { + server = new MockWebServer(); + server.start(InetAddress.getByName("127.0.0.1"), 0); + } + + @AfterEach + void tearDown() throws IOException { + server.shutdown(); + } + + @Test + @DisplayName("default (retransmitRequestBodies true): a body severed on a pooled connection " + + "is retransmitted and the call succeeds") + void retransmitEnabled_bodyIsRetransmittedOnPooledConnection_callSucceeds() { + server.enqueue(new MockResponse().setBody("{}")); // warm-up: pools the connection + server.enqueue(new MockResponse().setSocketPolicy(SocketPolicy.DISCONNECT_AFTER_REQUEST)); + server.enqueue(new MockResponse().setBody("{}")); // served only if retransmitted + + var client = new ConductorClient( + ConductorClient.builder().basePath(basePath()).retransmitRequestBodies(true)); + + assertDoesNotThrow(() -> client.execute(warmupRequest())); + + assertDoesNotThrow(() -> client.execute(postRequest()), + "retransmission must hide the severed attempt from the caller"); + assertEquals(3, server.getRequestCount(), "warm-up + severed attempt + retransmit"); + } + + @Test + @DisplayName("opt-in (retransmitRequestBodies(false)): a body severed on a pooled connection " + + "is not retransmitted and the call fails") + void optOut_bodyIsNotRetransmittedOnPooledConnection_callFails() { + server.enqueue(new MockResponse().setBody("{}")); // warm-up: pools the connection + server.enqueue(new MockResponse().setSocketPolicy(SocketPolicy.DISCONNECT_AFTER_REQUEST)); + server.enqueue(new MockResponse().setBody("{}")); // must never be reached + + var client = new ConductorClient( + ConductorClient.builder().basePath(basePath()).retransmitRequestBodies(false)); + + assertDoesNotThrow(() -> client.execute(warmupRequest())); + + var e = assertThrows(ConductorClientException.class, () -> client.execute(postRequest())); + + assertEquals(2, server.getRequestCount(), "warm-up + severed attempt, no retransmit"); + System.out.println("optOut_bodyIsNotRetransmittedOnPooledConnection_callFails threw: " + e); + } + + @Test + @DisplayName("contract pin: the built request body is one-shot only when opted out") + void requestBody_isOneShot_onlyWhenRetransmitDisabled() { + var defaultClient = new ConductorClient(ConductorClient.builder().basePath(basePath())); + var optOutClient = new ConductorClient( + ConductorClient.builder().basePath(basePath()).retransmitRequestBodies(false)); + + assertFalse(builtRequestBody(defaultClient).isOneShot()); + assertTrue(builtRequestBody(optOutClient).isOneShot()); + } + + @Test + @DisplayName("a pre-send connection failure still falls back to another route") + void preSendConnectFailure_fallsBackToAnotherRoute() { + server.enqueue(new MockResponse().setBody("{}")); + + Dns twoRouteDns = hostname -> List.of( + InetAddress.getByName("127.0.0.2"), InetAddress.getByName("127.0.0.1")); + + var client = new ConductorClient( + ConductorClient.builder() + .basePath("http://multi-route.invalid:" + server.getPort() + "/api") + .connectTimeout(1000) + .configureOkHttp(b -> b.dns(twoRouteDns))); + + assertDoesNotThrow(() -> client.execute(warmupRequest()), + "the dead first route must not fail the call: OkHttp should fall back to the live second route"); + assertEquals(1, server.getRequestCount(), "request must have reached the server via the fallback route"); + } + + private String basePath() { + return "http://127.0.0.1:" + server.getPort() + "/api"; + } + + private static ConductorClientRequest warmupRequest() { + return ConductorClientRequest.builder() + .method(ConductorClientRequest.Method.GET) + .path("/workflow") + .build(); + } + + private static ConductorClientRequest postRequest() { + return ConductorClientRequest.builder() + .method(ConductorClientRequest.Method.POST) + .path("/workflow") + .body("{\"name\":\"test\"}") + .build(); + } + + private static RequestBody builtRequestBody(ConductorClient client) { + return client.buildRequest("POST", "/workflow", List.of(), List.of(), Map.of(), "{\"name\":\"test\"}").body(); + } +} From 759a8d3ad4a994743f5d89b41675e7e2a30feee6 Mon Sep 17 00:00:00 2001 From: Miguel Prieto Date: Sat, 3 Oct 2026 15:49:50 -0300 Subject: [PATCH 5/5] fix(client): address review on one-shot default flip Fix preSendConnectFailure_fallsBackToAnotherRoute: it was built with the default (retransmitRequestBodies true, nothing wrapped) and sent a bodyless GET, so it only pinned stock OkHttp route fallback. Rebuild it with retransmitRequestBodies(false) and a POST body so it actually pins that pre-send fallback survives the one-shot wrapper; drop connectTimeout to 250ms since the dead route blackholes rather than refusing. Expand the setter javadoc to cover 408 and 421 follow-ups alongside 307/308, note what the caller sees (a ConductorClientException carrying the blocked status), and soften "already delivered" to "once transmission has begun" to match requestSendStarted. Settle the opt-in/opt-out naming on "retransmitRequestBodies(false) opts in to the protection" throughout. Replace the opt-in test's debug println with an assertion that the failure cause is an IOException and not an InterruptedIOException, guarding against a call-timeout masquerading as the expected failure. Track clients created in each test and evict their connection pools in tearDown. Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_01JvEmz7CHxsyLEmhWwPWm7F --- .../client/http/ConductorClient.java | 13 +++-- .../http/RequestBodyRetransmissionTest.java | 55 ++++++++++++------- 2 files changed, 44 insertions(+), 24 deletions(-) diff --git a/conductor-client/src/main/java/com/netflix/conductor/client/http/ConductorClient.java b/conductor-client/src/main/java/com/netflix/conductor/client/http/ConductorClient.java index 2a89145a7..089848454 100644 --- a/conductor-client/src/main/java/com/netflix/conductor/client/http/ConductorClient.java +++ b/conductor-client/src/main/java/com/netflix/conductor/client/http/ConductorClient.java @@ -695,11 +695,14 @@ public T proxy(Proxy proxy) { } /** - * Pass {@code false} to mark request bodies one-shot, so a request already delivered to - * the server is never sent again on a retried connection. Doing so also stops OkHttp from - * following 307 and 308 redirects for requests with a body (301/302/303 are unaffected, - * since those convert to a bodyless GET). Defaults to {@code true}, OkHttp's stock - * retransmitting behaviour. + * Pass {@code false} to mark request bodies one-shot, opting in to the protection: once + * transmission of a body has begun, it is never sent again on a retried connection. This + * also stops OkHttp replaying the request on 307/308 redirects, 408 responses and 421 + * misdirected-request responses (301/302/303 are unaffected, since those convert to a + * bodyless GET). A blocked 307/308 surfaces as a {@link ConductorClientException} + * carrying that status code rather than a transparent redirect, since + * {@link ConductorClient#handleResponse} treats 3xx as unsuccessful. Defaults to + * {@code true}, OkHttp's stock retransmit/replay behaviour. */ public T retransmitRequestBodies(boolean retransmitRequestBodies) { this.retransmitRequestBodies = retransmitRequestBodies; diff --git a/conductor-client/src/test/java/com/netflix/conductor/client/http/RequestBodyRetransmissionTest.java b/conductor-client/src/test/java/com/netflix/conductor/client/http/RequestBodyRetransmissionTest.java index 8332b6fa1..1ddb8b1e8 100644 --- a/conductor-client/src/test/java/com/netflix/conductor/client/http/RequestBodyRetransmissionTest.java +++ b/conductor-client/src/test/java/com/netflix/conductor/client/http/RequestBodyRetransmissionTest.java @@ -13,7 +13,9 @@ package com.netflix.conductor.client.http; import java.io.IOException; +import java.io.InterruptedIOException; import java.net.InetAddress; +import java.util.ArrayList; import java.util.List; import java.util.Map; @@ -33,6 +35,7 @@ import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -40,7 +43,7 @@ * A request body that is not one-shot can be sent twice: OkHttp retransmits it after a * recoverable connection failure, and the caller is then told about the second attempt rather * than the first. For a non-idempotent call that means the work happened and the error says it - * did not. + * did not. {@code retransmitRequestBodies(false)} opts in to the one-shot protection. * *

The decisive pair below makes one successful call first, so the connection is pooled, then * severs the next one: a connection taken from the pool retries on its own address, which is the @@ -51,6 +54,7 @@ class RequestBodyRetransmissionTest { private MockWebServer server; + private final List clients = new ArrayList<>(); @BeforeEach void setUp() throws IOException { @@ -60,6 +64,8 @@ void setUp() throws IOException { @AfterEach void tearDown() throws IOException { + // Pools are per-client here, but evict explicitly so no pooled connection outlives its server. + clients.forEach(c -> c.okHttpClient.connectionPool().evictAll()); server.shutdown(); } @@ -71,9 +77,9 @@ void retransmitEnabled_bodyIsRetransmittedOnPooledConnection_callSucceeds() { server.enqueue(new MockResponse().setSocketPolicy(SocketPolicy.DISCONNECT_AFTER_REQUEST)); server.enqueue(new MockResponse().setBody("{}")); // served only if retransmitted - var client = new ConductorClient( - ConductorClient.builder().basePath(basePath()).retransmitRequestBodies(true)); + var client = newClient(ConductorClient.builder().basePath(basePath()).retransmitRequestBodies(true)); + // Must run to completion (response body drained) or the connection is never pooled and the retry never fires. assertDoesNotThrow(() -> client.execute(warmupRequest())); assertDoesNotThrow(() -> client.execute(postRequest()), @@ -84,52 +90,63 @@ void retransmitEnabled_bodyIsRetransmittedOnPooledConnection_callSucceeds() { @Test @DisplayName("opt-in (retransmitRequestBodies(false)): a body severed on a pooled connection " + "is not retransmitted and the call fails") - void optOut_bodyIsNotRetransmittedOnPooledConnection_callFails() { + void optIn_bodyIsNotRetransmittedOnPooledConnection_callFails() { server.enqueue(new MockResponse().setBody("{}")); // warm-up: pools the connection server.enqueue(new MockResponse().setSocketPolicy(SocketPolicy.DISCONNECT_AFTER_REQUEST)); server.enqueue(new MockResponse().setBody("{}")); // must never be reached - var client = new ConductorClient( - ConductorClient.builder().basePath(basePath()).retransmitRequestBodies(false)); + var client = newClient(ConductorClient.builder().basePath(basePath()).retransmitRequestBodies(false)); + // Must run to completion (response body drained) or the connection is never pooled and the retry never fires. assertDoesNotThrow(() -> client.execute(warmupRequest())); - var e = assertThrows(ConductorClientException.class, () -> client.execute(postRequest())); + ConductorClientException e = assertThrows(ConductorClientException.class, () -> client.execute(postRequest())); assertEquals(2, server.getRequestCount(), "warm-up + severed attempt, no retransmit"); - System.out.println("optOut_bodyIsNotRetransmittedOnPooledConnection_callFails threw: " + e); + assertInstanceOf(IOException.class, e.getCause(), "got: " + e.getCause()); + assertFalse(e.getCause() instanceof InterruptedIOException, + "must be the raw severed-socket failure, not a call timeout masquerading as it: " + e.getCause()); } @Test - @DisplayName("contract pin: the built request body is one-shot only when opted out") - void requestBody_isOneShot_onlyWhenRetransmitDisabled() { - var defaultClient = new ConductorClient(ConductorClient.builder().basePath(basePath())); - var optOutClient = new ConductorClient( - ConductorClient.builder().basePath(basePath()).retransmitRequestBodies(false)); + @DisplayName("contract pin: the built request body is one-shot only when opted in " + + "(retransmitRequestBodies(false))") + void requestBody_isOneShot_onlyWhenOptedIn() { + var defaultClient = newClient(ConductorClient.builder().basePath(basePath())); + var optInClient = newClient(ConductorClient.builder().basePath(basePath()).retransmitRequestBodies(false)); assertFalse(builtRequestBody(defaultClient).isOneShot()); - assertTrue(builtRequestBody(optOutClient).isOneShot()); + assertTrue(builtRequestBody(optInClient).isOneShot()); } @Test - @DisplayName("a pre-send connection failure still falls back to another route") + @DisplayName("a pre-send connection failure still falls back to another route, even for a " + + "one-shot POST body") void preSendConnectFailure_fallsBackToAnotherRoute() { server.enqueue(new MockResponse().setBody("{}")); Dns twoRouteDns = hostname -> List.of( InetAddress.getByName("127.0.0.2"), InetAddress.getByName("127.0.0.1")); - var client = new ConductorClient( + var client = newClient( ConductorClient.builder() .basePath("http://multi-route.invalid:" + server.getPort() + "/api") - .connectTimeout(1000) + .retransmitRequestBodies(false) + .connectTimeout(250) .configureOkHttp(b -> b.dns(twoRouteDns))); - assertDoesNotThrow(() -> client.execute(warmupRequest()), - "the dead first route must not fail the call: OkHttp should fall back to the live second route"); + assertDoesNotThrow(() -> client.execute(postRequest()), + "a one-shot POST must still fall back pre-send: recover() never consults " + + "requestIsOneShot before the body starts sending"); assertEquals(1, server.getRequestCount(), "request must have reached the server via the fallback route"); } + private ConductorClient newClient(ConductorClient.Builder builder) { + ConductorClient client = new ConductorClient(builder); + clients.add(client); + return client; + } + private String basePath() { return "http://127.0.0.1:" + server.getPort() + "/api"; }