This is an automated email from the ASF dual-hosted git repository.

tballison pushed a commit to branch TIKA-4809-stage-2
in repository https://gitbox.apache.org/repos/asf/tika.git

commit 6910b7ba912ba2660978c646f885f5995171c657
Author: tallison <[email protected]>
AuthorDate: Fri Aug 7 11:38:40 2026 -0400

    TIKA-4809: Merge /tika+/rmeta+/unpack and /pipes onto one shared PipesParser
---
 .../apache/tika/server/core/TikaServerProcess.java | 37 +++++++++++-----------
 .../server/core/resource/PipesParsingHelper.java   | 14 ++++++++
 .../tika/server/core/resource/PipesResource.java   | 36 ++++++++++-----------
 .../org/apache/tika/server/core/TikaPipesTest.java | 21 +++++++++---
 .../apache/tika/server/standard/TikaPipesTest.java | 20 +++++++++---
 5 files changed, 83 insertions(+), 45 deletions(-)

diff --git 
a/tika-server/tika-server-core/src/main/java/org/apache/tika/server/core/TikaServerProcess.java
 
b/tika-server/tika-server-core/src/main/java/org/apache/tika/server/core/TikaServerProcess.java
index ea569bceb1..351ae9a0bc 100644
--- 
a/tika-server/tika-server-core/src/main/java/org/apache/tika/server/core/TikaServerProcess.java
+++ 
b/tika-server/tika-server-core/src/main/java/org/apache/tika/server/core/TikaServerProcess.java
@@ -180,11 +180,12 @@ public class TikaServerProcess {
 
         ServerStatus serverStatus = new ServerStatus();
 
-        // Initialize pipes-based parsing only if /tika or /rmeta endpoints 
are enabled
+        // Initialize pipes-based parsing (and its shared PipesParser) only if 
any
+        // pipes-backed endpoint is enabled.
         PipesParsingHelper pipesParsingHelper = null;
         if (needsPipesParsingHelper(tikaServerConfig)) {
             pipesParsingHelper = initPipesParsingHelper(tikaServerConfig);
-            LOG.info("Pipes-based parsing enabled for /tika and /rmeta 
endpoints");
+            LOG.info("Pipes-based parsing enabled for /tika, /rmeta, /unpack, 
and /pipes endpoints");
         }
 
         TikaResource tikaResource = new TikaResource(tikaLoader, serverStatus, 
pipesParsingHelper,
@@ -424,17 +425,12 @@ public class TikaServerProcess {
             resourceProviders.add(new 
SingletonResourceProvider(localAsyncResource));
         }
         if (addPipesResource) {
-            final PipesResource localPipesResource = new 
PipesResource(tikaServerConfig.getConfigPath());
-            Runtime
-                    .getRuntime()
-                    .addShutdownHook(new Thread(() -> {
-                        try {
-                            localPipesResource.close();
-                        } catch (Exception e) {
-                            LOG.warn("exception closing local pipes resource", 
e);
-                        }
-                    }));
-            resourceProviders.add(new 
SingletonResourceProvider(localPipesResource));
+            // /pipes shares its PipesParser with /tika+/rmeta+/unpack (see
+            // needsPipesParsingHelper) -- non-null here is guaranteed by that 
check.
+            // Lifecycle (shutdown/close) is owned by whoever built the shared 
parser,
+            // not by PipesResource.
+            PipesParsingHelper helper = tikaResource.getPipesParsingHelper();
+            resourceProviders.add(new SingletonResourceProvider(new 
PipesResource(helper.getPipesParser())));
         }
         resourceProviders.addAll(loadResourceServices(serverStatus));
         return resourceProviders;
@@ -458,17 +454,22 @@ public class TikaServerProcess {
     }
 
     /**
-     * Determines if PipesParsingHelper is needed based on configured 
endpoints.
-     * It's needed when /tika or /rmeta endpoints are enabled (either 
explicitly or by default).
+     * Determines if the shared PipesParser (wrapped in PipesParsingHelper) is 
needed
+     * based on configured endpoints. It's needed when /tika, /rmeta, /unpack, 
or /pipes
+     * are enabled (either explicitly or by default) -- all four now share one 
parser.
+     * (Note: unlike the others, /pipes also requires allowPipes to actually 
start; if
+     * it's listed without allowPipes, loadCoreProviders will refuse to start 
regardless
+     * of whether this method already triggered building the shared parser.)
      */
     private static boolean needsPipesParsingHelper(TikaServerConfig 
tikaServerConfig) {
         List<String> endpoints = tikaServerConfig.getEndpoints();
-        // If no endpoints specified, all default endpoints are loaded 
(including tika and rmeta)
+        // If no endpoints specified, all default endpoints are loaded 
(including
+        // tika, rmeta, and unpack; pipes too when allowPipes is set)
         if (endpoints == null || endpoints.isEmpty()) {
             return true;
         }
-        // Check if tika or rmeta are in the configured endpoints
-        return endpoints.contains("tika") || endpoints.contains("rmeta");
+        return endpoints.contains("tika") || endpoints.contains("rmeta")
+                || endpoints.contains("unpack") || endpoints.contains("pipes");
     }
 
     /**
diff --git 
a/tika-server/tika-server-core/src/main/java/org/apache/tika/server/core/resource/PipesParsingHelper.java
 
b/tika-server/tika-server-core/src/main/java/org/apache/tika/server/core/resource/PipesParsingHelper.java
index 7a0660d639..78913c6875 100644
--- 
a/tika-server/tika-server-core/src/main/java/org/apache/tika/server/core/resource/PipesParsingHelper.java
+++ 
b/tika-server/tika-server-core/src/main/java/org/apache/tika/server/core/resource/PipesParsingHelper.java
@@ -42,6 +42,8 @@ import org.apache.tika.pipes.api.PipesResult;
 import org.apache.tika.pipes.api.emitter.EmitData;
 import org.apache.tika.pipes.api.emitter.EmitKey;
 import org.apache.tika.pipes.api.fetcher.FetchKey;
+import org.apache.tika.pipes.core.EmitStrategy;
+import org.apache.tika.pipes.core.EmitStrategyConfig;
 import org.apache.tika.pipes.core.PipesConfig;
 import org.apache.tika.pipes.core.PipesException;
 import org.apache.tika.pipes.core.PipesParser;
@@ -145,6 +147,12 @@ public class PipesParsingHelper {
             // Set parse mode in context
             parseContext.set(ParseMode.class, parseMode);
 
+            // This parser is shared with /pipes, whose own default is 
EMIT_ALL. No
+            // emitter is configured for /tika/rmeta/unpack requests 
(EmitKey.NO_EMIT
+            // below) -- results must come back over the socket, so set 
PASSBACK_ALL
+            // explicitly per-request rather than relying on the parser-level 
default.
+            parseContext.set(EmitStrategyConfig.class, new 
EmitStrategyConfig(EmitStrategy.PASSBACK_ALL));
+
             // Create FetchEmitTuple with relative filename (basePath is 
configured in fetcher)
             FetchKey fetchKey = new FetchKey(DEFAULT_FETCHER_ID, relativeName);
 
@@ -358,6 +366,12 @@ public class PipesParsingHelper {
             // Set parse mode to UNPACK
             parseContext.set(ParseMode.class, ParseMode.UNPACK);
 
+            // Shared parser (see parse() above) -- PASSBACK_ALL is also 
required here
+            // for correctness: with UNPACK mode, EmitHandler.shouldEmit() 
only skips
+            // re-emitting metadata (already emitted as part of the zip) when 
the
+            // effective strategy is PASSBACK_ALL.
+            parseContext.set(EmitStrategyConfig.class, new 
EmitStrategyConfig(EmitStrategy.PASSBACK_ALL));
+
             // Configure UnpackConfig - use existing or create new
             UnpackConfig unpackConfig = parseContext.get(UnpackConfig.class);
             if (unpackConfig == null) {
diff --git 
a/tika-server/tika-server-core/src/main/java/org/apache/tika/server/core/resource/PipesResource.java
 
b/tika-server/tika-server-core/src/main/java/org/apache/tika/server/core/resource/PipesResource.java
index 2c6316e081..960935e910 100644
--- 
a/tika-server/tika-server-core/src/main/java/org/apache/tika/server/core/resource/PipesResource.java
+++ 
b/tika-server/tika-server-core/src/main/java/org/apache/tika/server/core/resource/PipesResource.java
@@ -33,15 +33,13 @@ import jakarta.ws.rs.core.UriInfo;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
-import org.apache.tika.config.loader.TikaJsonConfig;
-import org.apache.tika.exception.TikaConfigException;
 import org.apache.tika.metadata.Metadata;
 import org.apache.tika.metadata.TikaCoreProperties;
+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.core.EmitStrategy;
 import org.apache.tika.pipes.core.EmitStrategyConfig;
-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.serialization.JsonFetchEmitTuple;
@@ -55,17 +53,13 @@ public class PipesResource {
 
     private final PipesParser pipesParser;
 
-    public PipesResource(java.nio.file.Path tikaConfig) throws 
TikaConfigException, IOException {
-        TikaJsonConfig tikaJsonConfig = TikaJsonConfig.load(tikaConfig);
-        PipesConfig pipesConfig = PipesConfig.load(tikaJsonConfig);
-        // The /pipes endpoint always emits from the child process; force 
EMIT_ALL.
-        if (pipesConfig.getEmitStrategy().getType() != EmitStrategy.EMIT_ALL) {
-            if (pipesConfig.getEmitStrategy().getType() != 
EmitStrategyConfig.DEFAULT_EMIT_STRATEGY) {
-                LOG.warn("resetting emit strategy to EMIT_ALL for pipes 
endpoint");
-            }
-            pipesConfig.setEmitStrategy(new 
EmitStrategyConfig(EmitStrategy.EMIT_ALL));
-        }
-        this.pipesParser = PipesParser.load(tikaJsonConfig, pipesConfig, 
tikaConfig);
+    /**
+     * @param pipesParser shared parser, also used by /tika, /rmeta, and 
/unpack.
+     *                     Lifecycle (construction, shutdown) is owned by 
whoever
+     *                     built it, not by this class.
+     */
+    public PipesResource(PipesParser pipesParser) {
+        this.pipesParser = pipesParser;
     }
 
 
@@ -98,7 +92,15 @@ public class PipesResource {
     }
 
     private Map<String, String> processTuple(FetchEmitTuple fetchEmitTuple) 
throws InterruptedException, PipesException, IOException {
-
+        // This parser is shared with /tika+/rmeta+/unpack, whose own default 
is
+        // PASSBACK_ALL. /pipes needs the child to emit via the client's 
configured
+        // emitter by default -- set EMIT_ALL explicitly per-request rather 
than
+        // relying on the parser-level default, but don't clobber a caller's 
own
+        // explicit override if they set one.
+        ParseContext parseContext = fetchEmitTuple.getParseContext();
+        if (parseContext.get(EmitStrategyConfig.class) == null) {
+            parseContext.set(EmitStrategyConfig.class, new 
EmitStrategyConfig(EmitStrategy.EMIT_ALL));
+        }
         PipesResult pipesResult = pipesParser.parse(fetchEmitTuple);
         if (pipesResult.isProcessCrash()) {
             return returnProcessCrash(pipesResult.status().toString());
@@ -144,8 +146,4 @@ public class PipesResource {
         statusMap.put("type", type);
         return statusMap;
     }
-
-    public void close() throws IOException {
-        pipesParser.close();
-    }
 }
diff --git 
a/tika-server/tika-server-core/src/test/java/org/apache/tika/server/core/TikaPipesTest.java
 
b/tika-server/tika-server-core/src/test/java/org/apache/tika/server/core/TikaPipesTest.java
index 38b1ec47a7..1f0ef15a71 100644
--- 
a/tika-server/tika-server-core/src/test/java/org/apache/tika/server/core/TikaPipesTest.java
+++ 
b/tika-server/tika-server-core/src/test/java/org/apache/tika/server/core/TikaPipesTest.java
@@ -49,6 +49,7 @@ import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.TestInstance;
 
+import org.apache.tika.config.loader.TikaJsonConfig;
 import org.apache.tika.exception.TikaConfigException;
 import org.apache.tika.metadata.Metadata;
 import org.apache.tika.metadata.TikaCoreProperties;
@@ -57,6 +58,10 @@ import org.apache.tika.pipes.api.FetchEmitTuple;
 import org.apache.tika.pipes.api.ParseMode;
 import org.apache.tika.pipes.api.emitter.EmitKey;
 import org.apache.tika.pipes.api.fetcher.FetchKey;
+import org.apache.tika.pipes.core.EmitStrategy;
+import org.apache.tika.pipes.core.EmitStrategyConfig;
+import org.apache.tika.pipes.core.PipesConfig;
+import org.apache.tika.pipes.core.PipesParser;
 import org.apache.tika.pipes.core.serialization.JsonFetchEmitTuple;
 import org.apache.tika.sax.BasicContentHandlerFactory;
 import org.apache.tika.sax.ContentHandlerFactory;
@@ -85,6 +90,7 @@ public class TikaPipesTest extends CXFTestBase {
     private static final String[] VALUE_ARRAY = new String[]{"my-value-1", 
"my-value-2", "my-value-3"};
 
     private PipesResource pipesResource;
+    private PipesParser pipesParser;
 
     @Override
     @BeforeAll
@@ -115,10 +121,11 @@ public class TikaPipesTest extends CXFTestBase {
     @Override
     @AfterAll
     public void tearDown() throws Exception {
-        if (pipesResource != null) {
-            pipesResource.close();
-            pipesResource = null;
+        if (pipesParser != null) {
+            pipesParser.close();
+            pipesParser = null;
         }
+        pipesResource = null;
         super.tearDown();
         if (tmpDir != null) {
             FileUtils.deleteDirectory(tmpDir.toFile());
@@ -140,7 +147,13 @@ public class TikaPipesTest extends CXFTestBase {
     protected void setUpResources(JAXRSServerFactoryBean sf) {
         List<ResourceProvider> rCoreProviders = new ArrayList<>();
         try {
-            pipesResource = new PipesResource(tikaConfigPath);
+            // Mirrors what PipesResource used to build internally, back when 
it
+            // constructed its own parser instead of sharing one.
+            TikaJsonConfig tikaJsonConfig = 
TikaJsonConfig.load(tikaConfigPath);
+            PipesConfig pipesConfig = PipesConfig.load(tikaJsonConfig);
+            pipesConfig.setEmitStrategy(new 
EmitStrategyConfig(EmitStrategy.EMIT_ALL));
+            pipesParser = PipesParser.load(tikaJsonConfig, pipesConfig, 
tikaConfigPath);
+            pipesResource = new PipesResource(pipesParser);
             rCoreProviders.add(new SingletonResourceProvider(pipesResource));
         } catch (IOException | TikaConfigException e) {
             throw new RuntimeException(e);
diff --git 
a/tika-server/tika-server-standard/src/test/java/org/apache/tika/server/standard/TikaPipesTest.java
 
b/tika-server/tika-server-standard/src/test/java/org/apache/tika/server/standard/TikaPipesTest.java
index 544382a2bd..d27b2e1715 100644
--- 
a/tika-server/tika-server-standard/src/test/java/org/apache/tika/server/standard/TikaPipesTest.java
+++ 
b/tika-server/tika-server-standard/src/test/java/org/apache/tika/server/standard/TikaPipesTest.java
@@ -59,6 +59,10 @@ import org.apache.tika.pipes.api.FetchEmitTuple;
 import org.apache.tika.pipes.api.ParseMode;
 import org.apache.tika.pipes.api.emitter.EmitKey;
 import org.apache.tika.pipes.api.fetcher.FetchKey;
+import org.apache.tika.pipes.core.EmitStrategy;
+import org.apache.tika.pipes.core.EmitStrategyConfig;
+import org.apache.tika.pipes.core.PipesConfig;
+import org.apache.tika.pipes.core.PipesParser;
 import org.apache.tika.pipes.core.extractor.UnpackConfig;
 import org.apache.tika.pipes.core.fetcher.FetcherManager;
 import org.apache.tika.pipes.core.serialization.JsonFetchEmitTuple;
@@ -93,6 +97,7 @@ public class TikaPipesTest extends CXFTestBase {
     private FetcherManager fetcherManager;
 
     private PipesResource pipesResource;
+    private PipesParser pipesParser;
 
     @Override
     @BeforeAll
@@ -134,7 +139,13 @@ public class TikaPipesTest extends CXFTestBase {
     protected void setUpResources(JAXRSServerFactoryBean sf) {
         List<ResourceProvider> rCoreProviders = new ArrayList<>();
         try {
-            pipesResource = new PipesResource(tikaConfigPath);
+            // Mirrors what PipesResource used to build internally, back when 
it
+            // constructed its own parser instead of sharing one.
+            TikaJsonConfig tikaJsonConfig = 
TikaJsonConfig.load(tikaConfigPath);
+            PipesConfig pipesConfig = PipesConfig.load(tikaJsonConfig);
+            pipesConfig.setEmitStrategy(new 
EmitStrategyConfig(EmitStrategy.EMIT_ALL));
+            pipesParser = PipesParser.load(tikaJsonConfig, pipesConfig, 
tikaConfigPath);
+            pipesResource = new PipesResource(pipesParser);
             rCoreProviders.add(new SingletonResourceProvider(pipesResource));
         } catch (IOException | TikaConfigException e) {
             throw new RuntimeException(e);
@@ -145,10 +156,11 @@ public class TikaPipesTest extends CXFTestBase {
     @Override
     @AfterAll
     public void tearDown() throws Exception {
-        if (pipesResource != null) {
-            pipesResource.close();
-            pipesResource = null;
+        if (pipesParser != null) {
+            pipesParser.close();
+            pipesParser = null;
         }
+        pipesResource = null;
         super.tearDown();
         if (tmpWorkingDir != null) {
             FileUtils.deleteDirectory(tmpWorkingDir.toFile());

Reply via email to