Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand All @@ -79,6 +80,7 @@ public class ConductorClient {
private final KeyManager[] keyManagers;
private final List<HeaderSupplier> headerSuppliers;
private final MetricsCollector metricsCollector;
private final boolean retransmitRequestBodies;

public static Builder<?> builder() {
return new Builder<>();
Expand All @@ -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();
Expand Down Expand Up @@ -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<Param> queryParams) {
Expand Down Expand Up @@ -603,6 +640,7 @@ public static class Builder<T extends Builder<T>> {
private Supplier<ObjectMapper> objectMapperSupplier = () -> new ObjectMapperProvider().getObjectMapper();
private final List<HeaderSupplier> headerSuppliers = new ArrayList<>();
MetricsCollector metricsCollector;
private boolean retransmitRequestBodies = true;

private boolean useEnvVariables = false;

Expand Down Expand Up @@ -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();
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,172 @@
/*
* Copyright 2026 Conductor Authors.
* <p>
* 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
* <p>
* http://www.apache.org/licenses/LICENSE-2.0
* <p>
* 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.
*
* <p>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<ConductorClient> 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();
}
}
Loading