Skip to content
Merged
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
10 changes: 10 additions & 0 deletions changelog/unreleased/SOLR-18472-tikaserver-httpclient-refcount.yml
Original file line number Diff line number Diff line change
@@ -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
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -67,11 +68,13 @@ public class TikaServerExtractionBackend implements ExtractionBackend {
private HashMap<String, Object> initArgsMap = new HashMap<>();
private final long maxCharsLimit;

// Singleton holder for the shared HttpClient/Executor resources (one per JVM)
private static volatile RefCounted<HttpClientResources> 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<HttpClientResources> SHARED_RESOURCES;
// Per-backend handle (same RefCounted instance as SHARED_RESOURCES) that this instance will
// decref() on close
private RefCounted<HttpClientResources> acquiredResourcesRef;
private final RefCounted<HttpClientResources> acquiredResourcesRef;
private final AtomicBoolean closed = new AtomicBoolean();

public TikaServerExtractionBackend(String baseUrl) {
this(baseUrl, DEFAULT_TIMEOUT_SECONDS, null, DEFAULT_MAXCHARS_LIMIT);
Expand Down Expand Up @@ -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";
Expand Down Expand Up @@ -367,6 +370,7 @@ private static final class ResourcesRef extends RefCounted<HttpClientResources>

@Override
protected void close() {
assert Thread.holdsLock(INIT_LOCK);
// stop client and shutdown executor
try {
if (resource.client != null) resource.client.stop();
Expand All @@ -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<HttpClientResources> initializeHttpClient() {
RefCounted<HttpClientResources> ref = SHARED_RESOURCES;
if (ref != null) return ref;
private static RefCounted<HttpClientResources> 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();
Expand All @@ -406,7 +406,7 @@ private static RefCounted<HttpClientResources> 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();
}
}

Expand Down Expand Up @@ -447,13 +447,10 @@ private void appendBackCompatTikaMetadata(ExtractionMetadata md) {

@Override
public void close() {
RefCounted<HttpClientResources> ref;
synchronized (INIT_LOCK) {
ref = acquiredResourcesRef;
acquiredResourcesRef = null;
}
if (ref != null) {
ref.decref();
if (closed.compareAndSet(false, true)) {
synchronized (INIT_LOCK) {
acquiredResourcesRef.decref();
Comment thread
janhoy marked this conversation as resolved.
}
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<TikaServerExtractionBackend> 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));
}
}
}
Loading