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 @@ -21,6 +21,7 @@
import java.io.InputStream;
import java.io.PipedInputStream;
import java.io.PipedOutputStream;
import java.io.UncheckedIOException;
import java.lang.invoke.MethodHandles;
import java.net.CookieHandler;
import java.net.InetSocketAddress;
Expand All @@ -37,12 +38,9 @@
import java.util.HashMap;
import java.util.Locale;
import java.util.Map;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Future;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.TimeUnit;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
import javax.net.ssl.SSLContext;
Expand Down Expand Up @@ -78,65 +76,70 @@ public class HttpJdkSolrClient extends HttpSolrClient {

protected HttpClient httpClient;

/**
* Executor used to stream (produce) request bodies into the pipe consumed by the JDK HttpClient.
* This is the "producer" side and may be supplied by the caller.
*/
protected ExecutorService executor;

/** Dedicated executor handed to the JDK HttpClient */
protected ExecutorService httpClientExecutor;

private boolean forceHttp11;

private final boolean shutdownExecutor;

protected HttpJdkSolrClient(String serverBaseUrl, HttpJdkSolrClient.Builder builder) {
super(serverBaseUrl, builder);
HttpClient.Builder b = HttpClient.newBuilder();
HttpClient.Builder httpClientBuilder = HttpClient.newBuilder();

HttpClient.Redirect followRedirects =
Boolean.TRUE.equals(builder.getFollowRedirects())
? HttpClient.Redirect.NORMAL
: HttpClient.Redirect.NEVER;
b.followRedirects(followRedirects);
httpClientBuilder.followRedirects(followRedirects);

b.connectTimeout(Duration.of(builder.getConnectionTimeoutMillis(), ChronoUnit.MILLIS));
httpClientBuilder.connectTimeout(
Duration.of(builder.getConnectionTimeoutMillis(), ChronoUnit.MILLIS));
// note: idle timeout isn't used for the JDK client
// note: request timeout is set per request

if (builder.sslContext != null) {
b.sslContext(builder.sslContext);
httpClientBuilder.sslContext(builder.sslContext);
}

if (builder.getExecutor() != null) {
this.executor = builder.getExecutor();
this.shutdownExecutor = false;
} else {
BlockingQueue<Runnable> queue = new LinkedBlockingQueue<>(1024);
this.executor =
new ExecutorUtil.MDCAwareThreadPoolExecutor(
4,
256,
60,
TimeUnit.SECONDS,
queue,
new SolrNamedThreadFactory(this.getClass().getSimpleName()));
ExecutorUtil.newMDCAwareCachedThreadPool(
new SolrNamedThreadFactory(this.getClass().getSimpleName() + "-reqBody"));
this.shutdownExecutor = true;
}
b.executor(this.executor);
this.httpClientExecutor =
ExecutorUtil.newMDCAwareCachedThreadPool(
new SolrNamedThreadFactory(this.getClass().getSimpleName() + "-http"));
httpClientBuilder.executor(this.httpClientExecutor);

if (builder.shouldUseHttp1_1()) {
this.forceHttp11 = true;
b.version(HttpClient.Version.HTTP_1_1);
httpClientBuilder.version(HttpClient.Version.HTTP_1_1);
}

if (builder.cookieHandler != null) {
b.cookieHandler(builder.cookieHandler);
httpClientBuilder.cookieHandler(builder.cookieHandler);
}

if (builder.getProxyHost() != null) {
if (builder.isProxyIsSocks4()) {
log.warn(
"Socks4 is likely not supported by this client. See https://bugs.openjdk.org/browse/JDK-8214516");
}
b.proxy(
httpClientBuilder.proxy(
ProxySelector.of(new InetSocketAddress(builder.getProxyHost(), builder.getProxyPort())));
}
this.httpClient = b.build();
this.httpClient = httpClientBuilder.build();

assert ObjectReleaseTracker.track(this);
}
Expand All @@ -147,7 +150,7 @@ protected CompletableFuture<HttpResponse<InputStream>> requestInputStreamAsync(
PreparedRequest pReq = prepareRequest(baseUrl, solrRequest, collection);
return httpClient
.sendAsync(pReq.reqb.build(), HttpResponse.BodyHandlers.ofInputStream())
.whenComplete((httpResponse, throwable) -> releaseContentWriting(pReq));
.whenComplete((httpResponse, throwable) -> pReq.releaseContentWriting());
} catch (Exception e) {
CompletableFuture<HttpResponse<InputStream>> cf = new CompletableFuture<>();
cf.completeExceptionally(e);
Expand All @@ -162,7 +165,7 @@ public CompletableFuture<NamedList<Object>> requestAsync(
PreparedRequest pReq = prepareRequest(null, solrRequest, collection);
return httpClient
.sendAsync(pReq.reqb.build(), HttpResponse.BodyHandlers.ofInputStream())
.whenComplete((httpResponse, throwable) -> releaseContentWriting(pReq))
.whenComplete((httpResponse, throwable) -> pReq.releaseContentWriting())
.thenApply(
httpResponse -> {
try {
Expand All @@ -179,21 +182,6 @@ public CompletableFuture<NamedList<Object>> requestAsync(
}
}

private void releaseContentWriting(PreparedRequest pReq) {
if (pReq.contentWritingFuture != null) {
pReq.contentWritingFuture.cancel(true);
}
// Closing the sink is what unblocks a writer already stuck in the pipe; cancel() alone does
// not.
if (pReq.contentWritingSink != null) {
try {
pReq.contentWritingSink.close();
} catch (IOException e) {
log.warn("Could not close content-writing pipe", e);
}
}
}

@Override
public NamedList<Object> requestWithBaseUrl(
String baseUrl, SolrRequest<?> solrRequest, String collection)
Expand All @@ -214,9 +202,7 @@ public NamedList<Object> requestWithBaseUrl(
} catch (RuntimeException e) {
throw new SolrServerException(e);
} finally {
if (pReq.contentWritingFuture != null) {
pReq.contentWritingFuture.cancel(true);
}
pReq.releaseContentWriting();

// See
// https://docs.oracle.com/en/java/javase/17/docs/api/java.net.http/java/net/http/HttpResponse.BodySubscribers.html#ofInputStream()
Expand Down Expand Up @@ -252,7 +238,7 @@ protected PreparedRequest prepareRequest(
ResponseParser parserToUse = responseParser(solrRequest);
ModifiableSolrParams queryParams = initializeSolrParams(solrRequest, parserToUse);
var reqb = HttpRequest.newBuilder();
PreparedRequest pReq = null;
PreparedRequest pReq;
try {
switch (solrRequest.getMethod()) {
case GET:
Expand Down Expand Up @@ -289,7 +275,7 @@ private PreparedRequest prepareGet(
reqb.GET();
decorateRequest(reqb, solrRequest);
reqb.uri(new URI(url + queryParams.toQueryString()));
return new PreparedRequest(reqb, null, null);
return new PreparedRequest(reqb);
}

private PreparedRequest preparePutOrPost(
Expand Down Expand Up @@ -320,28 +306,16 @@ private PreparedRequest preparePutOrPost(
}

HttpRequest.BodyPublisher bodyPublisher;
Future<?> contentWritingFuture = null;
PipedInputStream contentWritingSink = null;
PreparedRequest pReq = new PreparedRequest(reqb);
if (contentWriter != null) {
boolean success = maybeTryHeadRequest(url);
if (!success) {
reqb.version(HttpClient.Version.HTTP_1_1);
}

final PipedOutputStream source = new PipedOutputStream();
contentWritingSink = new PipedInputStream(source);
final PipedInputStream sink = contentWritingSink;
bodyPublisher = HttpRequest.BodyPublishers.ofInputStream(() -> sink);

contentWritingFuture =
executor.submit(
() -> {
try (source) {
contentWriter.write(source);
} catch (Exception e) {
log.error("Cannot write Content Stream", e);
}
});
bodyPublisher =
HttpRequest.BodyPublishers.ofInputStream(
() -> pReq.beginContentWriting(contentWriter, this.executor));
} else if (streams != null && streams.size() == 1) {
boolean success = maybeTryHeadRequest(url);
if (!success) {
Expand Down Expand Up @@ -372,25 +346,62 @@ private PreparedRequest preparePutOrPost(
URI uriWithQueryParams = new URI(url + queryParams.toQueryString());
reqb.uri(uriWithQueryParams);

return new PreparedRequest(reqb, contentWritingFuture, contentWritingSink);
return pReq;
}

protected static class PreparedRequest {
Future<?> contentWritingFuture;
PipedInputStream contentWritingSink;
HttpRequest.Builder reqb;
final HttpRequest.Builder reqb;

ResponseParser parserToUse;

String url;

PreparedRequest(
HttpRequest.Builder reqb,
Future<?> contentWritingFuture,
PipedInputStream contentWritingSink) {
// Both remain null if the request has no streamed content, or if the body is never requested
// (e.g. the connection failed before sending it). Filled in lazily by
// beginContentWriting once the JDK HttpClient actually requests the body.
private PipedInputStream contentWritingSink;
private Future<?> contentWritingFuture;

PreparedRequest(HttpRequest.Builder reqb) {
this.reqb = reqb;
this.contentWritingFuture = contentWritingFuture;
this.contentWritingSink = contentWritingSink;
}

synchronized PipedInputStream beginContentWriting(
RequestWriter.ContentWriter contentWriter, ExecutorService bodyExecutor) {
final PipedOutputStream source = new PipedOutputStream();
try {
contentWritingSink = new PipedInputStream(source);
} catch (IOException e) {
throw new UncheckedIOException(e);
}

contentWritingFuture =
bodyExecutor.submit(
() -> {
// note: doesn't need to synchronize with PreparedRequest.this
try (source) {
contentWriter.write(source);
} catch (Exception e) {
log.error("Cannot write Content Stream", e);
}
});
return contentWritingSink;
}

synchronized void releaseContentWriting() {
if (contentWritingFuture != null) {
contentWritingFuture.cancel(true);
}

// Closing the sink is what unblocks a writer already stuck in the pipe; cancel() alone does
// not.
if (contentWritingSink != null) {
try {
contentWritingSink.close();
} catch (IOException e) {
log.warn("Could not close content-writing pipe", e);
}
}
}
}

Expand Down Expand Up @@ -545,6 +556,12 @@ public void close() throws IOException {
}
executor = null;

// The http client executor is always created and owned by this instance.
if (httpClientExecutor != null) {
ExecutorUtil.shutdownAndAwaitTermination(httpClientExecutor);
httpClientExecutor = null;
}

assert ObjectReleaseTracker.release(this);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@
import org.apache.solr.client.solrj.request.SolrQuery;
import org.apache.solr.client.solrj.request.UpdateRequest;
import org.apache.solr.client.solrj.request.XMLRequestWriter;
import org.apache.solr.client.solrj.request.json.JsonQueryRequest;
import org.apache.solr.client.solrj.response.JavaBinResponseParser;
import org.apache.solr.client.solrj.response.ResponseParser;
import org.apache.solr.client.solrj.response.SolrPingResponse;
Expand Down Expand Up @@ -637,6 +638,54 @@ public void testMaybeTryHeadRequestHasContentType() throws Exception {
}
}

@Test(timeout = 30000)
public void testConcurrentStreamedBodiesDoNotDeadlockWithHttp1() throws Exception {
DebugServlet.clear();
DebugServlet.addResponseHeader("Content-Type", "application/octet-stream");
DebugServlet.responseBodyByQueryFragment.put("", javabinResponse());
String url = solrTestRule.getBaseUrl() + DEBUG_SERVLET_PATH;

int concurrency = 8;
ExecutorService callers =
ExecutorUtil.newMDCAwareFixedThreadPool(concurrency, new NamedThreadFactory("test-caller"));

try (HttpJdkSolrClient client = builder(url).useHttp1_1(true).build()) {
List<CompletableFuture<Void>> futures = new ArrayList<>(concurrency);
for (int i = 0; i < concurrency; i++) {
futures.add(
CompletableFuture.runAsync(
() -> {
JsonQueryRequest q = buildLargeBodyQuery();
try {
q.process(client);
} catch (SolrServerException | IOException e) {
throw new RuntimeException(e);
}
},
callers));
}
CompletableFuture.allOf(futures.toArray(new CompletableFuture<?>[0]))
.get(30, TimeUnit.SECONDS);
} finally {
ExecutorUtil.shutdownAndAwaitTermination(callers);
}
}

private static JsonQueryRequest buildLargeBodyQuery() {
StringBuilder filter = new StringBuilder("id:(");
for (int i = 0; i < 400; i++) {
if (i > 0) {
filter.append(" OR ");
}
filter.append("value_").append(i);
}
filter.append(')');
JsonQueryRequest q = new JsonQueryRequest();
q.setQuery("*:*");
q.withFilter(filter.toString());
return q;
}

/**
* This is not required for any test, but there appears to be a bug in the JDK client where it
* does not release all threads if the client has not performed any queries, even after a forced
Expand Down
Loading