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());
