diff --git a/changelog/unreleased/SOLR-18472-tikaserver-httpclient-refcount.yml b/changelog/unreleased/SOLR-18472-tikaserver-httpclient-refcount.yml new file mode 100644 index 000000000000..5caf01ba1fa0 --- /dev/null +++ b/changelog/unreleased/SOLR-18472-tikaserver-httpclient-refcount.yml @@ -0,0 +1,10 @@ +title: > + Fixed a race in the Tika Server extraction backend where closing or reloading one core could leave + another core using a stopped HTTP client. +type: fixed +authors: + - name: Jan Høydahl + nick: janhoy +links: + - name: SOLR-18472 + url: https://issues.apache.org/jira/browse/SOLR-18472 diff --git a/solr/modules/extraction/src/java/org/apache/solr/handler/extraction/TikaServerExtractionBackend.java b/solr/modules/extraction/src/java/org/apache/solr/handler/extraction/TikaServerExtractionBackend.java index a6d08bb11f77..2c422f7ebfe5 100644 --- a/solr/modules/extraction/src/java/org/apache/solr/handler/extraction/TikaServerExtractionBackend.java +++ b/solr/modules/extraction/src/java/org/apache/solr/handler/extraction/TikaServerExtractionBackend.java @@ -33,6 +33,7 @@ import java.util.concurrent.ThreadFactory; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicBoolean; import org.apache.solr.common.SolrException; import org.apache.solr.common.util.ExecutorUtil; import org.apache.solr.common.util.NamedList; @@ -67,11 +68,13 @@ public class TikaServerExtractionBackend implements ExtractionBackend { private HashMap initArgsMap = new HashMap<>(); private final long maxCharsLimit; - // Singleton holder for the shared HttpClient/Executor resources (one per JVM) - private static volatile RefCounted SHARED_RESOURCES; + // Singleton holder for the shared HttpClient/Executor resources (one per JVM), guarded by + // INIT_LOCK so that acquiring a reference cannot race with the last decref() closing it + private static RefCounted SHARED_RESOURCES; // Per-backend handle (same RefCounted instance as SHARED_RESOURCES) that this instance will // decref() on close - private RefCounted acquiredResourcesRef; + private final RefCounted acquiredResourcesRef; + private final AtomicBoolean closed = new AtomicBoolean(); public TikaServerExtractionBackend(String baseUrl) { this(baseUrl, DEFAULT_TIMEOUT_SECONDS, null, DEFAULT_MAXCHARS_LIMIT); @@ -117,7 +120,7 @@ public TikaServerExtractionBackend( Duration.ofSeconds(timeoutSeconds > 0 ? timeoutSeconds : DEFAULT_TIMEOUT_SECONDS); // Acquire a reference to the shared resources; keep a handle so we can decref() on close - acquiredResourcesRef = initializeHttpClient().incref(); + acquiredResourcesRef = acquireSharedResources(); } public static final String NAME = "tikaserver"; @@ -367,6 +370,7 @@ private static final class ResourcesRef extends RefCounted @Override protected void close() { + assert Thread.holdsLock(INIT_LOCK); // stop client and shutdown executor try { if (resource.client != null) resource.client.stop(); @@ -376,20 +380,16 @@ protected void close() { if (resource.executor != null) resource.executor.shutdownNow(); } catch (Throwable ignore) { } - synchronized (INIT_LOCK) { - // clear the shared reference when closed - if (SHARED_RESOURCES == this) { - SHARED_RESOURCES = null; - } + // clear the shared reference when closed + if (SHARED_RESOURCES == this) { + SHARED_RESOURCES = null; } } } - private static RefCounted initializeHttpClient() { - RefCounted ref = SHARED_RESOURCES; - if (ref != null) return ref; + private static RefCounted acquireSharedResources() { synchronized (INIT_LOCK) { - if (SHARED_RESOURCES != null) return SHARED_RESOURCES; + if (SHARED_RESOURCES != null) return SHARED_RESOURCES.incref(); ThreadFactory tf = new SolrNamedThreadFactory("TikaServerHttpClient"); ExecutorService exec = ExecutorUtil.newMDCAwareCachedThreadPool(tf); HttpClient client = new HttpClient(); @@ -406,7 +406,7 @@ private static RefCounted initializeHttpClient() { SolrException.ErrorCode.SERVER_ERROR, "Failed to start shared Jetty HttpClient", e); } SHARED_RESOURCES = new ResourcesRef(new HttpClientResources(client, exec)); - return SHARED_RESOURCES; + return SHARED_RESOURCES.incref(); } } @@ -447,13 +447,10 @@ private void appendBackCompatTikaMetadata(ExtractionMetadata md) { @Override public void close() { - RefCounted ref; - synchronized (INIT_LOCK) { - ref = acquiredResourcesRef; - acquiredResourcesRef = null; - } - if (ref != null) { - ref.decref(); + if (closed.compareAndSet(false, true)) { + synchronized (INIT_LOCK) { + acquiredResourcesRef.decref(); + } } } } diff --git a/solr/modules/extraction/src/test/org/apache/solr/handler/extraction/TikaServerExtractionBackendTest.java b/solr/modules/extraction/src/test/org/apache/solr/handler/extraction/TikaServerExtractionBackendTest.java index 326ab818596c..4505a346fe4d 100644 --- a/solr/modules/extraction/src/test/org/apache/solr/handler/extraction/TikaServerExtractionBackendTest.java +++ b/solr/modules/extraction/src/test/org/apache/solr/handler/extraction/TikaServerExtractionBackendTest.java @@ -24,10 +24,15 @@ import java.util.List; import java.util.Locale; import java.util.Map; +import java.util.concurrent.CyclicBarrier; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Future; import org.apache.lucene.tests.util.QuickPatchThreadsFilter; import org.apache.solr.SolrIgnoredThreadsFilter; import org.apache.solr.SolrTestCaseJ4; import org.apache.solr.common.SolrException; +import org.apache.solr.common.util.ExecutorUtil; +import org.apache.solr.common.util.SolrNamedThreadFactory; import org.apache.solr.handler.extraction.fromtika.ToXMLContentHandler; import org.junit.ClassRule; import org.junit.Test; @@ -195,4 +200,64 @@ public void testMaxCharsLimitEnforcedWithSaxHandler() throws Exception { } } } + + /** + * Races the last backend's close() (which stops the shared HttpClient) against construction of a + * new backend, which must end up with a working client rather than the one being stopped. + */ + @Test + public void testConcurrentCloseAndAcquireSharedResources() throws Exception { + String baseUrl = tikaContainer.getBaseUrl(); + int iterations = atLeast(25); + ExecutorService exec = + ExecutorUtil.newMDCAwareFixedThreadPool(2, new SolrNamedThreadFactory("tikaRace")); + try { + for (int i = 0; i < iterations; i++) { + TikaServerExtractionBackend last = new TikaServerExtractionBackend(baseUrl); + CyclicBarrier barrier = new CyclicBarrier(2); + Future closer = + exec.submit( + () -> { + barrier.await(); + last.close(); + return null; + }); + Future opener = + exec.submit( + () -> { + barrier.await(); + return new TikaServerExtractionBackend(baseUrl); + }); + closer.get(); + try (TikaServerExtractionBackend fresh = opener.get()) { + assertExtracts(fresh, "iteration " + i); + } + } + } finally { + ExecutorUtil.shutdownAndAwaitTermination(exec); + } + } + + @Test + public void testCloseIsIdempotent() throws Exception { + String baseUrl = tikaContainer.getBaseUrl(); + TikaServerExtractionBackend first = new TikaServerExtractionBackend(baseUrl); + try (TikaServerExtractionBackend second = new TikaServerExtractionBackend(baseUrl)) { + first.close(); + first.close(); + assertExtracts(second, "after double close"); + } + try (TikaServerExtractionBackend third = new TikaServerExtractionBackend(baseUrl)) { + assertExtracts(third, "after all closed"); + } + } + + private void assertExtracts(TikaServerExtractionBackend backend, String marker) throws Exception { + byte[] data = ("Hello " + marker).getBytes(StandardCharsets.UTF_8); + try (ByteArrayInputStream in = new ByteArrayInputStream(data)) { + ExtractionResult res = backend.extract(in, newRequest("test.txt", "text/plain", "text")); + assertNotNull(res); + assertTrue(res.getContent().contains("Hello " + marker)); + } + } }