diff --git a/solr/solrj/src/java/org/apache/solr/client/solrj/impl/HttpJdkSolrClient.java b/solr/solrj/src/java/org/apache/solr/client/solrj/impl/HttpJdkSolrClient.java index 111bfd1bc92..64ae67c3182 100644 --- a/solr/solrj/src/java/org/apache/solr/client/solrj/impl/HttpJdkSolrClient.java +++ b/solr/solrj/src/java/org/apache/solr/client/solrj/impl/HttpJdkSolrClient.java @@ -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; @@ -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; @@ -78,54 +76,59 @@ 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 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) { @@ -133,10 +136,10 @@ protected HttpJdkSolrClient(String serverBaseUrl, HttpJdkSolrClient.Builder buil 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); } @@ -147,7 +150,7 @@ protected CompletableFuture> 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> cf = new CompletableFuture<>(); cf.completeExceptionally(e); @@ -162,7 +165,7 @@ public CompletableFuture> 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 { @@ -179,21 +182,6 @@ public CompletableFuture> 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 requestWithBaseUrl( String baseUrl, SolrRequest solrRequest, String collection) @@ -214,9 +202,7 @@ public NamedList 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() @@ -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: @@ -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( @@ -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) { @@ -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); + } + } } } @@ -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); } diff --git a/solr/solrj/src/test/org/apache/solr/client/solrj/impl/HttpJdkSolrClientTest.java b/solr/solrj/src/test/org/apache/solr/client/solrj/impl/HttpJdkSolrClientTest.java index 79bfdb63e9c..9f9233f375e 100644 --- a/solr/solrj/src/test/org/apache/solr/client/solrj/impl/HttpJdkSolrClientTest.java +++ b/solr/solrj/src/test/org/apache/solr/client/solrj/impl/HttpJdkSolrClientTest.java @@ -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; @@ -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> 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