This is an automated email from the ASF dual-hosted git repository. tballison pushed a commit to branch TIKA-4848-pipes in repository https://gitbox.apache.org/repos/asf/tika.git
commit cd685c80742a66128e4aa1512919db992c611fd9 Author: tallison <[email protected]> AuthorDate: Thu Aug 27 20:33:23 2026 -0400 TIKA-4848: apply exception-reporting policy in pipes worker --- .../tika/pipes/core/server/ConnectionHandler.java | 3 +- .../apache/tika/pipes/core/server/EmitHandler.java | 4 +- .../tika/pipes/core/server/FetchHandler.java | 9 +- .../tika/pipes/core/server/ParseHandler.java | 6 +- .../apache/tika/pipes/core/server/PipesServer.java | 7 +- .../apache/tika/pipes/core/server/PipesWorker.java | 14 +-- .../tika/pipes/core/server/ServerProtocolIO.java | 10 ++- .../pipes/core/server/SharedServerResources.java | 16 +++- .../pipes/core/server/ServerProtocolIOTest.java | 18 +++- .../pipes/core/ExceptionReportingPipesTest.java | 99 ++++++++++++++++++++++ .../configs/tika-config-exception-reporting.json | 59 +++++++++++++ .../src/test/resources/test-documents/mock-npe.xml | 23 +++++ 12 files changed, 244 insertions(+), 24 deletions(-) diff --git a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/ConnectionHandler.java b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/ConnectionHandler.java index b1b7c1861e..a1755be2ab 100644 --- a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/ConnectionHandler.java +++ b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/ConnectionHandler.java @@ -105,7 +105,8 @@ public class ConnectionHandler implements Runnable, Closeable { this.resources = resources; this.pipesConfig = pipesConfig; this.heartbeatIntervalMillis = pipesConfig.getHeartbeatIntervalMillis(); - this.protocolIO = new ServerProtocolIO(input, output, pipesConfig.getMaxIpcPayloadBytes()); + this.protocolIO = new ServerProtocolIO(input, output, pipesConfig.getMaxIpcPayloadBytes(), + resources.getExceptionReporting()); } @Override diff --git a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/EmitHandler.java b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/EmitHandler.java index 8458e7ceec..f7d956882a 100644 --- a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/EmitHandler.java +++ b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/EmitHandler.java @@ -161,7 +161,7 @@ class EmitHandler { return new PipesResult(PipesResult.RESULT_STATUS.EMITTER_NOT_FOUND, noEmitterMsg); } catch (IOException | TikaException e) { LOG.warn("Couldn't initialize emitter for task id '" + taskId + "'", e); - return new PipesResult(PipesResult.RESULT_STATUS.EMITTER_INITIALIZATION_EXCEPTION, ExceptionUtils.getStackTrace(e)); + return new PipesResult(PipesResult.RESULT_STATUS.EMITTER_INITIALIZATION_EXCEPTION, ExceptionUtils.format(e, parseContext)); } try { ParseMode parseMode = parseContext.get(ParseMode.class); @@ -175,7 +175,7 @@ class EmitHandler { } } catch (IOException e) { LOG.warn("emit exception", e); - String msg = ExceptionUtils.getStackTrace(e); + String msg = ExceptionUtils.format(e, parseContext); //for now, we're hiding the parse exception if there was also an emit exception return new PipesResult(PipesResult.RESULT_STATUS.EMIT_EXCEPTION, msg); } diff --git a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/FetchHandler.java b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/FetchHandler.java index e5116b63bb..5a0820bc00 100644 --- a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/FetchHandler.java +++ b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/FetchHandler.java @@ -46,7 +46,7 @@ class FetchHandler { } public TisOrResult fetch(FetchEmitTuple fetchEmitTuple, Metadata metadata, ParseContext parseContext) { - FetcherOrResult fetcherResult = getFetcher(fetchEmitTuple); + FetcherOrResult fetcherResult = getFetcher(fetchEmitTuple, parseContext); if (fetcherResult.pipesResult != null) { return new TisOrResult(null, fetcherResult.pipesResult); } @@ -55,11 +55,12 @@ class FetchHandler { fetchEmitTuple.getFetchKey().getFetchKey(), metadata, parseContext); return new TisOrResult(tis, null); } catch (IOException | TikaException e) { - return new TisOrResult(null, new PipesResult(PipesResult.RESULT_STATUS.FETCH_EXCEPTION, ExceptionUtils.getStackTrace(e))); + return new TisOrResult(null, new PipesResult(PipesResult.RESULT_STATUS.FETCH_EXCEPTION, + ExceptionUtils.format(e, parseContext))); } } - private FetcherOrResult getFetcher(FetchEmitTuple t) { + private FetcherOrResult getFetcher(FetchEmitTuple t, ParseContext parseContext) { String fetcherId = t.getFetchKey().getFetcherId(); // Built in, not configured: the bytes come with the request, so there is nothing for an // operator to point at and no reason for it to occupy an id in the ConfigStore. @@ -75,7 +76,7 @@ class FetchHandler { } catch (IOException | TikaException e) { LOG.warn("Couldn't initialize fetcher for fetch id={}", t.getId(), e); return new FetcherOrResult(null, new PipesResult(PipesResult.RESULT_STATUS.FETCHER_INITIALIZATION_EXCEPTION, - ExceptionUtils.getStackTrace(e))); + ExceptionUtils.format(e, parseContext))); } } diff --git a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/ParseHandler.java b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/ParseHandler.java index e8aa9375fb..2c055ef1d0 100644 --- a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/ParseHandler.java +++ b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/ParseHandler.java @@ -244,19 +244,19 @@ class ParseHandler { try { autoDetectParser.parse(stream, handler, metadata, parseContext); } catch (SAXException e) { - containerException = ExceptionUtils.getStackTrace(e); + containerException = ExceptionUtils.format(ExceptionUtils.unwrapTikaException(e), parseContext); LOG.warn("sax problem:" + fetchEmitTuple.getId(), e); if (WriteLimitReachedException.isWriteLimitReached(e)) { writeLimitReached = true; } } catch (EncryptedDocumentException e) { - containerException = ExceptionUtils.getStackTrace(e); + containerException = ExceptionUtils.format(ExceptionUtils.unwrapTikaException(e), parseContext); LOG.warn("encrypted document:" + fetchEmitTuple.getId(), e); } catch (SecurityException e) { LOG.warn("security exception:" + fetchEmitTuple.getId(), e); throw e; } catch (Exception e) { - containerException = ExceptionUtils.getStackTrace(e); + containerException = ExceptionUtils.format(ExceptionUtils.unwrapTikaException(e), parseContext); LOG.warn("parse exception: " + fetchEmitTuple.getId(), e); } finally { // BasicContentHandlerFactory's "ignore" handler's toString() returns "" (not diff --git a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/PipesServer.java b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/PipesServer.java index a92602242d..17b16eb7ca 100644 --- a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/PipesServer.java +++ b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/PipesServer.java @@ -45,6 +45,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.xml.sax.SAXException; +import org.apache.tika.config.ExceptionReporting; import org.apache.tika.config.ParseTimeout; import org.apache.tika.config.TimeoutLimits; import org.apache.tika.config.loader.TikaJsonConfig; @@ -218,7 +219,8 @@ public class PipesServer implements AutoCloseable { } catch (Exception e) { LOG.error("Failed to start up", e); try { - String msg = ExceptionUtils.getStackTrace(e); + // Config may be what failed to load, so no ExceptionReporting policy is available. + String msg = ExceptionUtils.format(e, ExceptionReporting.DEFAULT); byte[] bytes = msg.getBytes(StandardCharsets.UTF_8); PipesMessage.startupFailed(bytes).write(dos); // pipesConfig may not have loaded successfully (that may be why we're @@ -252,7 +254,8 @@ public class PipesServer implements AutoCloseable { validateHeartbeatInterval(pipesConfig); emitStrategy = pipesConfig.getEmitStrategy().getType(); - this.protocolIO = new ServerProtocolIO(input, output, pipesConfig.getMaxIpcPayloadBytes()); + this.protocolIO = new ServerProtocolIO(input, output, pipesConfig.getMaxIpcPayloadBytes(), + ExceptionReporting.get(tikaLoader.loadParseContext())); } diff --git a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/PipesWorker.java b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/PipesWorker.java index f08f3fed3c..2f0b49eeec 100644 --- a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/PipesWorker.java +++ b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/PipesWorker.java @@ -198,7 +198,7 @@ class PipesWorker implements Callable<PipesResult> { } catch (IOException e) { LOG.warn("Failed to create zip file", e); return new PipesResult(PipesResult.RESULT_STATUS.EMIT_EXCEPTION, - "Failed to create zip file: " + ExceptionUtils.getStackTrace(e)); + "Failed to create zip file: " + ExceptionUtils.format(e, parseContext)); } // Emit the zip file @@ -208,7 +208,7 @@ class PipesWorker implements Callable<PipesResult> { } catch (IOException e) { LOG.warn("Failed to emit zip file", e); return new PipesResult(PipesResult.RESULT_STATUS.EMIT_EXCEPTION, - "Failed to emit zip file: " + ExceptionUtils.getStackTrace(e)); + "Failed to emit zip file: " + ExceptionUtils.format(e, parseContext)); } LOG.debug("Successfully zipped and emitted {} embedded files to {}", @@ -289,7 +289,7 @@ class PipesWorker implements Callable<PipesResult> { } catch (IOException e) { LOG.warn("Failed to create Frictionless zip file", e); return new PipesResult(PipesResult.RESULT_STATUS.EMIT_EXCEPTION, - "Failed to create Frictionless zip file: " + ExceptionUtils.getStackTrace(e)); + "Failed to create Frictionless zip file: " + ExceptionUtils.format(e, parseContext)); } // Emit the zip file @@ -299,7 +299,7 @@ class PipesWorker implements Callable<PipesResult> { } catch (IOException e) { LOG.warn("Failed to emit Frictionless zip file", e); return new PipesResult(PipesResult.RESULT_STATUS.EMIT_EXCEPTION, - "Failed to emit Frictionless zip file: " + ExceptionUtils.getStackTrace(e)); + "Failed to emit Frictionless zip file: " + ExceptionUtils.format(e, parseContext)); } LOG.debug("Successfully emitted Frictionless package with {} resources to {}", @@ -351,7 +351,7 @@ class PipesWorker implements Callable<PipesResult> { } catch (IOException e) { LOG.warn("Failed to emit Frictionless directory output", e); return new PipesResult(PipesResult.RESULT_STATUS.EMIT_EXCEPTION, - "Failed to emit Frictionless directory output: " + ExceptionUtils.getStackTrace(e)); + "Failed to emit Frictionless directory output: " + ExceptionUtils.format(e, parseContext)); } LOG.debug("Successfully emitted Frictionless package with {} resources (directory mode) to {}", @@ -484,7 +484,7 @@ class PipesWorker implements Callable<PipesResult> { } catch (IOException e) { LOG.warn("fetcher initialization exception id={}", fetchEmitTuple.getId(), e); return new ParseDataOrPipesResult(null, - new PipesResult(PipesResult.RESULT_STATUS.FETCHER_INITIALIZATION_EXCEPTION, ExceptionUtils.getStackTrace(e))); + new PipesResult(PipesResult.RESULT_STATUS.FETCHER_INITIALIZATION_EXCEPTION, ExceptionUtils.format(e, parseContext))); } // Use newMetadata() to apply any configured write limits Metadata metadata = localContext.newMetadata(); @@ -503,7 +503,7 @@ class PipesWorker implements Callable<PipesResult> { } catch (TikaException | IOException e) { LOG.warn("fetch exception id={}", fetchEmitTuple.getId(), e); return new ParseDataOrPipesResult(null, - new PipesResult(PipesResult.RESULT_STATUS.UNSPECIFIED_CRASH, ExceptionUtils.getStackTrace(e))); + new PipesResult(PipesResult.RESULT_STATUS.UNSPECIFIED_CRASH, ExceptionUtils.format(e, parseContext))); } } diff --git a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/ServerProtocolIO.java b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/ServerProtocolIO.java index e94407d1c5..09ba6ed797 100644 --- a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/ServerProtocolIO.java +++ b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/ServerProtocolIO.java @@ -26,6 +26,7 @@ import java.util.Locale; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.apache.tika.config.ExceptionReporting; import org.apache.tika.config.TimeoutLimits; import org.apache.tika.exception.TikaConfigException; import org.apache.tika.metadata.Metadata; @@ -82,8 +83,15 @@ public class ServerProtocolIO { private final DataInputStream input; private final DataOutputStream output; private final int maxIpcPayloadBytes; + private final ExceptionReporting exceptionReporting; public ServerProtocolIO(DataInputStream input, DataOutputStream output, int maxIpcPayloadBytes) { + this(input, output, maxIpcPayloadBytes, ExceptionReporting.DEFAULT); + } + + public ServerProtocolIO(DataInputStream input, DataOutputStream output, int maxIpcPayloadBytes, + ExceptionReporting exceptionReporting) { + this.exceptionReporting = exceptionReporting; if (maxIpcPayloadBytes < MIN_FALLBACK_PAYLOAD_BYTES) { throw new IllegalArgumentException(String.format(Locale.ROOT, "maxIpcPayloadBytes %d is below the minimum %d required to carry a PAYLOAD_LIMIT_EXCEEDED response", @@ -196,7 +204,7 @@ public class ServerProtocolIO { * @throws IOException on serialization, I/O, or unexpected ACK response */ public void writeCrash(PipesMessageType crashType, Throwable t) throws IOException { - String msg = (t != null) ? ExceptionUtils.getStackTrace(t) : ""; + String msg = (t != null) ? ExceptionUtils.format(t, exceptionReporting) : ""; BoundedOutputStream bos = new BoundedOutputStream(maxIpcPayloadBytes); try { JsonPipesIpc.toStream(msg, bos); diff --git a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/SharedServerResources.java b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/SharedServerResources.java index a488a4922e..29338b2823 100644 --- a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/SharedServerResources.java +++ b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/server/SharedServerResources.java @@ -20,6 +20,7 @@ import java.io.IOException; import org.xml.sax.SAXException; +import org.apache.tika.config.ExceptionReporting; import org.apache.tika.config.loader.TikaJsonConfig; import org.apache.tika.config.loader.TikaLoader; import org.apache.tika.detect.Detector; @@ -64,6 +65,7 @@ public class SharedServerResources { private final MetadataWriteLimiterFactory defaultMetadataWriteLimiterFactory; private final EmitStrategy emitStrategy; private final ConfigStore configStore; + private final ExceptionReporting exceptionReporting; private SharedServerResources(TikaLoader tikaLoader, PipesConfig pipesConfig, AutoDetectParser autoDetectParser, Detector detector, @@ -71,7 +73,9 @@ public class SharedServerResources { EmitterManager emitterManager, MetadataFilter defaultMetadataFilter, ContentHandlerFactory defaultContentHandlerFactory, MetadataWriteLimiterFactory defaultMetadataWriteLimiterFactory, - EmitStrategy emitStrategy, ConfigStore configStore) { + EmitStrategy emitStrategy, ConfigStore configStore, + ExceptionReporting exceptionReporting) { + this.exceptionReporting = exceptionReporting; this.tikaLoader = tikaLoader; this.pipesConfig = pipesConfig; this.autoDetectParser = autoDetectParser; @@ -117,14 +121,16 @@ public class SharedServerResources { // Load filters and factories MetadataFilter metadataFilter = tikaLoader.loadMetadataFilters(); ContentHandlerFactory contentHandlerFactory = tikaLoader.loadContentHandlerFactory(); + ParseContext configContext = tikaLoader.loadParseContext(); MetadataWriteLimiterFactory metadataWriteLimiterFactory = - tikaLoader.loadParseContext().get(MetadataWriteLimiterFactory.class); + configContext.get(MetadataWriteLimiterFactory.class); EmitStrategy emitStrategy = pipesConfig.getEmitStrategy().getType(); return new SharedServerResources(tikaLoader, pipesConfig, autoDetectParser, detector, rMetaParser, fetcherManager, emitterManager, metadataFilter, contentHandlerFactory, - metadataWriteLimiterFactory, emitStrategy, configStore); + metadataWriteLimiterFactory, emitStrategy, configStore, + ExceptionReporting.get(configContext)); } private static ConfigStore createConfigStore(PipesConfig pipesConfig, TikaPluginManager tikaPluginManager) @@ -206,6 +212,10 @@ public class SharedServerResources { return emitStrategy; } + public ExceptionReporting getExceptionReporting() { + return exceptionReporting; + } + public ConfigStore getConfigStore() { return configStore; } diff --git a/tika-pipes/tika-pipes-core/src/test/java/org/apache/tika/pipes/core/server/ServerProtocolIOTest.java b/tika-pipes/tika-pipes-core/src/test/java/org/apache/tika/pipes/core/server/ServerProtocolIOTest.java index 016f630ee4..bbb456bd4d 100644 --- a/tika-pipes/tika-pipes-core/src/test/java/org/apache/tika/pipes/core/server/ServerProtocolIOTest.java +++ b/tika-pipes/tika-pipes-core/src/test/java/org/apache/tika/pipes/core/server/ServerProtocolIOTest.java @@ -17,6 +17,7 @@ package org.apache.tika.pipes.core.server; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -29,6 +30,7 @@ import java.util.concurrent.atomic.AtomicReference; import org.junit.jupiter.api.Test; +import org.apache.tika.config.ExceptionReporting; import org.apache.tika.config.TimeoutLimits; import org.apache.tika.metadata.Metadata; import org.apache.tika.parser.ParseContext; @@ -293,6 +295,11 @@ class ServerProtocolIOTest { * sends ACK, and returns the deserialized stack trace string. */ private String exchangeCrash(Throwable t, int maxPayloadBytes) throws Exception { + return exchangeCrash(t, maxPayloadBytes, ExceptionReporting.DEFAULT); + } + + private String exchangeCrash(Throwable t, int maxPayloadBytes, ExceptionReporting reporting) + throws Exception { PipedOutputStream serverOutPipe = new PipedOutputStream(); PipedInputStream clientInPipe = new PipedInputStream(serverOutPipe, 1024 * 1024); PipedOutputStream clientOutPipe = new PipedOutputStream(); @@ -320,7 +327,7 @@ class ServerProtocolIOTest { ServerProtocolIO io = new ServerProtocolIO( new DataInputStream(serverInPipe), new DataOutputStream(serverOutPipe), - maxPayloadBytes); + maxPayloadBytes, reporting); io.writeCrash(PipesMessageType.UNSPECIFIED_CRASH, t); clientThread.join(5000); @@ -358,4 +365,13 @@ class ServerProtocolIOTest { // Empty string fallback: the trace was too large, we get an empty payload. assertEquals("", returned); } + + @Test + void testCrashHonorsExceptionReporting() throws Exception { + RuntimeException ex = new RuntimeException("something failed"); + String returned = exchangeCrash(ex, PipesMessage.MAX_PAYLOAD_BYTES, + new ExceptionReporting(ExceptionReporting.Level.MESSAGE_REDACTED, -1)); + assertTrue(returned.contains("java.lang.RuntimeException")); + assertFalse(returned.contains("something failed")); + } } diff --git a/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/ExceptionReportingPipesTest.java b/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/ExceptionReportingPipesTest.java new file mode 100644 index 0000000000..89094f94f0 --- /dev/null +++ b/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/ExceptionReportingPipesTest.java @@ -0,0 +1,99 @@ +/* + * 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.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +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.metadata.TikaCoreProperties; +import org.apache.tika.parser.ParseContext; +import org.apache.tika.pipes.api.FetchEmitTuple; +import org.apache.tika.pipes.api.ParseMode; +import org.apache.tika.pipes.api.PipesResult; +import org.apache.tika.pipes.api.emitter.EmitKey; +import org.apache.tika.pipes.api.fetcher.FetchKey; + +/** + * The forked worker must apply the config's exception-reporting policy to the container + * exception it records and to the PipesResult message it returns on fetch failure. + */ +public class ExceptionReportingPipesTest { + + private static final String FETCHER_NAME = "fsf"; + private static final String NPE_DOC = "mock-npe.xml"; + private static final String MESSAGE = "secret null pointer message"; + + private PipesClient init(Path tmp, String template) throws Exception { + Path tikaConfigPath = PluginsTestHelper.getFileSystemFetcherConfig( + template, tmp, tmp.resolve("input"), tmp.resolve("output"), false); + PluginsTestHelper.copyTestFilesToTmpInput(tmp, NPE_DOC); + PipesConfig pipesConfig = PipesConfig.load(TikaJsonConfig.load(tikaConfigPath)); + return new PipesClient(pipesConfig, tikaConfigPath); + } + + private static FetchEmitTuple tuple(String fetchKey) { + ParseContext parseContext = new ParseContext(); + parseContext.set(ParseMode.class, ParseMode.RMETA); + return new FetchEmitTuple(fetchKey, new FetchKey(FETCHER_NAME, fetchKey), new EmitKey(), + new Metadata(), parseContext, FetchEmitTuple.ON_PARSE_EXCEPTION.EMIT); + } + + @Test + public void containerExceptionRedacted(@TempDir Path tmp) throws Exception { + try (PipesClient client = init(tmp, "tika-config-exception-reporting.json")) { + PipesResult result = client.process(tuple(NPE_DOC)); + assertEquals(PipesResult.RESULT_STATUS.PARSE_SUCCESS_WITH_EXCEPTION, result.status()); + String trace = result.emitData().getMetadataList().get(0) + .get(TikaCoreProperties.CONTAINER_EXCEPTION); + assertTrue(trace.contains("java.lang.NullPointerException"), trace); + assertTrue(trace.contains("\tat "), trace); + assertFalse(trace.contains(MESSAGE), trace); + assertEquals(trace, result.emitData().getContainerStackTrace()); + } + } + + @Test + public void containerExceptionFullByDefault(@TempDir Path tmp) throws Exception { + try (PipesClient client = init(tmp, "tika-config-basic.json")) { + PipesResult result = client.process(tuple(NPE_DOC)); + assertEquals(PipesResult.RESULT_STATUS.PARSE_SUCCESS_WITH_EXCEPTION, result.status()); + String trace = result.emitData().getMetadataList().get(0) + .get(TikaCoreProperties.CONTAINER_EXCEPTION); + assertTrue(trace.contains(MESSAGE), trace); + } + } + + @Test + public void fetchExceptionRedacted(@TempDir Path tmp) throws Exception { + try (PipesClient client = init(tmp, "tika-config-exception-reporting.json")) { + PipesResult result = client.process(tuple("does-not-exist.xml")); + assertEquals(PipesResult.RESULT_STATUS.FETCH_EXCEPTION, result.status()); + String msg = result.message(); + assertTrue(msg.contains("Exception"), msg); + // the fetcher's message names the missing path; the policy must strip it + assertFalse(msg.contains("does-not-exist"), msg); + } + } +} diff --git a/tika-pipes/tika-pipes-integration-tests/src/test/resources/configs/tika-config-exception-reporting.json b/tika-pipes/tika-pipes-integration-tests/src/test/resources/configs/tika-config-exception-reporting.json new file mode 100644 index 0000000000..fbbaa4c906 --- /dev/null +++ b/tika-pipes/tika-pipes-integration-tests/src/test/resources/configs/tika-config-exception-reporting.json @@ -0,0 +1,59 @@ +{ + "content-handler-factory": { + "basic-content-handler-factory": { + "type": "TEXT", + "writeLimit": -1, + "throwOnWriteLimitReached": true + } + }, + "fetchers": { + "fsf": { + "file-system-fetcher": { + "basePath": "FETCHER_BASE_PATH", + "extractFileSystemMetadata": false + } + } + }, + "emitters": { + "fse": { + "file-system-emitter": { + "basePath": "EMITTER_BASE_PATH", + "fileExtension": "json", + "onExists": "EXCEPTION" + } + } + }, + "pipes-iterator": { + "file-system-pipes-iterator": { + "basePath": "FETCHER_BASE_PATH", + "countTotal": true, + "fetcherId": "fsf", + "emitterId": "fse" + } + }, + "pipes": { + "parseMode": "RMETA", + "onParseException": "EMIT", + "numClients": 4, + "emitIntermediateResults": "EMIT_INTERMEDIATE_RESULTS", + "forkedJvmArgs": ["-Xmx512m"], + "emitStrategy": { + "type": "DYNAMIC", + "thresholdBytes": 1000000 + } + }, + "auto-detect-parser": { + "throwOnZeroBytes": false + }, + "parse-context": { + "mock-digester-factory": {}, + "exception-reporting": { + "level": "MESSAGE_REDACTED", + "maxLength": 10000 + }, + "timeout-limits": { + "progressTimeoutMillis": 5000 + } + }, + "plugin-roots": "PLUGINS_PATHS" +} diff --git a/tika-pipes/tika-pipes-integration-tests/src/test/resources/test-documents/mock-npe.xml b/tika-pipes/tika-pipes-integration-tests/src/test/resources/test-documents/mock-npe.xml new file mode 100644 index 0000000000..f8f99a3490 --- /dev/null +++ b/tika-pipes/tika-pipes-integration-tests/src/test/resources/test-documents/mock-npe.xml @@ -0,0 +1,23 @@ +<?xml version="1.0" encoding="UTF-8" ?> +<!-- + 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. +--> +<mock> + <write element="p">main_content</write> + <throw class="java.lang.NullPointerException">secret null pointer message</throw> +</mock>
