This is an automated email from the ASF dual-hosted git repository.
tballison pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/tika.git
The following commit(s) were added to refs/heads/main by this push:
new 4abbb31ee1 TIKA-4953: add PipesParser.start() and
PipesForkParser.start() (#3299)
4abbb31ee1 is described below
commit 4abbb31ee11b59de4dc8f4c2052a3eabaada25b9
Author: Tim Allison <[email protected]>
AuthorDate: Mon Oct 5 19:55:10 2026 -0400
TIKA-4953: add PipesParser.start() and PipesForkParser.start() (#3299)
---
CHANGES.txt | 3 +
.../ROOT/pages/using-tika/java-api/pipes.adoc | 4 +
.../org/apache/tika/pipes/core/PipesClient.java | 20 +++++
.../org/apache/tika/pipes/core/PipesParser.java | 65 +++++++++++++++
.../apache/tika/pipes/fork/PipesForkParser.java | 18 +++++
.../tika/pipes/fork/PipesForkParserTest.java | 27 +++++++
.../apache/tika/pipes/core/PipesStartupTest.java | 94 ++++++++++++++++++++++
7 files changed, 231 insertions(+)
diff --git a/CHANGES.txt b/CHANGES.txt
index d00798a75c..4d4009589e 100644
--- a/CHANGES.txt
+++ b/CHANGES.txt
@@ -1,5 +1,8 @@
Release 4.2.0 - unreleased
+ * PipesParser.start() and PipesForkParser.start() bring the forks up before
the
+ first parse and fail fast if one can't start (TIKA-4953).
+
* Refactor date parsing for more consistent, accurate dates. File-system
times from archive entries and mail attachments move to fs:created/
fs:modified. MailDateParser is superseded by TikaDates and MailUtil
diff --git a/docs/modules/ROOT/pages/using-tika/java-api/pipes.adoc
b/docs/modules/ROOT/pages/using-tika/java-api/pipes.adoc
index c332561dca..b7516d78ce 100644
--- a/docs/modules/ROOT/pages/using-tika/java-api/pipes.adoc
+++ b/docs/modules/ROOT/pages/using-tika/java-api/pipes.adoc
@@ -331,6 +331,10 @@ problems instead of returning them.
== Lifecycle
+Forks start on the first parse. To start them up front, call `start()` on
`PipesParser` or
+`PipesForkParser`: it waits until every fork is ready and throws if one can't
start (bad config,
+a missing plugin, bad `javaPath`), so a broken setup shows up before any work
is queued.
+
All three classes are `Closeable`. Closing kills the forked JVMs and removes
the temporary files
the parent created for them. Each fork also restarts itself after
`maxFilesProcessedPerProcess` documents to
bound slow leaks in the parsing libraries. Sizing the pool — `numClients`,
heap per fork, CPU
diff --git
a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/PipesClient.java
b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/PipesClient.java
index da8722ff8b..aeff9c8d64 100644
---
a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/PipesClient.java
+++
b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/PipesClient.java
@@ -142,6 +142,26 @@ public class PipesClient implements Closeable {
this.ownsServerManager = true;
}
+ /**
+ * Brings the server up and waits until it is ready, rather than on the
first
+ * {@link #process}. Does nothing if it is already up.
+ *
+ * @throws ServerInitializationException if the server can't start
+ */
+ public void start() throws InterruptedException,
ServerInitializationException {
+ try {
+ maybeInit();
+ } catch (InterruptedException e) {
+ serverManager.connectionAbandoned();
+ closeConnection();
+ throw e;
+ } catch (ServerInitializationException e) {
+ serverManager.markServerForRestart(RestartReason.CRASH,
connectionGeneration);
+ closeConnection();
+ throw e;
+ }
+ }
+
public int getFilesProcessed() {
return filesProcessed;
}
diff --git
a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/PipesParser.java
b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/PipesParser.java
index d154a5962e..b123517ac5 100644
---
a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/PipesParser.java
+++
b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/PipesParser.java
@@ -21,6 +21,10 @@ import java.io.IOException;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.List;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
import java.util.concurrent.LinkedBlockingDeque;
import java.util.concurrent.TimeUnit;
@@ -136,6 +140,67 @@ public class PipesParser implements Closeable {
}
}
+ /**
+ * Starts every idle client's server now and waits until each is ready,
instead of on the
+ * first {@link #parse}. Optional: it surfaces a server that can't start
(bad config, a
+ * missing plugin) before work is queued. Clients busy in a {@link #parse}
are skipped;
+ * they are already up.
+ *
+ * @throws ServerInitializationException if a server fails to start
+ */
+ public void start() throws InterruptedException,
ServerInitializationException {
+ List<PipesClient> idle = new ArrayList<>();
+ PipesClient client;
+ while ((client = clientQueue.pollFirst()) != null) {
+ idle.add(client);
+ }
+ if (idle.isEmpty()) {
+ return;
+ }
+ ExecutorService executor = Executors.newFixedThreadPool(idle.size());
+ try {
+ List<Future<Void>> futures = new ArrayList<>();
+ for (PipesClient c : idle) {
+ futures.add(executor.submit(() -> {
+ c.start();
+ return null;
+ }));
+ }
+ ServerInitializationException failure = null;
+ for (Future<Void> future : futures) {
+ try {
+ future.get();
+ } catch (ExecutionException e) {
+ if (failure == null) {
+ failure = e.getCause() instanceof
ServerInitializationException sie ? sie
+ : new ServerInitializationException("server
failed to start", e.getCause());
+ }
+ }
+ }
+ if (failure != null) {
+ throw failure;
+ }
+ } finally {
+ // A client goes back in the queue only once nothing is using it;
each start is
+ // bounded by the startup timeouts.
+ executor.shutdownNow();
+ boolean interrupted = false;
+ while (true) {
+ try {
+ if (executor.awaitTermination(1, TimeUnit.SECONDS)) {
+ break;
+ }
+ } catch (InterruptedException e) {
+ interrupted = true;
+ }
+ }
+ idle.forEach(clientQueue::offerFirst);
+ if (interrupted) {
+ Thread.currentThread().interrupt();
+ }
+ }
+ }
+
@Override
public void close() throws IOException {
List<IOException> exceptions = new ArrayList<>();
diff --git
a/tika-pipes/tika-pipes-fork-parser/src/main/java/org/apache/tika/pipes/fork/PipesForkParser.java
b/tika-pipes/tika-pipes-fork-parser/src/main/java/org/apache/tika/pipes/fork/PipesForkParser.java
index f9a6c35da4..1719601228 100644
---
a/tika-pipes/tika-pipes-fork-parser/src/main/java/org/apache/tika/pipes/fork/PipesForkParser.java
+++
b/tika-pipes/tika-pipes-fork-parser/src/main/java/org/apache/tika/pipes/fork/PipesForkParser.java
@@ -39,6 +39,7 @@ import org.apache.tika.pipes.core.EmitStrategy;
import org.apache.tika.pipes.core.PipesConfig;
import org.apache.tika.pipes.core.PipesException;
import org.apache.tika.pipes.core.PipesParser;
+import org.apache.tika.pipes.core.ServerInitializationException;
import org.apache.tika.pipes.core.config.ConfigMerger;
import org.apache.tika.pipes.core.config.ConfigOverrides;
import org.apache.tika.pipes.core.config.DefaultPluginsDir;
@@ -150,6 +151,23 @@ public class PipesForkParser implements Closeable {
this.pipesParser = PipesParser.load(tikaConfigPath);
}
+ /**
+ * Starts the forked processes now and waits until each is ready, instead
of on the first
+ * {@link #parse}. Optional: call it to find out that the forks can't
start (bad config,
+ * missing plugins, bad JVM args) before any work is queued.
+ *
+ * @throws PipesForkParserException with status {@code
FAILED_TO_INITIALIZE} if a forked
+ * process fails to start
+ */
+ public void start() throws InterruptedException, PipesForkParserException {
+ try {
+ pipesParser.start();
+ } catch (ServerInitializationException e) {
+ throw new
PipesForkParserException(PipesResult.RESULT_STATUS.FAILED_TO_INITIALIZE,
+ "Failed to start forked process: " + e.getMessage(), e);
+ }
+ }
+
/**
* Parse a file in a forked JVM process.
*
diff --git
a/tika-pipes/tika-pipes-fork-parser/src/test/java/org/apache/tika/pipes/fork/PipesForkParserTest.java
b/tika-pipes/tika-pipes-fork-parser/src/test/java/org/apache/tika/pipes/fork/PipesForkParserTest.java
index b0a1c60075..be5b24e640 100644
---
a/tika-pipes/tika-pipes-fork-parser/src/test/java/org/apache/tika/pipes/fork/PipesForkParserTest.java
+++
b/tika-pipes/tika-pipes-fork-parser/src/test/java/org/apache/tika/pipes/fork/PipesForkParserTest.java
@@ -143,6 +143,33 @@ public class PipesForkParserTest {
}
/** TIKA-4931: a code setting equal to the default still overrides the
user config file. */
+ @Test
+ public void testStartBeforeFirstParse() throws Exception {
+ Path testFile = tempDir.resolve("test.txt");
+ Files.writeString(testFile, "hello");
+ PipesForkParserConfig config = new
PipesForkParserConfig().setPluginsDir(PLUGINS_DIR);
+
+ try (PipesForkParser parser = new PipesForkParser(config);
+ TikaInputStream tis = TikaInputStream.get(testFile)) {
+ parser.start();
+ assertTrue(parser.parse(tis).isSuccess());
+ }
+ }
+
+ /** start() reports a fork that can't start, without a parse to provoke
it. */
+ @Test
+ public void testStartFailsFast() throws Exception {
+ PipesForkParserConfig config = new PipesForkParserConfig()
+ .setPluginsDir(PLUGINS_DIR)
+ .setJavaPath(tempDir.resolve("no-such-java").toString());
+
+ try (PipesForkParser parser = new PipesForkParser(config)) {
+ PipesForkParserException e =
assertThrows(PipesForkParserException.class, parser::start);
+ assertEquals(PipesResult.RESULT_STATUS.FAILED_TO_INITIALIZE,
e.getStatus());
+ assertTrue(e.getMessage().contains("no-such-java"),
e.getMessage());
+ }
+ }
+
@Test
public void testExplicitDefaultJavaPathBeatsUserConfig() throws Exception {
Path userConfig = tempDir.resolve("user-config.json");
diff --git
a/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/PipesStartupTest.java
b/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/PipesStartupTest.java
new file mode 100644
index 0000000000..7ef55526f8
--- /dev/null
+++
b/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/PipesStartupTest.java
@@ -0,0 +1,94 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You 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
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * 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 org.apache.tika.pipes.core;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+
+import java.nio.file.Path;
+
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+import org.apache.tika.config.loader.TikaJsonConfig;
+import org.apache.tika.metadata.Metadata;
+import org.apache.tika.parser.ParseContext;
+import org.apache.tika.pipes.api.FetchEmitTuple;
+import org.apache.tika.pipes.api.PipesResult;
+import org.apache.tika.pipes.api.emitter.EmitKey;
+import org.apache.tika.pipes.api.fetcher.FetchKey;
+
+/** {@link PipesParser#start()} brings the forks up before the first parse. */
+public class PipesStartupTest {
+
+ private static final String TEST_DOC = "testOverlappingText.pdf";
+
+ private static Path config(Path tmp) throws Exception {
+ Path tikaConfigPath = PluginsTestHelper.getFileSystemFetcherConfig(
+ tmp, tmp.resolve("input"), tmp.resolve("output"));
+ PluginsTestHelper.copyTestFilesToTmpInput(tmp, TEST_DOC);
+ return tikaConfigPath;
+ }
+
+ private static PipesResult parse(PipesParser parser) throws Exception {
+ return parser.parse(new FetchEmitTuple(TEST_DOC, new FetchKey("fsf",
TEST_DOC),
+ new EmitKey(), new Metadata(), new ParseContext(),
+ FetchEmitTuple.ON_PARSE_EXCEPTION.SKIP));
+ }
+
+ @Test
+ public void startBringsUpEveryFork(@TempDir Path tmp) throws Exception {
+ Path tikaConfigPath = config(tmp);
+ TikaJsonConfig tikaJsonConfig = TikaJsonConfig.load(tikaConfigPath);
+ PipesConfig pipesConfig = PipesConfig.load(tikaJsonConfig);
+ pipesConfig.setNumClients(2);
+ try (PipesParser parser = PipesParser.load(tikaJsonConfig,
pipesConfig, tikaConfigPath)) {
+ assertEquals(0, parser.startedServerCount());
+ parser.start();
+ assertEquals(2, parser.startedServerCount());
+ assertEquals(2, parser.getIdleClientCount(), "start() must return
every client");
+ assertEquals(PipesResult.RESULT_STATUS.PARSE_SUCCESS,
parse(parser).status());
+ }
+ }
+
+ @Test
+ public void startInSharedMode(@TempDir Path tmp) throws Exception {
+ Path tikaConfigPath = config(tmp);
+ TikaJsonConfig tikaJsonConfig = TikaJsonConfig.load(tikaConfigPath);
+ PipesConfig pipesConfig = PipesConfig.load(tikaJsonConfig);
+ pipesConfig.setNumClients(2);
+ pipesConfig.setUseSharedServer(true);
+ try (PipesParser parser = PipesParser.load(tikaJsonConfig,
pipesConfig, tikaConfigPath)) {
+ parser.start();
+ assertEquals(1, parser.startedServerCount());
+ assertEquals(PipesResult.RESULT_STATUS.PARSE_SUCCESS,
parse(parser).status());
+ }
+ }
+
+ @Test
+ public void startFailsOnBadConfig(@TempDir Path tmp) throws Exception {
+ Path tikaConfigPath = PluginsTestHelper.getFileSystemFetcherConfig(
+ "tika-config-bad-class.json", tmp);
+ TikaJsonConfig tikaJsonConfig = TikaJsonConfig.load(tikaConfigPath);
+ PipesConfig pipesConfig = PipesConfig.load(tikaJsonConfig);
+ try (PipesParser parser = PipesParser.load(tikaJsonConfig,
pipesConfig, tikaConfigPath)) {
+ assertThrows(ServerInitializationException.class, parser::start);
+ assertEquals(pipesConfig.getNumClients(),
parser.getIdleClientCount(),
+ "a failed start() must still return every client");
+ }
+ }
+}