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");
+        }
+    }
+}

Reply via email to