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..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 @@ -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,47 @@ 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); + } + + @Override + public boolean isDuplex() { + return delegate.isDuplex(); + } + + // 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 +640,7 @@ public static class Builder> { private Supplier objectMapperSupplier = () -> new ObjectMapperProvider().getObjectMapper(); private final List headerSuppliers = new ArrayList<>(); MetricsCollector metricsCollector; + private boolean retransmitRequestBodies = true; private boolean useEnvVariables = false; @@ -656,6 +694,21 @@ public T proxy(Proxy proxy) { return self(); } + /** + * 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; + 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/RequestBodyRetransmissionTest.java b/conductor-client/src/test/java/com/netflix/conductor/client/http/RequestBodyRetransmissionTest.java new file mode 100644 index 000000000..1ddb8b1e8 --- /dev/null +++ b/conductor-client/src/test/java/com/netflix/conductor/client/http/RequestBodyRetransmissionTest.java @@ -0,0 +1,172 @@ +/* + * 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.io.InterruptedIOException; +import java.net.InetAddress; +import java.util.ArrayList; +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.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. {@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 + * 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; + private final List clients = new ArrayList<>(); + + @BeforeEach + void setUp() throws IOException { + server = new MockWebServer(); + server.start(InetAddress.getByName("127.0.0.1"), 0); + } + + @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(); + } + + @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 = 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()), + "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 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 = 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())); + + ConductorClientException e = assertThrows(ConductorClientException.class, () -> client.execute(postRequest())); + + assertEquals(2, server.getRequestCount(), "warm-up + severed attempt, no retransmit"); + 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 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(optInClient).isOneShot()); + } + + @Test + @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 = newClient( + ConductorClient.builder() + .basePath("http://multi-route.invalid:" + server.getPort() + "/api") + .retransmitRequestBodies(false) + .connectTimeout(250) + .configureOkHttp(b -> b.dns(twoRouteDns))); + + 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"; + } + + 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(); + } +}