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 1468136e9e TIKA-4835 -- tika-core primitives for spill-less parsing
(#3073)
1468136e9e is described below
commit 1468136e9e9f89bc1d98d648dcb26709ae2b9037
Author: Tim Allison <[email protected]>
AuthorDate: Wed Aug 26 16:27:08 2026 -0400
TIKA-4835 -- tika-core primitives for spill-less parsing (#3073)
---
CHANGES.txt | 28 ++
docs/modules/ROOT/pages/pipes/configuration.adoc | 27 +-
.../src/main/java/org/apache/tika/gui/TikaGUI.java | 4 +-
.../apache/tika/digest/BufferingDigestSink.java | 108 +++++++
.../org/apache/tika/digest/CompositeDigester.java | 54 ++++
.../java/org/apache/tika/digest/DigestHelper.java | 41 ++-
.../java/org/apache/tika/digest/DigestSink.java | 82 +++++
.../main/java/org/apache/tika/digest/Digester.java | 19 ++
.../apache/tika/digest/InputStreamDigester.java | 30 ++
.../java/org/apache/tika/io/ByteArraySource.java | 12 +-
.../org/apache/tika/io/CachingInputStream.java | 14 +-
.../java/org/apache/tika/io/CachingSource.java | 35 +-
.../main/java/org/apache/tika/io/FileSource.java | 5 +
.../apache/tika/io/MemorySeekableByteChannel.java | 7 +
.../java/org/apache/tika/io/ReopenableSource.java | 5 +
.../main/java/org/apache/tika/io/StreamCache.java | 5 +
.../org/apache/tika/io/TemporaryResources.java | 47 ++-
.../java/org/apache/tika/io/TikaInputSource.java | 9 +
.../java/org/apache/tika/io/TikaInputStream.java | 42 ++-
.../org/apache/tika/digest/DigestHelperTest.java | 128 ++++++++
.../org/apache/tika/digest/DigestSinkTest.java | 356 +++++++++++++++++++++
.../apache/tika/digest/FailingTestTranslator.java | 63 ++++
.../apache/tika/io/InMemoryContentViewTest.java | 266 +++++++++++++++
.../org/apache/tika/io/StreamCacheBudgetTest.java | 59 ++++
.../org/apache/tika/io/TemporaryResourcesTest.java | 32 ++
....apache.tika.extractor.EmbeddedStreamTranslator | 15 +
.../apache/tika/pipes/core/server/PipesServer.java | 40 ++-
.../core/server/CacheMemoryBudgetSeedingTest.java | 24 +-
.../apache/tika/pipes/core/FontCacheWarmer.java | 55 ++++
.../tika/pipes/core/SharedServerModeTest.java | 2 +
.../pipes/core/async/AsyncChaosMonkeyTest.java | 2 +
.../org.junit.jupiter.api.extension.Extension | 15 +
.../resources/configs/tika-config-passback.json | 2 +-
.../resources/configs/tika-config-truncate.json | 61 ----
.../resources/configs/tika-config-uppercasing.json | 2 +-
.../configs/tika-config-write-limiter.json | 2 +-
.../src/test/resources/junit-platform.properties | 17 +
.../server/core/TikaServerIntegrationTest.java | 4 +-
38 files changed, 1565 insertions(+), 154 deletions(-)
diff --git a/CHANGES.txt b/CHANGES.txt
index d7907ba138..5ae3270114 100644
--- a/CHANGES.txt
+++ b/CHANGES.txt
@@ -1,5 +1,33 @@
Release 4.1.0 - unreleased
+ * tika-pipes: the cache memory budget (how much rewindable content a forked
+ worker keeps in memory before spilling to disk; new since 4.0.0, which had
+ no budget at all) defaults to a quarter of the fork's heap, so raising
+ -Xmx raises it. It is one pool per forked JVM shared by all of its
threads.
+ -Dtika.pipes.cacheMemoryBudgetBytes in forkedJvmArgs overrides it (below
+ the quarter-heap ceiling; <=0 disables); the fork logs the value and its
+ source at startup. TikaInputStream.hasFile() now also reports content the
+ stream cache spilled on its own, not only content a getPath() call put on
+ disk; note getPath() may still have to drain the rest of the source into
+ that file. TikaInputStream.toString() no longer forces a spill, so logging
+ or debugger-inspecting a stream is side-effect-free.
+ TikaInputStream.inMemoryContent(channel) gives a zero-copy read-only view
+ of cached content for consumers that need random access. Digester
+ gains digestSink(), a DigestSink that digests as it is written; nothing is
+ written to the metadata unless the producer calls commit(), so any failure
+ -- exception, Error, or a producer that closes the sink itself --
publishes
+ no digest rather than a digest of the bytes that happened to arrive. A
+ translator that claims a stream and writes nothing likewise publishes
+ nothing: embedded PST mail items, whose translator is still a stub, no
+ longer carry the digest of zero bytes (the same value for every one of
+ them) and instead carry no digest at all. DigestHelper uses it for
+ translated embedded streams, which no longer touch a temp file when the
+ digester implements digestSink (all of Tika's do; one that only implements
+ digest() still buffers).
+ TemporaryResources.closeAll(Closeable...) closes every argument even when
+ one throws unchecked; TemporaryResources, CachingSource,
CachingInputStream
+ and CompositeDigester use it (TIKA-4835).
+
* Documentation: corrected a batch of pages and javadocs that contradicted
the code. Notably: the ES/OpenSearch attachmentStrategy has no default
(unset means embedded documents get neither the parent field nor the
diff --git a/docs/modules/ROOT/pages/pipes/configuration.adoc
b/docs/modules/ROOT/pages/pipes/configuration.adoc
index 02a309661e..f356d23c8c 100644
--- a/docs/modules/ROOT/pages/pipes/configuration.adoc
+++ b/docs/modules/ROOT/pages/pipes/configuration.adoc
@@ -68,23 +68,28 @@ rather than disk, and `/dev/shm` is commonly sized at half
of RAM -- so if you p
=== Embedded-object cache memory budget
-Each fork holds a process-wide in-memory budget for stream caching (chiefly
the rewind
-buffers used when digesting embedded documents), so small embedded objects
stay in RAM
-instead of spilling to a temp file at the per-object 1MB threshold. The
default is 256MB per
-fork, clamped to a quarter of the fork's max heap; the effective value is
logged at fork
-startup. Tune it with a system property in `forkedJvmArgs` (plain bytes, no
unit suffix;
-`<=0` disables the budget and restores the per-object threshold):
+Each fork holds a process-wide in-memory budget for stream caching (the rewind
buffers used
+when digesting embedded documents, and the in-memory views parsers read
random-access
+content from), so small embedded objects stay in RAM instead of spilling to a
temp file at
+the per-object 1MB threshold. The default is a quarter of the fork's max heap,
so raising
+`-Xmx` raises it; a value set via the system property is clamped to that same
quarter-heap
+ceiling. The effective value and where it came from are logged at fork
startup. Tune it with
+a system property in `forkedJvmArgs` (plain bytes, no unit suffix; `<=0`
disables the
+budget and restores the per-object threshold):
[source,json]
----
"forkedJvmArgs": ["-Xmx1g", "-Dtika.pipes.cacheMemoryBudgetBytes=134217728"]
----
-Size `-Xmx` with the budget in mind: the budget is additional heap the fork
may use on top
-of its parsing working set, and in per-client mode every fork holds its own
budget
-(`numClients` x budget in aggregate). In shared-server mode all concurrent
parses in the
-single forked server share one budget, so each in-flight document gets a
smaller slice of
-the same value.
+Size `-Xmx` so that three quarters of it covers the parsing working set; the
remaining
+quarter is what the budget may hold. It is a ceiling on live cached bytes, not
a
+preallocation. In per-client mode every fork holds its own budget
(`numClients` x a quarter
+of each fork's heap in aggregate). In shared-server mode all concurrent parses
in the single
+forked server share one pool, so a single large document can take most of it
and push its
+siblings to disk for a while. With no explicit heap flag the fork's heap is
sized from the
+host (see `forkedJvmArgs`), and the budget scales with it -- on a large host
that is
+several GB per fork by default.
== Timeouts
diff --git a/tika-app/src/main/java/org/apache/tika/gui/TikaGUI.java
b/tika-app/src/main/java/org/apache/tika/gui/TikaGUI.java
index 9efc00a489..a1634a75f0 100644
--- a/tika-app/src/main/java/org/apache/tika/gui/TikaGUI.java
+++ b/tika-app/src/main/java/org/apache/tika/gui/TikaGUI.java
@@ -338,7 +338,9 @@ public class TikaGUI extends JFrame implements
ActionListener, HyperlinkListener
context.set(DocumentSelector.class, new ImageDocumentSelector());
int mark = -1;
- if (tis.hasFile()) {
+ // hasLength(), not hasFile(): getLength() on a stream whose length is
unknown spools
+ // the whole stream to measure it
+ if (tis.hasLength()) {
mark = (int) tis.getLength();
}
if (mark == -1) {
diff --git
a/tika-core/src/main/java/org/apache/tika/digest/BufferingDigestSink.java
b/tika-core/src/main/java/org/apache/tika/digest/BufferingDigestSink.java
new file mode 100644
index 0000000000..4cc8bee603
--- /dev/null
+++ b/tika-core/src/main/java/org/apache/tika/digest/BufferingDigestSink.java
@@ -0,0 +1,108 @@
+/*
+ * 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.digest;
+
+import java.io.IOException;
+import java.nio.file.Files;
+import java.nio.file.Path;
+
+import org.apache.commons.io.output.DeferredFileOutputStream;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import org.apache.tika.io.TemporaryResources;
+import org.apache.tika.io.TikaInputStream;
+import org.apache.tika.metadata.Metadata;
+import org.apache.tika.parser.ParseContext;
+
+/**
+ * The {@link Digester#digestSink} default for digesters that only implement
+ * {@link Digester#digest}: buffers the written bytes (memory below the
threshold, a temp file
+ * above) and runs the pull-style digest over them on close.
+ */
+class BufferingDigestSink extends DigestSink {
+
+ private static final Logger LOG =
LoggerFactory.getLogger(BufferingDigestSink.class);
+
+ static final int MEMORY_THRESHOLD = 1024 * 1024;
+
+ private final Digester digester;
+ private final Metadata metadata;
+ private final ParseContext context;
+ // the DeferredFileOutputStream is eager; the temp file behind it is
created only if
+ // the content crosses the threshold
+ private final DeferredFileOutputStream buffer =
DeferredFileOutputStream.builder()
+ .setThreshold(MEMORY_THRESHOLD)
+ .setPrefix("apache-tika-")
+ .setSuffix(".tmp")
+ .get();
+
+ BufferingDigestSink(Digester digester, Metadata metadata, ParseContext
context) {
+ this.digester = digester;
+ this.metadata = metadata;
+ this.context = context;
+ }
+
+ @Override
+ public void write(int b) throws IOException {
+ ensureOpen();
+ buffer.write(b);
+ }
+
+ @Override
+ public void write(byte[] b, int off, int len) throws IOException {
+ ensureOpen();
+ buffer.write(b, off, len);
+ }
+
+ @Override
+ protected void finish(boolean publish) throws IOException {
+ // read before close(): getPath() is non-null exactly when a file was
created, and
+ // unlike isInMemory() it is a field read that cannot throw and cannot
be stale
+ Path spilled = buffer.getPath();
+ try {
+ buffer.close();
+ if (!publish) {
+ return;
+ }
+ // A re-openable source: the pull digester's
enableRewind()/rewind() re-open
+ // instead of copying the content into a second cache.
+ Path path = spilled;
+ try (TemporaryResources tmp = new TemporaryResources();
+ TikaInputStream tis = path == null ?
+ TikaInputStream.get(buffer::toInputStream, tmp, null)
:
+ TikaInputStream.get(() -> Files.newInputStream(path),
tmp, null)) {
+ digester.digest(tis, metadata, context);
+ }
+ } finally {
+ deleteQuietly(spilled);
+ }
+ }
+
+ // a failed delete must not mask a digest that succeeded
+ private static void deleteQuietly(Path spilled) {
+ if (spilled == null) {
+ return;
+ }
+ try {
+ Files.deleteIfExists(spilled);
+ } catch (IOException | RuntimeException e) {
+ LOG.warn("could not delete {}; will delete on exit", spilled, e);
+ spilled.toFile().deleteOnExit();
+ }
+ }
+}
diff --git
a/tika-core/src/main/java/org/apache/tika/digest/CompositeDigester.java
b/tika-core/src/main/java/org/apache/tika/digest/CompositeDigester.java
index 445cd51749..7cd7b7a6c1 100644
--- a/tika-core/src/main/java/org/apache/tika/digest/CompositeDigester.java
+++ b/tika-core/src/main/java/org/apache/tika/digest/CompositeDigester.java
@@ -18,6 +18,7 @@ package org.apache.tika.digest;
import java.io.IOException;
+import org.apache.tika.io.TemporaryResources;
import org.apache.tika.io.TikaInputStream;
import org.apache.tika.metadata.Metadata;
import org.apache.tika.parser.ParseContext;
@@ -37,4 +38,57 @@ public class CompositeDigester implements Digester {
digester.digest(tis, m, parseContext);
}
}
+
+ /** Fans each write out to every child's sink; commit and close reach them
all. */
+ @Override
+ public DigestSink digestSink(Metadata m, ParseContext parseContext) throws
IOException {
+ DigestSink[] sinks = new DigestSink[digesters.length];
+ try {
+ for (int i = 0; i < digesters.length; i++) {
+ sinks[i] = digesters[i].digestSink(m, parseContext);
+ }
+ } catch (Throwable e) {
+ // uncommitted, so closing publishes nothing
+ try {
+ TemporaryResources.closeAll(sinks);
+ } catch (Throwable t) {
+ if (t != e) {
+ e.addSuppressed(t);
+ }
+ }
+ throw e;
+ }
+ return new DigestSink() {
+ @Override
+ public void write(int b) throws IOException {
+ ensureOpen();
+ for (DigestSink sink : sinks) {
+ sink.write(b);
+ }
+ }
+
+ @Override
+ public void write(byte[] b, int off, int len) throws IOException {
+ ensureOpen();
+ for (DigestSink sink : sinks) {
+ sink.write(b, off, len);
+ }
+ }
+
+ @Override
+ protected void finish(boolean publish) throws IOException {
+ try {
+ if (publish) {
+ for (DigestSink sink : sinks) {
+ sink.commit();
+ }
+ }
+ } finally {
+ // a child that refuses to commit must not strand the
others: this sink is
+ // already marked closed, so nothing else will ever close
them
+ TemporaryResources.closeAll(sinks);
+ }
+ }
+ };
+ }
}
diff --git a/tika-core/src/main/java/org/apache/tika/digest/DigestHelper.java
b/tika-core/src/main/java/org/apache/tika/digest/DigestHelper.java
index 359ebf9f27..d2cd1e68b7 100644
--- a/tika-core/src/main/java/org/apache/tika/digest/DigestHelper.java
+++ b/tika-core/src/main/java/org/apache/tika/digest/DigestHelper.java
@@ -17,14 +17,13 @@
package org.apache.tika.digest;
import java.io.IOException;
-import java.io.OutputStream;
-import java.nio.file.Files;
-import java.nio.file.Path;
+
+import org.apache.commons.io.output.CloseShieldOutputStream;
+import org.apache.commons.io.output.CountingOutputStream;
import org.apache.tika.extractor.DefaultEmbeddedStreamTranslator;
import org.apache.tika.extractor.EmbeddedStreamTranslator;
import org.apache.tika.io.CacheMemoryBudget;
-import org.apache.tika.io.TemporaryResources;
import org.apache.tika.io.TikaInputStream;
import org.apache.tika.metadata.Metadata;
import org.apache.tika.metadata.TikaCoreProperties;
@@ -83,15 +82,35 @@ public class DigestHelper {
// The translator consumes `tis` (e.g. OLE2), so enableRewind() before
and rewind()
// after -- otherwise the caller would see an exhausted stream.
+ // Safe to tee here: the translator only writes, so no mark/reset/skip
can
+ // desynchronize the digest from the bytes.
if (EMBEDDED_STREAM_TRANSLATOR.shouldTranslate(tis, metadata)) {
tis.enableRewind(context.get(CacheMemoryBudget.class));
- try (TemporaryResources tmp = new TemporaryResources()) {
- Path tmpBytes = tmp.createTempFile();
- try (OutputStream os = Files.newOutputStream(tmpBytes)) {
- EMBEDDED_STREAM_TRANSLATOR.translate(tis, metadata, os);
- }
- try (TikaInputStream translated =
TikaInputStream.get(tmpBytes)) {
- digester.digest(translated, metadata, context);
+ try {
+ DigestSink sink = digester.digestSink(metadata, context);
+ try {
+ // Close-shielded so a translator that closes the stream
cannot publish on
+ // our behalf, and counted because "returned normally" is
not "produced the
+ // content": a translator that claims the stream and
writes nothing (see
+ // PSTEmailStreamTranslator) would otherwise publish the
digest of zero
+ // bytes -- the same value for every such object.
+ CountingOutputStream counted =
+ new
CountingOutputStream(CloseShieldOutputStream.wrap(sink));
+ EMBEDDED_STREAM_TRANSLATOR.translate(tis, metadata,
counted);
+ if (counted.getByteCount() > 0) {
+ sink.commit();
+ }
+ sink.close();
+ } catch (Throwable t) {
+ // close() can fail too; that must not erase why the
translation failed
+ try {
+ sink.close();
+ } catch (Throwable closeFailure) {
+ if (closeFailure != t) {
+ t.addSuppressed(closeFailure);
+ }
+ }
+ throw t;
}
} finally {
tis.rewind();
diff --git a/tika-core/src/main/java/org/apache/tika/digest/DigestSink.java
b/tika-core/src/main/java/org/apache/tika/digest/DigestSink.java
new file mode 100644
index 0000000000..5cb44ac8d0
--- /dev/null
+++ b/tika-core/src/main/java/org/apache/tika/digest/DigestSink.java
@@ -0,0 +1,82 @@
+/*
+ * 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.digest;
+
+import java.io.IOException;
+import java.io.OutputStream;
+
+/**
+ * A sink returned by {@link Digester#digestSink}: bytes written to it are
digested, and the
+ * value(s) reach the metadata only if the producer calls {@link #commit()}.
+ * <p>
+ * Commit is explicit because the alternative -- publish on close unless
something remembered
+ * to cancel -- makes every unanticipated exit (a checked exception, an {@link
Error}, a
+ * producer that closes the sink itself) publish a digest of whatever bytes
happened to arrive.
+ * A digest of a partial write is worse than no digest: it is wrong, and it is
stably wrong,
+ * so every failure of the same shape produces the same plausible value.
+ * <p>
+ * Contract:
+ * <ul>
+ * <li>Write the content, then {@link #commit()}, then {@link #close()} --
close in a
+ * {@code finally} or via try-with-resources.</li>
+ * <li>{@link #close()} releases resources either way, and publishes only if
+ * {@code commit()} was called first. It is idempotent.</li>
+ * <li>{@link #commit()} after {@code close()} throws {@link
IllegalStateException}: the
+ * chance to publish is gone, and failing loudly beats a silently
missing digest.</li>
+ * <li>Writing after {@code close()} throws {@link IOException}.</li>
+ * <li>All methods must be called from the producing thread.</li>
+ * </ul>
+ */
+public abstract class DigestSink extends OutputStream {
+
+ private boolean committed;
+ private boolean closed;
+
+ /**
+ * Marks the content complete, so {@link #close()} publishes it.
Idempotent.
+ *
+ * @throws IllegalStateException if this sink is already closed
+ */
+ public final void commit() {
+ if (closed) {
+ throw new IllegalStateException("cannot commit a closed digest
sink");
+ }
+ committed = true;
+ }
+
+ @Override
+ public final void close() throws IOException {
+ if (closed) {
+ return;
+ }
+ closed = true;
+ finish(committed);
+ }
+
+ /**
+ * Called once, from the first {@link #close()}. Release resources either
way; set the
+ * metadata only when {@code publish} is true.
+ */
+ protected abstract void finish(boolean publish) throws IOException;
+
+ /** For subclasses' write methods. */
+ protected final void ensureOpen() throws IOException {
+ if (closed) {
+ throw new IOException("digest sink is closed");
+ }
+ }
+}
diff --git a/tika-core/src/main/java/org/apache/tika/digest/Digester.java
b/tika-core/src/main/java/org/apache/tika/digest/Digester.java
index 133d5dce09..59ffb4aa1d 100644
--- a/tika-core/src/main/java/org/apache/tika/digest/Digester.java
+++ b/tika-core/src/main/java/org/apache/tika/digest/Digester.java
@@ -42,4 +42,23 @@ public interface Digester {
* @throws IOException on I/O error
*/
void digest(TikaInputStream tis, Metadata m, ParseContext parseContext)
throws IOException;
+
+ /**
+ * A sink that digests whatever is written to it and sets the value(s) in
the metadata
+ * when committed and closed. For producers that only write (an
embedded-stream
+ * translator, say) this digests the bytes as they are produced, with no
buffer and no
+ * temp file. See {@link DigestSink} for the commit/close contract.
+ * <p>
+ * The default buffers what is written -- in memory below a fixed
threshold, in a temp
+ * file above it -- and runs {@link #digest} over it when committed, so
one that only
+ * overrides {@code digest} keeps working unchanged. It does <em>not</em>
get the
+ * streaming behaviour: override this method to provide it.
+ *
+ * @param m Metadata the values are set on when the sink is
committed
+ * @param parseContext ParseContext
+ * @return a sink; the caller must close it, and the values are set only
if it was committed
+ */
+ default DigestSink digestSink(Metadata m, ParseContext parseContext)
throws IOException {
+ return new BufferingDigestSink(this, m, parseContext);
+ }
}
diff --git
a/tika-core/src/main/java/org/apache/tika/digest/InputStreamDigester.java
b/tika-core/src/main/java/org/apache/tika/digest/InputStreamDigester.java
index 00941602d3..ec46af2839 100644
--- a/tika-core/src/main/java/org/apache/tika/digest/InputStreamDigester.java
+++ b/tika-core/src/main/java/org/apache/tika/digest/InputStreamDigester.java
@@ -137,4 +137,34 @@ public class InputStreamDigester implements Digester {
tis.rewind();
}
+ /** Streams: every write goes straight into the MessageDigest; close sets
the values. */
+ @Override
+ public DigestSink digestSink(Metadata metadata, ParseContext parseContext)
{
+ MessageDigest messageDigest = newMessageDigest();
+ return new DigestSink() {
+ private long total;
+
+ @Override
+ public void write(int b) throws IOException {
+ ensureOpen();
+ messageDigest.update((byte) b);
+ total++;
+ }
+
+ @Override
+ public void write(byte[] b, int off, int len) throws IOException {
+ ensureOpen();
+ messageDigest.update(b, off, len);
+ total += len;
+ }
+
+ @Override
+ protected void finish(boolean publish) {
+ if (publish) {
+ setContentLength(total, metadata);
+ metadata.set(metadataProperty,
encoder.encode(messageDigest.digest()));
+ }
+ }
+ };
+ }
}
diff --git a/tika-core/src/main/java/org/apache/tika/io/ByteArraySource.java
b/tika-core/src/main/java/org/apache/tika/io/ByteArraySource.java
index 72c94f728e..d5eee0b07d 100644
--- a/tika-core/src/main/java/org/apache/tika/io/ByteArraySource.java
+++ b/tika-core/src/main/java/org/apache/tika/io/ByteArraySource.java
@@ -98,6 +98,11 @@ class ByteArraySource extends InputStream implements
TikaInputSource {
this.position = (int) newPosition;
}
+ @Override
+ public Path materializedPath() {
+ return spilledPath;
+ }
+
@Override
public boolean hasPath() {
return spilledPath != null;
@@ -107,10 +112,13 @@ class ByteArraySource extends InputStream implements
TikaInputSource {
public Path getPath(String suffix) throws IOException {
if (spilledPath == null) {
// Spill to temp file on first call
- spilledPath = tmp.createTempFile(suffix);
- try (OutputStream out = Files.newOutputStream(spilledPath)) {
+ // assign only on success: a failed write must not leave a
truncated file that
+ // a later getPath() would hand back as if it were complete
+ Path p = tmp.createTempFile(suffix);
+ try (OutputStream out = Files.newOutputStream(p)) {
out.write(data, 0, length);
}
+ spilledPath = p;
}
return spilledPath;
}
diff --git a/tika-core/src/main/java/org/apache/tika/io/CachingInputStream.java
b/tika-core/src/main/java/org/apache/tika/io/CachingInputStream.java
index ff97c31485..4d743f00e1 100644
--- a/tika-core/src/main/java/org/apache/tika/io/CachingInputStream.java
+++ b/tika-core/src/main/java/org/apache/tika/io/CachingInputStream.java
@@ -201,10 +201,20 @@ class CachingInputStream extends InputStream {
return cache.isFileBacked();
}
+ /**
+ * The spill file once it holds the whole source, or null. A mid-stream
spill file is
+ * still growing, so naming it would be misleading; this is for
diagnostics only and the
+ * very tail may still be buffered until {@link #spillToFile(String)}
flushes it.
+ */
+ Path completeSpillFile() {
+ return sourceExhausted ? cache.spillFilePath() : null;
+ }
+
@Override
public void close() throws IOException {
- source.close();
- cache.close();
+ // the cache is registered with tmp only once it spills, so a throwing
source
+ // must not skip it: its budget reservation would leak for the fork's
life
+ TemporaryResources.closeAll(source, cache);
}
/**
diff --git a/tika-core/src/main/java/org/apache/tika/io/CachingSource.java
b/tika-core/src/main/java/org/apache/tika/io/CachingSource.java
index 70d4fc6a63..8b424696c7 100644
--- a/tika-core/src/main/java/org/apache/tika/io/CachingSource.java
+++ b/tika-core/src/main/java/org/apache/tika/io/CachingSource.java
@@ -17,7 +17,6 @@
package org.apache.tika.io;
import java.io.BufferedInputStream;
-import java.io.Closeable;
import java.io.IOException;
import java.io.InputStream;
import java.nio.channels.FileChannel;
@@ -258,9 +257,20 @@ class CachingSource extends InputStream implements
TikaInputSource {
}
}
+ @Override
+ public Path materializedPath() {
+ if (spilledPath != null) {
+ return spilledPath;
+ }
+ // the cache spills on its own; offer that file only once it holds the
whole source,
+ // since a mid-stream spill file is still growing
+ return cachingStream == null ? null :
cachingStream.completeSpillFile();
+ }
+
@Override
public boolean hasPath() {
- return spilledPath != null;
+ // the cache can spill on its own, without any getPath() call
+ return spilledPath != null || (cachingStream != null &&
cachingStream.isFileBacked());
}
@Override
@@ -326,24 +336,7 @@ class CachingSource extends InputStream implements
TikaInputSource {
@Override
public void close() throws IOException {
- IOException exception = null;
- for (Closeable closeable :
- new Closeable[]{this::closeFileStream, cachingStream,
passthroughStream, spilledSource}) {
- if (closeable == null) {
- continue;
- }
- try {
- closeable.close();
- } catch (IOException e) {
- if (exception == null) {
- exception = e;
- } else {
- exception.addSuppressed(e);
- }
- }
- }
- if (exception != null) {
- throw exception;
- }
+ TemporaryResources.closeAll(this::closeFileStream, cachingStream,
passthroughStream,
+ spilledSource);
}
}
diff --git a/tika-core/src/main/java/org/apache/tika/io/FileSource.java
b/tika-core/src/main/java/org/apache/tika/io/FileSource.java
index fbd06d7a0d..257cbacbc3 100644
--- a/tika-core/src/main/java/org/apache/tika/io/FileSource.java
+++ b/tika-core/src/main/java/org/apache/tika/io/FileSource.java
@@ -99,6 +99,11 @@ class FileSource extends InputStream implements
TikaInputSource {
this.position = newPosition;
}
+ @Override
+ public Path materializedPath() {
+ return path;
+ }
+
@Override
public boolean hasPath() {
return true;
diff --git
a/tika-core/src/main/java/org/apache/tika/io/MemorySeekableByteChannel.java
b/tika-core/src/main/java/org/apache/tika/io/MemorySeekableByteChannel.java
index 71d71d93ed..56da4e77a0 100644
--- a/tika-core/src/main/java/org/apache/tika/io/MemorySeekableByteChannel.java
+++ b/tika-core/src/main/java/org/apache/tika/io/MemorySeekableByteChannel.java
@@ -46,6 +46,13 @@ class MemorySeekableByteChannel implements
SeekableByteChannel {
this.onClose = onClose;
}
+ /** Read-only view of the whole content, independent of this channel's
position. */
+ ByteBuffer buffer() throws IOException {
+ ensureOpen();
+ // slice(): capacity == limit == length, so clear() cannot expose the
array's slack
+ return ByteBuffer.wrap(data, 0, length).slice().asReadOnlyBuffer();
+ }
+
@Override
public int read(ByteBuffer dst) throws IOException {
ensureOpen();
diff --git a/tika-core/src/main/java/org/apache/tika/io/ReopenableSource.java
b/tika-core/src/main/java/org/apache/tika/io/ReopenableSource.java
index f758f1a295..332e56eba3 100644
--- a/tika-core/src/main/java/org/apache/tika/io/ReopenableSource.java
+++ b/tika-core/src/main/java/org/apache/tika/io/ReopenableSource.java
@@ -134,6 +134,11 @@ class ReopenableSource extends InputStream implements
TikaInputSource {
this.position = newPosition;
}
+ @Override
+ public Path materializedPath() {
+ return spilledPath;
+ }
+
@Override
public boolean hasPath() {
return spilledPath != null;
diff --git a/tika-core/src/main/java/org/apache/tika/io/StreamCache.java
b/tika-core/src/main/java/org/apache/tika/io/StreamCache.java
index d3e4aaae7d..606910abcd 100644
--- a/tika-core/src/main/java/org/apache/tika/io/StreamCache.java
+++ b/tika-core/src/main/java/org/apache/tika/io/StreamCache.java
@@ -308,6 +308,11 @@ class StreamCache implements Closeable {
/**
* Whether the cache has spilled to a file.
*/
+ /** The spill file, or null while the content is still in memory. Never
creates one. */
+ Path spillFilePath() {
+ return spillFile;
+ }
+
boolean isFileBacked() {
return spillFile != null;
}
diff --git a/tika-core/src/main/java/org/apache/tika/io/TemporaryResources.java
b/tika-core/src/main/java/org/apache/tika/io/TemporaryResources.java
index c1565ab86d..21bc80ec01 100644
--- a/tika-core/src/main/java/org/apache/tika/io/TemporaryResources.java
+++ b/tika-core/src/main/java/org/apache/tika/io/TemporaryResources.java
@@ -91,7 +91,7 @@ public class TemporaryResources implements Closeable {
addResource(() -> {
try {
Files.delete(path);
- } catch (IOException e) {
+ } catch (IOException | RuntimeException e) {
// delete when exit if current delete fail
LOG.warn("delete tmp file fail, will delete it on exit");
path.toFile().deleteOnExit();
@@ -169,24 +169,41 @@ public class TemporaryResources implements Closeable {
* could not be closed
*/
public void close() throws IOException {
- // Release all resources and keep track of any exceptions
- IOException exception = null;
- for (Closeable resource : resources) {
+ try {
+ closeAll(resources.toArray(new Closeable[0]));
+ } finally {
+ resources.clear();
+ }
+ }
+
+ /**
+ * Closes every closeable even when one throws, unchecked included. The
first failure
+ * propagates with the rest attached as suppressed; nulls are skipped.
+ */
+ public static void closeAll(Closeable... closeables) throws IOException {
+ Throwable first = null;
+ for (Closeable closeable : closeables) {
+ if (closeable == null) {
+ continue;
+ }
try {
- resource.close();
- } catch (IOException e) {
- if (exception == null) {
- exception = e;
- } else {
- exception.addSuppressed(e);
+ closeable.close();
+ } catch (Throwable t) {
+ if (first == null) {
+ first = t;
+ } else if (t != first) { // addSuppressed rejects
self-suppression
+ first.addSuppressed(t);
}
}
}
- resources.clear();
-
- // Throw any exceptions that were captured from above
- if (exception != null) {
- throw exception;
+ if (first instanceof IOException) {
+ throw (IOException) first;
+ } else if (first instanceof RuntimeException) {
+ throw (RuntimeException) first;
+ } else if (first instanceof Error) {
+ throw (Error) first;
+ } else if (first != null) {
+ throw new IOException(first);
}
}
diff --git a/tika-core/src/main/java/org/apache/tika/io/TikaInputSource.java
b/tika-core/src/main/java/org/apache/tika/io/TikaInputSource.java
index 53d8942399..9e7b4a8d30 100644
--- a/tika-core/src/main/java/org/apache/tika/io/TikaInputSource.java
+++ b/tika-core/src/main/java/org/apache/tika/io/TikaInputSource.java
@@ -42,6 +42,15 @@ interface TikaInputSource extends Closeable {
*/
boolean hasPath();
+ /**
+ * The file this source is already associated with, or {@code null}. Never
creates one
+ * and never reads from the source, so it is safe where {@link
#getPath(String)} is not
+ * -- logging, diagnostics, {@code toString()}. For diagnostics only: the
file may since
+ * have been deleted, and for a stream cache the last bytes may still be
buffered until
+ * {@link #getPath(String)} completes it.
+ */
+ Path materializedPath();
+
/**
* Gets the file path, potentially spilling to a temp file if needed.
* @param suffix file suffix for temp files
diff --git a/tika-core/src/main/java/org/apache/tika/io/TikaInputStream.java
b/tika-core/src/main/java/org/apache/tika/io/TikaInputStream.java
index 6de71cfa8e..b58cd548db 100644
--- a/tika-core/src/main/java/org/apache/tika/io/TikaInputStream.java
+++ b/tika-core/src/main/java/org/apache/tika/io/TikaInputStream.java
@@ -25,6 +25,7 @@ import java.net.URI;
import java.net.URISyntaxException;
import java.net.URL;
import java.net.URLConnection;
+import java.nio.ByteBuffer;
import java.nio.channels.FileChannel;
import java.nio.channels.SeekableByteChannel;
import java.nio.file.Files;
@@ -400,6 +401,14 @@ public class TikaInputStream extends TaggedInputStream {
tmp.addResource(closeable);
}
+ /**
+ * Whether the content is already on disk: a real file, or a stream cache
that spilled.
+ * {@link #getPath()} then returns that file without re-copying anything
already written
+ * to it -- but it is not free, and it is not a getter: for a cache that
spilled
+ * mid-stream it first drains the rest of the source into the file,
switches this stream
+ * to reading from that file, and sets {@code Content-Length} on the
Metadata this stream
+ * was created with. Use {@link #hasLength()} if you only need the size.
+ */
public boolean hasFile() {
TikaInputSource source = inputSource();
return source != null && source.hasPath();
@@ -552,18 +561,33 @@ public class TikaInputStream extends TaggedInputStream {
return source.getSeekableByteChannel();
}
+ /**
+ * Zero-copy, read-only view of the content behind a channel from
+ * {@link #getSeekableByteChannel()}, or {@code null} when that content is
on disk. Lets a
+ * consumer that wants random access (PDFBox, metadata-extractor) read
what is already in
+ * memory without a second copy; when this returns null the caller should
use the file.
+ * <p>
+ * The view aliases the cache's own array and is valid exactly while
{@code channel} is
+ * open: the channel pins the array, and the content is fully drained
before any channel
+ * is handed out. Keep the channel open for as long as the view is in use,
then close it
+ * -- a view that outlives its channel still reads correctly but is no
longer counted
+ * against the memory budget.
+ */
+ public static ByteBuffer inMemoryContent(SeekableByteChannel channel)
throws IOException {
+ if (channel instanceof MemorySeekableByteChannel) {
+ return ((MemorySeekableByteChannel) channel).buffer();
+ }
+ return null;
+ }
+
@Override
public String toString() {
String str = "TikaInputStream of ";
- if (hasFile()) {
- try {
- str += getPath().toString();
- } catch (IOException e) {
- str += "unknown path";
- }
- } else {
- str += in.toString();
- }
+ // materializedPath(), never getPath(): on a spilled cache the latter
drains the
+ // source, reopens it and writes metadata -- toString() must not do
that
+ TikaInputSource source = inputSource();
+ Path materialized = source == null ? null : source.materializedPath();
+ str += materialized != null ? materialized.toString() : in.toString();
if (openContainer != null) {
str += " (in " + openContainer + ")";
}
diff --git
a/tika-core/src/test/java/org/apache/tika/digest/DigestHelperTest.java
b/tika-core/src/test/java/org/apache/tika/digest/DigestHelperTest.java
new file mode 100644
index 0000000000..7e7ad245c4
--- /dev/null
+++ b/tika-core/src/test/java/org/apache/tika/digest/DigestHelperTest.java
@@ -0,0 +1,128 @@
+/*
+ * 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.digest;
+
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+
+import java.io.IOException;
+import java.nio.charset.StandardCharsets;
+import java.util.HexFormat;
+
+import org.junit.jupiter.api.Test;
+
+import org.apache.tika.io.TikaInputStream;
+import org.apache.tika.metadata.HttpHeaders;
+import org.apache.tika.metadata.Metadata;
+import org.apache.tika.parser.ParseContext;
+
+/** The translate branch of DigestHelper, through the service-registered test
translator. */
+public class DigestHelperTest {
+
+ private static final String MD5_KEY = "tk:digest:MD5";
+ private static final Encoder HEX = bytes ->
HexFormat.of().formatHex(bytes);
+ private static final byte[] SOURCE = "hello, digest
helper".getBytes(StandardCharsets.US_ASCII);
+ private static final byte[] TRANSLATED = "HELLO, DIGEST
HELPER".getBytes(StandardCharsets.US_ASCII);
+
+ private static ParseContext contextWithDigester() {
+ ParseContext context = new ParseContext();
+ context.set(DigesterFactory.class, new DigesterFactory() {
+ @Override
+ public Digester build() {
+ return new InputStreamDigester("MD5", MD5_KEY, HEX);
+ }
+
+ @Override
+ public boolean isSkipContainerDocumentDigest() {
+ return false;
+ }
+ });
+ return context;
+ }
+
+ private static String md5Of(byte[] bytes) throws IOException {
+ Metadata m = new Metadata();
+ try (TikaInputStream tis = TikaInputStream.get(bytes)) {
+ new InputStreamDigester("MD5", MD5_KEY, HEX).digest(tis, m, new
ParseContext());
+ }
+ return m.get(MD5_KEY);
+ }
+
+ @Test
+ public void testDigestIsOfTheTranslatedBytesAndStreamIsRewound() throws
Exception {
+ Metadata metadata = new Metadata();
+ metadata.set(FailingTestTranslator.MODE, "upper");
+ try (TikaInputStream tis = TikaInputStream.get(SOURCE, new
Metadata())) {
+ DigestHelper.maybeDigest(tis, metadata, contextWithDigester());
+ assertEquals(md5Of(TRANSLATED), metadata.get(MD5_KEY), "digest of
what the translator wrote");
+ assertEquals(0, tis.getPosition(), "caller gets the stream back at
0");
+ assertArrayEquals(SOURCE, tis.readAllBytes(), "and intact");
+ }
+ }
+
+ /** A failed translation must publish nothing -- not a digest of the
fragment. */
+ @Test
+ public void testFailedTranslationPublishesNoDigest() throws Exception {
+ Metadata metadata = new Metadata();
+ metadata.set(FailingTestTranslator.MODE, "fail");
+ try (TikaInputStream tis = TikaInputStream.get(SOURCE, new
Metadata())) {
+ IOException e = assertThrows(IOException.class,
+ () -> DigestHelper.maybeDigest(tis, metadata,
contextWithDigester()));
+ assertEquals("translator gave up half way", e.getMessage());
+ assertNull(metadata.get(MD5_KEY), "no digest of a partial
translation");
+ assertNull(metadata.get(HttpHeaders.CONTENT_LENGTH));
+ assertEquals(0, tis.getPosition(), "stream still rewound on
failure");
+ }
+ }
+
+ @Test
+ public void testNoTranslationDigestsTheStreamItself() throws Exception {
+ Metadata metadata = new Metadata();
+ try (TikaInputStream tis = TikaInputStream.get(SOURCE, new
Metadata())) {
+ DigestHelper.maybeDigest(tis, metadata, contextWithDigester());
+ assertEquals(md5Of(SOURCE), metadata.get(MD5_KEY));
+ }
+ }
+
+ /**
+ * A translator that claims the stream and writes nothing must publish
nothing. Otherwise
+ * every such object gets the digest of zero bytes -- the same wrong value
for all of them.
+ */
+ @Test
+ public void testTranslatorThatWritesNothingPublishesNoDigest() throws
Exception {
+ Metadata metadata = new Metadata();
+ metadata.set(FailingTestTranslator.MODE, "silent");
+ try (TikaInputStream tis = TikaInputStream.get(SOURCE, new
Metadata())) {
+ DigestHelper.maybeDigest(tis, metadata, contextWithDigester());
+ assertNull(metadata.get(MD5_KEY), "no digest of the zero bytes it
produced");
+ assertNull(metadata.get(HttpHeaders.CONTENT_LENGTH));
+ }
+ }
+
+ /** A translator that closes the sink must not publish on our behalf, nor
break the parse. */
+ @Test
+ public void testTranslatorClosingTheStreamStillDigestsWhatItWrote() throws
Exception {
+ Metadata metadata = new Metadata();
+ metadata.set(FailingTestTranslator.MODE, "closes");
+ try (TikaInputStream tis = TikaInputStream.get(SOURCE, new
Metadata())) {
+ DigestHelper.maybeDigest(tis, metadata, contextWithDigester());
+ assertEquals(md5Of(TRANSLATED), metadata.get(MD5_KEY));
+ }
+ }
+}
diff --git a/tika-core/src/test/java/org/apache/tika/digest/DigestSinkTest.java
b/tika-core/src/test/java/org/apache/tika/digest/DigestSinkTest.java
new file mode 100644
index 0000000000..0390b29cbb
--- /dev/null
+++ b/tika-core/src/test/java/org/apache/tika/digest/DigestSinkTest.java
@@ -0,0 +1,356 @@
+/*
+ * 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.digest;
+
+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.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+
+import java.io.IOException;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.util.HexFormat;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.stream.Stream;
+
+import org.junit.jupiter.api.Test;
+
+import org.apache.tika.io.TikaInputStream;
+import org.apache.tika.metadata.HttpHeaders;
+import org.apache.tika.metadata.Metadata;
+import org.apache.tika.parser.ParseContext;
+
+/**
+ * The sink must produce exactly what the pull-style digest produces for the
same bytes,
+ * through every Digester shape, and must honour the DigestSink contract on
every exit.
+ */
+public class DigestSinkTest {
+
+ private static final String MD5_KEY = "tk:digest:MD5";
+ private static final String SHA_KEY = "tk:digest:SHA-256";
+ private static final Encoder HEX = bytes ->
HexFormat.of().formatHex(bytes);
+
+ private static byte[] data(int len) {
+ byte[] d = new byte[len];
+ for (int i = 0; i < len; i++) {
+ d[i] = (byte) (i * 131 + 17);
+ }
+ return d;
+ }
+
+ private static Digester pullOnly(InputStreamDigester inner) {
+ return (tis, m, ctx) -> inner.digest(tis, m, ctx);
+ }
+
+ private static Metadata viaStream(Digester digester, byte[] data) throws
IOException {
+ Metadata m = new Metadata();
+ try (TikaInputStream tis = TikaInputStream.get(data)) {
+ digester.digest(tis, m, new ParseContext());
+ }
+ return m;
+ }
+
+ /** Writes in awkward chunk sizes, including single bytes, to exercise
both overloads. */
+ private static Metadata viaSink(Digester digester, byte[] data) throws
IOException {
+ Metadata m = new Metadata();
+ try (DigestSink sink = digester.digestSink(m, new ParseContext())) {
+ int pos = 0;
+ int[] chunks = {1, 7, 1000, 1, 65536};
+ int c = 0;
+ while (pos < data.length) {
+ int len = Math.min(chunks[c++ % chunks.length], data.length -
pos);
+ if (len == 1) {
+ sink.write(data[pos]);
+ } else {
+ sink.write(data, pos, len);
+ }
+ pos += len;
+ }
+ sink.commit();
+ }
+ return m;
+ }
+
+ private static long tempFiles() throws IOException {
+ try (Stream<Path> s =
Files.list(Path.of(System.getProperty("java.io.tmpdir")))) {
+ return s.filter(p ->
p.getFileName().toString().startsWith("apache-tika-")).count();
+ }
+ }
+
+ @Test
+ public void testStreamingSinkMatchesPullDigest() throws Exception {
+ Digester d = new InputStreamDigester("MD5", MD5_KEY, HEX);
+ byte[] data = data(200_000);
+ Metadata expected = viaStream(d, data);
+ Metadata actual = viaSink(d, data);
+ assertNotNull(expected.get(MD5_KEY));
+ assertEquals(expected.get(MD5_KEY), actual.get(MD5_KEY));
+ assertEquals(Integer.toString(data.length),
actual.get(HttpHeaders.CONTENT_LENGTH));
+ }
+
+ @Test
+ public void testCompositeFansOut() throws Exception {
+ Digester d = new CompositeDigester(
+ new InputStreamDigester("MD5", MD5_KEY, HEX),
+ new InputStreamDigester("SHA-256", SHA_KEY, HEX));
+ byte[] data = data(50_000);
+ Metadata expected = viaStream(d, data);
+ Metadata actual = viaSink(d, data);
+ assertEquals(expected.get(MD5_KEY), actual.get(MD5_KEY));
+ assertEquals(expected.get(SHA_KEY), actual.get(SHA_KEY));
+ }
+
+ @Test
+ public void testMixedCompositeStreamingAndPullOnly() throws Exception {
+ InputStreamDigester sha = new InputStreamDigester("SHA-256", SHA_KEY,
HEX);
+ Digester d = new CompositeDigester(new InputStreamDigester("MD5",
MD5_KEY, HEX), pullOnly(sha));
+ byte[] data = data(30_000);
+ Metadata expected = viaStream(d, data);
+ Metadata actual = viaSink(d, data);
+ assertEquals(expected.get(MD5_KEY), actual.get(MD5_KEY));
+ assertEquals(expected.get(SHA_KEY), actual.get(SHA_KEY));
+ }
+
+ @Test
+ public void testDefaultBuffersLargeContentThroughATempFile() throws
Exception {
+ InputStreamDigester inner = new InputStreamDigester("SHA-256",
SHA_KEY, HEX);
+ long before = tempFiles();
+ byte[] data = data(BufferingDigestSink.MEMORY_THRESHOLD + 12_345);
+ Metadata m = new Metadata();
+ try (DigestSink sink = pullOnly(inner).digestSink(m, new
ParseContext())) {
+ sink.write(data, 0, data.length);
+ assertEquals(before + 1, tempFiles(), "past the threshold the
content is on disk");
+ sink.commit();
+ }
+ assertEquals(before, tempFiles(), "spill file must be deleted on
close");
+ assertEquals(viaStream(inner, data).get(SHA_KEY), m.get(SHA_KEY));
+ assertEquals(Integer.toString(data.length),
m.get(HttpHeaders.CONTENT_LENGTH));
+ }
+
+ @Test
+ public void testDefaultBuffersSmallContentWithNoTempFile() throws
Exception {
+ InputStreamDigester inner = new InputStreamDigester("MD5", MD5_KEY,
HEX);
+ byte[] data = data(1234);
+ Metadata m = new Metadata();
+ long before = tempFiles();
+ try (DigestSink sink = pullOnly(inner).digestSink(m, new
ParseContext())) {
+ sink.write(data, 0, data.length);
+ assertEquals(before, tempFiles(), "under the threshold nothing may
touch disk");
+ sink.commit();
+ }
+ assertEquals(viaStream(inner, data).get(MD5_KEY), m.get(MD5_KEY));
+ }
+
+ @Test
+ public void testValuesSetOnlyOnCloseAndCloseIsIdempotent() throws
Exception {
+ Digester d = new InputStreamDigester("MD5", MD5_KEY, HEX);
+ Metadata m = new Metadata();
+ DigestSink sink = d.digestSink(m, new ParseContext());
+ sink.write(data(100), 0, 100);
+ assertNull(m.get(MD5_KEY), "writing alone must not publish");
+ sink.commit();
+ assertNull(m.get(MD5_KEY), "committing alone must not publish either");
+ sink.close();
+ String first = m.get(MD5_KEY);
+ assertNotNull(first);
+ sink.close();
+ assertEquals(first, m.get(MD5_KEY));
+ }
+
+ @Test
+ public void testCommitAfterCloseThrows() throws Exception {
+ Digester d = new InputStreamDigester("MD5", MD5_KEY, HEX);
+ Metadata m = new Metadata();
+ DigestSink sink = d.digestSink(m, new ParseContext());
+ sink.write(data(100), 0, 100);
+ sink.close();
+ assertThrows(IllegalStateException.class, sink::commit,
+ "the chance to publish is gone; failing loudly beats a missing
digest");
+ assertNull(m.get(MD5_KEY));
+ }
+
+ @Test
+ public void testUncommittedPublishesNothing() throws Exception {
+ for (Digester d : new Digester[]{
+ new InputStreamDigester("MD5", MD5_KEY, HEX),
+ pullOnly(new InputStreamDigester("MD5", MD5_KEY, HEX)),
+ new CompositeDigester(new InputStreamDigester("MD5", MD5_KEY,
HEX),
+ pullOnly(new InputStreamDigester("SHA-256", SHA_KEY,
HEX)))}) {
+ Metadata m = new Metadata();
+ DigestSink sink = d.digestSink(m, new ParseContext());
+ sink.write(data(5000), 0, 5000);
+ sink.close();
+ assertNull(m.get(MD5_KEY), "uncommitted sink must not publish: " +
d.getClass());
+ assertNull(m.get(SHA_KEY));
+ assertNull(m.get(HttpHeaders.CONTENT_LENGTH));
+ }
+ }
+
+ @Test
+ public void testWriteAfterCloseThrows() throws Exception {
+ for (Digester d : new Digester[]{
+ new InputStreamDigester("MD5", MD5_KEY, HEX),
+ pullOnly(new InputStreamDigester("MD5", MD5_KEY, HEX)),
+ new CompositeDigester(new InputStreamDigester("MD5", MD5_KEY,
HEX))}) {
+ DigestSink closed = d.digestSink(new Metadata(), new
ParseContext());
+ closed.close();
+ assertThrows(IOException.class, () -> closed.write(1));
+ assertThrows(IOException.class, () -> closed.write(new byte[3], 0,
3));
+ }
+ }
+
+ /** A child whose close() throws unchecked must not leave the children
after it open. */
+ @Test
+ public void testCompositeClosesEveryChildWhenOneThrows() throws Exception {
+ AtomicInteger closedChildren = new AtomicInteger();
+ Digester counting = new Digester() {
+ @Override
+ public void digest(TikaInputStream tis, Metadata m, ParseContext
ctx) {
+ }
+
+ @Override
+ public DigestSink digestSink(Metadata m, ParseContext ctx) {
+ return new DigestSink() {
+ @Override
+ public void write(int b) {
+ }
+
+ @Override
+ protected void finish(boolean publish) {
+ closedChildren.incrementAndGet();
+ }
+ };
+ }
+ };
+ Digester throwing = (tis, m, ctx) -> {
+ throw new IllegalStateException("boom");
+ };
+ // the thrower is a pull-only digester, so its BufferingDigestSink
throws from close()
+ Digester d = new CompositeDigester(counting, throwing, counting);
+ DigestSink sink = d.digestSink(new Metadata(), new ParseContext());
+ sink.write(1);
+ sink.commit(); // publishing is what runs the pull digester that
throws
+ assertThrows(IllegalStateException.class, sink::close);
+ assertEquals(2, closedChildren.get(), "children after the throwing one
still closed");
+ }
+
+ /**
+ * A child sink that cannot be created must not leave the already-created
children
+ * publishing a digest of the zero bytes they received.
+ */
+ @Test
+ public void testCompositeCleansUpWhenAChildSinkCannotBeCreated() throws
Exception {
+ AtomicInteger closedChildren = new AtomicInteger();
+ AtomicBoolean publishedAnything = new AtomicBoolean();
+ Digester ok = new Digester() {
+ @Override
+ public void digest(TikaInputStream tis, Metadata m, ParseContext
ctx) {
+ }
+
+ @Override
+ public DigestSink digestSink(Metadata m, ParseContext ctx) {
+ return new DigestSink() {
+ @Override
+ public void write(int b) {
+ }
+
+ @Override
+ protected void finish(boolean publish) {
+ closedChildren.incrementAndGet();
+ publishedAnything.compareAndSet(false, publish);
+ }
+ };
+ }
+ };
+ Digester real = new InputStreamDigester("MD5", MD5_KEY, HEX);
+ Digester failing = new Digester() {
+ @Override
+ public void digest(TikaInputStream tis, Metadata m, ParseContext
ctx) {
+ }
+
+ @Override
+ public DigestSink digestSink(Metadata m, ParseContext ctx) throws
IOException {
+ throw new IOException("cannot open");
+ }
+ };
+ Metadata m = new Metadata();
+ IOException e = assertThrows(IOException.class,
+ () -> new CompositeDigester(ok, real, failing).digestSink(m,
new ParseContext()));
+ assertEquals("cannot open", e.getMessage(), "the original failure
propagates");
+ assertEquals(1, closedChildren.get(), "already-created sinks are
closed");
+ assertEquals(0, e.getSuppressed().length, "clean closes add nothing");
+ assertFalse(publishedAnything.get(), "no child may publish on the
failure path");
+ assertNull(m.get(MD5_KEY), "no digest of the zero bytes the children
received");
+ assertNull(m.get(HttpHeaders.CONTENT_LENGTH), "and no Content-Length:
0");
+ }
+
+ /** An uncommitted sink that spilled must still take its temp file with
it. */
+ @Test
+ public void testUncommittedSpillIsDeleted() throws Exception {
+ InputStreamDigester inner = new InputStreamDigester("MD5", MD5_KEY,
HEX);
+ long before = tempFiles();
+ Metadata m = new Metadata();
+ try (DigestSink sink = pullOnly(inner).digestSink(m, new
ParseContext())) {
+ sink.write(data(BufferingDigestSink.MEMORY_THRESHOLD + 4096), 0,
+ BufferingDigestSink.MEMORY_THRESHOLD + 4096);
+ assertEquals(before + 1, tempFiles(), "content is on disk");
+ // no commit
+ }
+ assertEquals(before, tempFiles(), "an uncommitted sink still deletes
its spill file");
+ assertNull(m.get(MD5_KEY));
+ }
+
+ /** A child that refuses to commit must not strand its siblings' temp
files. */
+ @Test
+ public void testCompositeClosesChildrenWhenACommitThrows() throws
Exception {
+ InputStreamDigester inner = new InputStreamDigester("MD5", MD5_KEY,
HEX);
+ Digester selfClosing = new Digester() {
+ @Override
+ public void digest(TikaInputStream tis, Metadata m, ParseContext
ctx) {
+ }
+
+ @Override
+ public DigestSink digestSink(Metadata m, ParseContext ctx) {
+ return new DigestSink() {
+ @Override
+ public void write(int b) {
+ }
+
+ @Override
+ public void write(byte[] b, int off, int len) throws
IOException {
+ close(); // legal, and it makes the parent's
commit() throw
+ }
+
+ @Override
+ protected void finish(boolean publish) {
+ }
+ };
+ }
+ };
+ long before = tempFiles();
+ DigestSink sink = new CompositeDigester(pullOnly(inner), selfClosing)
+ .digestSink(new Metadata(), new ParseContext());
+ sink.write(data(BufferingDigestSink.MEMORY_THRESHOLD + 4096), 0,
+ BufferingDigestSink.MEMORY_THRESHOLD + 4096);
+ sink.commit();
+ assertThrows(IllegalStateException.class, sink::close);
+ assertEquals(before, tempFiles(), "the sibling's spill file was still
cleaned up");
+ }
+}
diff --git
a/tika-core/src/test/java/org/apache/tika/digest/FailingTestTranslator.java
b/tika-core/src/test/java/org/apache/tika/digest/FailingTestTranslator.java
new file mode 100644
index 0000000000..5d2b683629
--- /dev/null
+++ b/tika-core/src/test/java/org/apache/tika/digest/FailingTestTranslator.java
@@ -0,0 +1,63 @@
+/*
+ * 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.digest;
+
+import java.io.IOException;
+import java.io.OutputStream;
+
+import org.apache.tika.extractor.EmbeddedStreamTranslator;
+import org.apache.tika.io.TikaInputStream;
+import org.apache.tika.metadata.Metadata;
+import org.apache.tika.metadata.Property;
+
+/**
+ * Service-registered for tests only. Translates when the metadata carries
{@link #MODE}:
+ * "upper" upper-cases the bytes; "fail" writes half of them and then throws,
the way a
+ * translator meets a truncated container; "silent" claims the stream and
writes nothing,
+ * like PSTEmailStreamTranslator; "closes" writes everything and then closes
the stream it
+ * was handed, which it is not supposed to do.
+ */
+public class FailingTestTranslator implements EmbeddedStreamTranslator {
+
+ public static final Property MODE =
Property.internalText("test:translator-mode");
+
+ @Override
+ public boolean shouldTranslate(TikaInputStream inputStream, Metadata
metadata) {
+ return metadata.get(MODE) != null;
+ }
+
+ @Override
+ public void translate(TikaInputStream inputStream, Metadata metadata,
OutputStream os)
+ throws IOException {
+ String mode = metadata.get(MODE);
+ if ("silent".equals(mode)) {
+ return;
+ }
+ byte[] all = inputStream.readAllBytes();
+ boolean fail = "fail".equals(mode);
+ int n = fail ? all.length / 2 : all.length;
+ for (int i = 0; i < n; i++) {
+ os.write(Character.toUpperCase((char) all[i]));
+ }
+ if (fail) {
+ throw new IOException("translator gave up half way");
+ }
+ if ("closes".equals(mode)) {
+ os.close();
+ }
+ }
+}
diff --git
a/tika-core/src/test/java/org/apache/tika/io/InMemoryContentViewTest.java
b/tika-core/src/test/java/org/apache/tika/io/InMemoryContentViewTest.java
new file mode 100644
index 0000000000..5f711113ef
--- /dev/null
+++ b/tika-core/src/test/java/org/apache/tika/io/InMemoryContentViewTest.java
@@ -0,0 +1,266 @@
+/*
+ * 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.io;
+
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+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.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.io.ByteArrayInputStream;
+import java.nio.ByteBuffer;
+import java.nio.ReadOnlyBufferException;
+import java.nio.channels.ClosedChannelException;
+import java.nio.channels.SeekableByteChannel;
+import java.nio.file.Files;
+import java.nio.file.Path;
+
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+import org.apache.tika.metadata.HttpHeaders;
+import org.apache.tika.metadata.Metadata;
+
+public class InMemoryContentViewTest {
+
+ @TempDir
+ Path tempDir;
+
+ private static byte[] data(int len) {
+ byte[] d = new byte[len];
+ for (int i = 0; i < len; i++) {
+ d[i] = (byte) (i * 31 + 7);
+ }
+ return d;
+ }
+
+ private static byte[] contents(ByteBuffer view) {
+ byte[] out = new byte[view.remaining()];
+ view.duplicate().get(out);
+ return out;
+ }
+
+ @Test
+ public void testViewIsExactAndReadOnly() throws Exception {
+ byte[] data = data(10_000);
+ try (TemporaryResources tmp = new TemporaryResources()) {
+ TikaInputStream tis = TikaInputStream.get(new
ByteArrayInputStream(data), tmp, new Metadata());
+ tis.enableRewind(null);
+ try (SeekableByteChannel channel = tis.getSeekableByteChannel()) {
+ ByteBuffer view = TikaInputStream.inMemoryContent(channel);
+ assertNotNull(view, "small stream-backed content must be
served from memory");
+ // exact bounds: the cache array is over-allocated, the view
must not be
+ assertEquals(data.length, view.remaining());
+ assertEquals(data.length, view.capacity(), "clear() must not
expose the array's slack");
+ assertArrayEquals(data, contents(view));
+ assertTrue(view.isReadOnly());
+ assertThrows(ReadOnlyBufferException.class, () -> view.put(0,
(byte) 1));
+ }
+ }
+ }
+
+ @Test
+ public void testViewIsIndependentOfChannelPosition() throws Exception {
+ byte[] data = data(5_000);
+ try (TikaInputStream tis = TikaInputStream.get(data);
+ SeekableByteChannel channel = tis.getSeekableByteChannel()) {
+ channel.position(4_000);
+ ByteBuffer view = TikaInputStream.inMemoryContent(channel);
+ assertEquals(0, view.position());
+ assertEquals(data.length, view.limit());
+ view.get(new byte[100]);
+ assertEquals(4_000, channel.position());
+ }
+ }
+
+ @Test
+ public void testViewSurvivesSpillWhileChannelOpen() throws Exception {
+ byte[] data = data(20_000);
+ try (TemporaryResources tmp = new TemporaryResources()) {
+ tmp.setTemporaryFileDirectory(tempDir);
+ TikaInputStream tis = TikaInputStream.get(new
ByteArrayInputStream(data), tmp, new Metadata());
+ tis.enableRewind(null);
+ try (SeekableByteChannel channel = tis.getSeekableByteChannel()) {
+ ByteBuffer view = TikaInputStream.inMemoryContent(channel);
+ assertNotNull(view);
+ // spill + close the memory cache underneath the live view
+ Path spilled = tis.getPath();
+ assertTrue(Files.exists(spilled));
+ assertEquals(data.length, Files.size(spilled));
+ assertArrayEquals(data, contents(view), "view must still read
the content");
+ tis.rewind();
+ assertArrayEquals(data, tis.readAllBytes());
+ assertArrayEquals(data, contents(view), "view unaffected by
the stream moving");
+ }
+ }
+ }
+
+ @Test
+ public void testNullOnceContentIsOnDisk() throws Exception {
+ // no budget => the cache spills past its 1MB per-object threshold
+ byte[] data = data(3 * 1024 * 1024);
+ try (TemporaryResources tmp = new TemporaryResources()) {
+ tmp.setTemporaryFileDirectory(tempDir);
+ TikaInputStream tis = TikaInputStream.get(new
ByteArrayInputStream(data), tmp, new Metadata());
+ tis.enableRewind(null);
+ try (SeekableByteChannel channel = tis.getSeekableByteChannel()) {
+ assertNull(TikaInputStream.inMemoryContent(channel), "spilled
content has no view");
+ assertTrue(tis.hasFile());
+ }
+ }
+ }
+
+ @Test
+ public void testNullWhenBudgetIsExhausted() throws Exception {
+ CacheMemoryBudget budget = new CacheMemoryBudget(4 * 1024 * 1024);
+ assertEquals(budget.getMaxBytes(),
budget.tryReserve(budget.getMaxBytes()), "precondition");
+ byte[] data = data(2 * 1024 * 1024);
+ try (TemporaryResources tmp = new TemporaryResources()) {
+ tmp.setTemporaryFileDirectory(tempDir);
+ TikaInputStream tis = TikaInputStream.get(new
ByteArrayInputStream(data), tmp, new Metadata());
+ tis.enableRewind(budget);
+ try (SeekableByteChannel channel = tis.getSeekableByteChannel()) {
+ assertNull(TikaInputStream.inMemoryContent(channel), "over
threshold with no room => on disk");
+ assertTrue(tis.hasFile());
+ }
+ }
+ }
+
+ @Test
+ public void testViewOverReopenableSource() throws Exception {
+ byte[] data = data(10_000);
+ try (TemporaryResources tmp = new TemporaryResources();
+ TikaInputStream tis = TikaInputStream.get(() -> new
ByteArrayInputStream(data), tmp, null);
+ SeekableByteChannel channel = tis.getSeekableByteChannel()) {
+ ByteBuffer view = TikaInputStream.inMemoryContent(channel);
+ assertNotNull(view, "a re-openable source buffers small content in
memory");
+ assertEquals(data.length, view.capacity());
+ assertArrayEquals(data, contents(view));
+ }
+ }
+
+ @Test
+ public void testNullForFileBackedInput() throws Exception {
+ Path file = tempDir.resolve("in.bin");
+ Files.write(file, data(100));
+ try (TikaInputStream tis = TikaInputStream.get(file);
+ SeekableByteChannel channel = tis.getSeekableByteChannel()) {
+ assertNull(TikaInputStream.inMemoryContent(channel));
+ }
+ }
+
+ @Test
+ public void testNoViewFromClosedChannel() throws Exception {
+ try (TikaInputStream tis = TikaInputStream.get(data(100))) {
+ SeekableByteChannel channel = tis.getSeekableByteChannel();
+ channel.close();
+ assertThrows(ClosedChannelException.class, () ->
TikaInputStream.inMemoryContent(channel));
+ }
+ }
+
+ @Test
+ public void testPinHoldsBudgetUntilChannelCloses() throws Exception {
+ byte[] data = data(2 * 1024 * 1024);
+ CacheMemoryBudget budget = new CacheMemoryBudget(64L * 1024 * 1024);
+ try (TemporaryResources tmp = new TemporaryResources()) {
+ TikaInputStream tis = TikaInputStream.get(new
ByteArrayInputStream(data), tmp, new Metadata());
+ tis.enableRewind(budget);
+ SeekableByteChannel channel = tis.getSeekableByteChannel();
+ assertNotNull(TikaInputStream.inMemoryContent(channel), "2MB under
a 64MB budget stays in memory");
+ assertTrue(budget.getReservedBytes() > 0);
+ tis.close();
+ assertTrue(budget.getReservedBytes() > 0, "open channel still pins
the reservation");
+ channel.close();
+ assertEquals(0, budget.getReservedBytes(), "last close releases");
+ }
+ }
+
+ /**
+ * toString() must never force a spill. The at-risk state is a cache that
spilled
+ * mid-stream: hasFile() is true but no getPath() has run, so getPath()
would drain the
+ * rest of the source, flip the source to file mode and write
Content-Length.
+ */
+ @Test
+ public void testToStringDoesNotMaterializeContent() throws Exception {
+ byte[] data = data(5 * 1024 * 1024);
+ Metadata metadata = new Metadata();
+ try (TemporaryResources tmp = new TemporaryResources()) {
+ tmp.setTemporaryFileDirectory(tempDir);
+ TikaInputStream tis = TikaInputStream.get(new
ByteArrayInputStream(data), tmp, metadata);
+ tis.enableRewind(null);
+ // read past the 1MB cache threshold but nowhere near EOF
+ assertEquals(2 * 1024 * 1024, tis.readNBytes(2 * 1024 *
1024).length);
+ assertTrue(tis.hasFile(), "the cache spilled on its own");
+
+ String rendered = tis.toString();
+ assertNull(metadata.get(HttpHeaders.CONTENT_LENGTH),
+ "toString must not drain the source into the spill file");
+ assertFalse(rendered.contains(tempDir.toString()),
+ "the cache file is still growing, so it is not named");
+
+ // the rest of the stream is still there to read
+ assertEquals(3 * 1024 * 1024, tis.readAllBytes().length);
+ }
+ }
+
+ /**
+ * Once the cache holds the whole source, toString() names its file even
though no
+ * getPath() has run -- that is what lets an operator find the spool from
a log line.
+ */
+ @Test
+ public void testToStringNamesTheCacheFileOnceComplete() throws Exception {
+ byte[] data = data(3 * 1024 * 1024);
+ Metadata metadata = new Metadata();
+ try (TemporaryResources tmp = new TemporaryResources()) {
+ tmp.setTemporaryFileDirectory(tempDir);
+ TikaInputStream tis = TikaInputStream.get(new
ByteArrayInputStream(data), tmp, metadata);
+ tis.enableRewind(null);
+ assertEquals(data.length, tis.readAllBytes().length);
+ assertTrue(tis.hasFile(), "the cache spilled");
+ assertTrue(tis.toString().contains(tempDir.toString()),
+ "a complete cache file must be nameable without calling
getPath()");
+ assertNull(metadata.get(HttpHeaders.CONTENT_LENGTH),
+ "and naming it must still not drain or stamp metadata");
+ }
+ }
+
+ /** Once a path really exists, toString() names it so an operator can find
the file. */
+ @Test
+ public void testToStringNamesTheSpillFileOnceItExists() throws Exception {
+ byte[] data = data(3 * 1024 * 1024);
+ try (TemporaryResources tmp = new TemporaryResources()) {
+ tmp.setTemporaryFileDirectory(tempDir);
+ TikaInputStream tis = TikaInputStream.get(new
ByteArrayInputStream(data), tmp, new Metadata());
+ tis.enableRewind(null);
+ Path spilled = tis.getPath();
+ assertTrue(tis.toString().contains(spilled.toString()));
+ }
+ }
+
+ /** A file-backed stream reports its path without any spill machinery. */
+ @Test
+ public void testToStringNamesARealFile() throws Exception {
+ Path file = tempDir.resolve("named.bin");
+ Files.write(file, data(50));
+ try (TikaInputStream tis = TikaInputStream.get(file)) {
+ assertTrue(tis.toString().contains(file.toString()));
+ }
+ }
+}
diff --git
a/tika-core/src/test/java/org/apache/tika/io/StreamCacheBudgetTest.java
b/tika-core/src/test/java/org/apache/tika/io/StreamCacheBudgetTest.java
index 0b2b7a941b..8e8b7b5f51 100644
--- a/tika-core/src/test/java/org/apache/tika/io/StreamCacheBudgetTest.java
+++ b/tika-core/src/test/java/org/apache/tika/io/StreamCacheBudgetTest.java
@@ -23,11 +23,17 @@ import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
+import java.io.ByteArrayInputStream;
import java.io.IOException;
+import java.io.InputStream;
+import java.nio.channels.SeekableByteChannel;
+import java.nio.file.Files;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
+import org.apache.tika.metadata.Metadata;
+
public class StreamCacheBudgetTest {
private static final int THRESHOLD = 1024;
@@ -155,4 +161,57 @@ public class StreamCacheBudgetTest {
assertThrows(IOException.class, cache::toFile);
assertNull(cache.getInMemorySeekableByteChannel());
}
+
+ /** A threshold spill inside the cache puts the content on disk; hasFile()
must say so. */
+ @Test
+ public void testHasFileAfterInternalCacheSpill() throws Exception {
+ // no budget => StreamCache's 1MB per-object threshold applies
+ byte[] big = new byte[3 * 1024 * 1024];
+ try (TemporaryResources resources = new TemporaryResources()) {
+ TikaInputStream tis = TikaInputStream.get(new
ByteArrayInputStream(big), resources, new Metadata());
+ assertFalse(tis.hasFile(), "nothing read yet, nothing spilled");
+ tis.enableRewind(null);
+ try (SeekableByteChannel c = tis.getSeekableByteChannel()) {
+ assertEquals(big.length, c.size());
+ }
+ assertTrue(tis.hasFile(), "cache spilled to disk; hasFile() must
report it");
+ assertEquals(big.length, Files.size(tis.getPath()));
+ }
+ }
+
+ @Test
+ public void testHasFileStaysFalseWhenCacheKeepsContentInMemory() throws
Exception {
+ byte[] small = new byte[1024];
+ try (TemporaryResources resources = new TemporaryResources()) {
+ TikaInputStream tis = TikaInputStream.get(new
ByteArrayInputStream(small), resources, new Metadata());
+ tis.enableRewind(null);
+ try (SeekableByteChannel c = tis.getSeekableByteChannel()) {
+ assertEquals(small.length, c.size());
+ }
+ assertFalse(tis.hasFile(), "content fit in memory; no file
exists");
+ }
+ }
+
+ /**
+ * The cache registers with TemporaryResources only once it spills, so a
source whose
+ * close() throws must not skip the cache's close -- that leaks its
reservation from the
+ * fork-wide pool for good.
+ */
+ @Test
+ public void testThrowingSourceCloseStillReleasesBudget() throws Exception {
+ byte[] data = new byte[2 * 1024 * 1024];
+ InputStream source = new ByteArrayInputStream(data) {
+ @Override
+ public void close() {
+ throw new IllegalStateException("fetcher abort");
+ }
+ };
+ CacheMemoryBudget budget = new CacheMemoryBudget(64L * 1024 * 1024);
+ TikaInputStream tis = TikaInputStream.get(source, new
TemporaryResources(), new Metadata());
+ tis.enableRewind(budget);
+ assertEquals(data.length, tis.readAllBytes().length);
+ assertTrue(budget.getReservedBytes() > 0, "over-threshold content is
charged");
+ assertThrows(IllegalStateException.class, tis::close);
+ assertEquals(0, budget.getReservedBytes(), "cache closed despite the
throwing source");
+ }
}
diff --git
a/tika-core/src/test/java/org/apache/tika/io/TemporaryResourcesTest.java
b/tika-core/src/test/java/org/apache/tika/io/TemporaryResourcesTest.java
index fffb3f3778..09bfd084b2 100644
--- a/tika-core/src/test/java/org/apache/tika/io/TemporaryResourcesTest.java
+++ b/tika-core/src/test/java/org/apache/tika/io/TemporaryResourcesTest.java
@@ -16,11 +16,16 @@
*/
package org.apache.tika.io;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
+import java.io.Closeable;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
+import java.util.concurrent.atomic.AtomicInteger;
import org.junit.jupiter.api.Test;
@@ -37,4 +42,31 @@ public class TemporaryResourcesTest {
"Temp file should not exist after TempResources is closed");
}
+ /**
+ * A resource whose close() throws unchecked must not leave the rest open.
Resources close
+ * in reverse registration order, so the ones registered BEFORE the
thrower are at risk.
+ */
+ @Test
+ public void testUncheckedThrowDoesNotAbandonRemainingResources() throws
IOException {
+ AtomicInteger closed = new AtomicInteger();
+ Closeable counting = closed::incrementAndGet;
+ IllegalStateException boom = new IllegalStateException("boom");
+ TemporaryResources tmp = new TemporaryResources();
+ tmp.addResource(counting);
+ Path tempFile = tmp.createTempFile();
+ tmp.addResource(() -> {
+ throw new IOException("checked");
+ });
+ tmp.addResource(counting);
+ tmp.addResource(() -> {
+ throw boom;
+ });
+ tmp.addResource(counting);
+
+ IllegalStateException thrown =
assertThrows(IllegalStateException.class, tmp::close);
+ assertSame(boom, thrown, "the first failure in close order
propagates");
+ assertEquals(1, thrown.getSuppressed().length, "the later checked
failure is suppressed");
+ assertEquals(3, closed.get(), "every counting resource closed,
including those after the throw");
+ assertTrue(Files.notExists(tempFile), "the temp file registered first
was still deleted");
+ }
}
diff --git
a/tika-core/src/test/resources/META-INF/services/org.apache.tika.extractor.EmbeddedStreamTranslator
b/tika-core/src/test/resources/META-INF/services/org.apache.tika.extractor.EmbeddedStreamTranslator
new file mode 100644
index 0000000000..5cd467a909
--- /dev/null
+++
b/tika-core/src/test/resources/META-INF/services/org.apache.tika.extractor.EmbeddedStreamTranslator
@@ -0,0 +1,15 @@
+# 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.
+org.apache.tika.digest.FailingTestTranslator
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 ee14e103f5..3bcbe0cb0a 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
@@ -94,35 +94,49 @@ public class PipesServer implements AutoCloseable {
// Process-wide budget bounding total in-memory stream caching
(embedded-object digest rewind
// buffers, etc.) so small embedded objects stay in RAM instead of
spilling per-object at 1MB.
- // NOTE: read in THIS (forked server) JVM -- set it via the config's
forkedJvmArgs, not on the
- // parent JVM. Tunable via -Dtika.pipes.cacheMemoryBudgetBytes; <=0
disables (falls back to
- // the 1MB per-object default). Clamped to a quarter of the fork's max
heap.
+ // One pool per forked JVM, shared by every thread in it. NOTE: read in
THIS (forked server)
+ // JVM -- set -Dtika.pipes.cacheMemoryBudgetBytes via the config's
forkedJvmArgs, not on the
+ // parent JVM.
static final String CACHE_MEMORY_BUDGET_BYTES_PROP =
"tika.pipes.cacheMemoryBudgetBytes";
- static final CacheMemoryBudget CACHE_MEMORY_BUDGET =
initCacheMemoryBudget();
+ // only reached when the JVM reports no heap limit at all (never on
HotSpot)
+ private static final long FALLBACK_BUDGET_BYTES = 256L * 1024 * 1024;
- private static CacheMemoryBudget initCacheMemoryBudget() {
- long bytes = 256L * 1024 * 1024;
- String val = System.getProperty(CACHE_MEMORY_BUDGET_BYTES_PROP);
+ static final CacheMemoryBudget CACHE_MEMORY_BUDGET =
+
initCacheMemoryBudget(System.getProperty(CACHE_MEMORY_BUDGET_BYTES_PROP),
+ Runtime.getRuntime().maxMemory());
+
+ // package-private and pure so every branch is a unit test
+ static CacheMemoryBudget initCacheMemoryBudget(String val, long maxMemory)
{
+ long clamp = maxMemory == Long.MAX_VALUE ? FALLBACK_BUDGET_BYTES :
maxMemory / 4;
+ long bytes = clamp;
+ String source = "default, a quarter of max heap";
if (val != null) {
try {
bytes = Long.parseLong(val.trim());
+ source = "-D" + CACHE_MEMORY_BUDGET_BYTES_PROP;
} catch (NumberFormatException e) {
LOG.warn("Could not parse -D{}={} (plain bytes required, no
unit suffix); " +
"using the default {}",
CACHE_MEMORY_BUDGET_BYTES_PROP, val, bytes);
}
}
- if (bytes <= 0) {
- LOG.info("Cache memory budget disabled; per-object 1MB spill
threshold applies");
+ if (bytes < 0) {
+ LOG.warn("-D{}={} is negative; treating it as 0 -- cache memory
budget disabled, " +
+ "per-object 1MB spill threshold applies",
CACHE_MEMORY_BUDGET_BYTES_PROP, val);
+ return null;
+ }
+ if (bytes == 0) {
+ LOG.info("Cache memory budget disabled ({}); per-object 1MB spill
threshold applies",
+ source);
return null;
}
- long clamp = Runtime.getRuntime().maxMemory() / 4;
if (bytes > clamp) {
- LOG.warn("Cache memory budget {} exceeds a quarter of max heap;
clamping to {}",
- bytes, clamp);
+ LOG.warn("-D{}={} exceeds a quarter of max heap; clamping to {}
(raise -Xmx to " +
+ "raise the ceiling)", CACHE_MEMORY_BUDGET_BYTES_PROP,
bytes, clamp);
bytes = clamp;
+ source += ", clamped to a quarter of max heap";
}
- LOG.info("Cache memory budget: {} bytes", bytes);
+ LOG.info("Cache memory budget: {} bytes ({})", bytes, source);
return new CacheMemoryBudget(bytes);
}
diff --git
a/tika-pipes/tika-pipes-core/src/test/java/org/apache/tika/pipes/core/server/CacheMemoryBudgetSeedingTest.java
b/tika-pipes/tika-pipes-core/src/test/java/org/apache/tika/pipes/core/server/CacheMemoryBudgetSeedingTest.java
index f44cae7f5c..cfb415841a 100644
---
a/tika-pipes/tika-pipes-core/src/test/java/org/apache/tika/pipes/core/server/CacheMemoryBudgetSeedingTest.java
+++
b/tika-pipes/tika-pipes-core/src/test/java/org/apache/tika/pipes/core/server/CacheMemoryBudgetSeedingTest.java
@@ -18,6 +18,7 @@ package org.apache.tika.pipes.core.server;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertSame;
import org.junit.jupiter.api.Test;
@@ -29,11 +30,26 @@ public class CacheMemoryBudgetSeedingTest {
@Test
public void testDefaultBudgetClampedToHeap() {
- // No -Dtika.pipes.cacheMemoryBudgetBytes in the surefire JVM -> the
256MB default,
- // clamped to a quarter of max heap
+ // No -Dtika.pipes.cacheMemoryBudgetBytes in the surefire JVM -> a
quarter of max heap
assertNotNull(PipesServer.CACHE_MEMORY_BUDGET);
- long expected = Math.min(256L * 1024 * 1024,
Runtime.getRuntime().maxMemory() / 4);
- assertEquals(expected, PipesServer.CACHE_MEMORY_BUDGET.getMaxBytes());
+ assertEquals(Runtime.getRuntime().maxMemory() / 4,
PipesServer.CACHE_MEMORY_BUDGET.getMaxBytes());
+ }
+
+ @Test
+ public void testInitBranches() {
+ long heap = 4L * 1024 * 1024 * 1024;
+ assertEquals(heap / 4, PipesServer.initCacheMemoryBudget(null,
heap).getMaxBytes());
+ assertEquals(512L * 1024 * 1024,
+ PipesServer.initCacheMemoryBudget("536870912",
heap).getMaxBytes(), "property honoured");
+ assertEquals(heap / 4,
+ PipesServer.initCacheMemoryBudget("9999999999",
heap).getMaxBytes(), "clamped");
+ assertEquals(heap / 4,
+ PipesServer.initCacheMemoryBudget("not-a-number",
heap).getMaxBytes(), "malformed => default");
+ assertNull(PipesServer.initCacheMemoryBudget("0", heap), "0 disables");
+ assertNull(PipesServer.initCacheMemoryBudget("-1", heap), "negative
disables");
+ assertEquals(256L * 1024 * 1024,
+ PipesServer.initCacheMemoryBudget(null,
Long.MAX_VALUE).getMaxBytes(),
+ "no heap limit reported => bounded fallback, not an unbounded
budget");
}
@Test
diff --git
a/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/FontCacheWarmer.java
b/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/FontCacheWarmer.java
new file mode 100644
index 0000000000..2ed2ebdd09
--- /dev/null
+++
b/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/FontCacheWarmer.java
@@ -0,0 +1,55 @@
+/*
+ * 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 java.util.concurrent.atomic.AtomicBoolean;
+
+import org.apache.pdfbox.pdmodel.font.FontMappers;
+import org.junit.jupiter.api.extension.BeforeAllCallback;
+import org.junit.jupiter.api.extension.ExtensionContext;
+
+/**
+ * Builds PDFBox's on-disk font cache once, before any test in this module
runs.
+ * <p>
+ * The first PDF parse in a JVM scans the system's fonts and writes {@code
~/.pdfbox.cache},
+ * which takes seconds on a cold machine. In a forked pipes worker that
happens inside a
+ * parse, where it counts as "no progress" against {@code
progressTimeoutMillis} -- so
+ * whichever test happens to parse the first PDF fails on a slow host, for a
reason that has
+ * nothing to do with what it tests. A full reactor build warms the cache in
the parser
+ * modules' own tests; CI runs this module in a shard where those never run.
+ * <p>
+ * The driver and the forks resolve the same cache file (PDFBox tries {@code
pdfbox.fontcache},
+ * then {@code user.home}, then {@code java.io.tmpdir}), so warming it here
covers the forks.
+ * Auto-registered via {@code META-INF/services} plus
+ * {@code junit.jupiter.extensions.autodetection.enabled}.
+ */
+public class FontCacheWarmer implements BeforeAllCallback {
+
+ private static final AtomicBoolean WARMED = new AtomicBoolean();
+
+ @Override
+ public void beforeAll(ExtensionContext context) {
+ if (!WARMED.compareAndSet(false, true)) {
+ return;
+ }
+ try {
+ FontMappers.instance().getTrueTypeFont("Helvetica", null);
+ } catch (Exception | LinkageError e) {
+ // warming is an optimization: a failure here must not fail the
build
+ }
+ }
+}
diff --git
a/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/SharedServerModeTest.java
b/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/SharedServerModeTest.java
index 59b8b4b3a5..65e3806a93 100644
---
a/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/SharedServerModeTest.java
+++
b/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/SharedServerModeTest.java
@@ -71,6 +71,8 @@ public class SharedServerModeTest {
"<throw class=\"java.lang.OutOfMemoryError\">oom message</throw>" +
"</mock>";
+ // hangs 60s with no runtime TimeoutLimits override, so
tika-config-shared-server.json's
+ // progressTimeoutMillis is what detects it: keep that short or this stops
testing anything
private static final String MOCK_TIMEOUT = "<?xml version=\"1.0\"
encoding=\"UTF-8\" ?>" +
"<mock>" +
"<metadata action=\"add\" name=\"dc:creator\">Timeout
Author</metadata>" +
diff --git
a/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/async/AsyncChaosMonkeyTest.java
b/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/async/AsyncChaosMonkeyTest.java
index 3e73c8e9c9..6613e65fe4 100644
---
a/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/async/AsyncChaosMonkeyTest.java
+++
b/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/async/AsyncChaosMonkeyTest.java
@@ -54,6 +54,8 @@ public class AsyncChaosMonkeyTest {
"<write element=\"p\">main_content</write>" +
"</mock>";
+ // hangs 60s and the expected timeout count is asserted exactly, so the
default config's
+ // (tika-config-basic.json) short progressTimeoutMillis is what detects
it: keep it short
private final String TIMEOUT = "<?xml version=\"1.0\" encoding=\"UTF-8\"
?>" + "<mock>" +
"<metadata action=\"add\" name=\"dc:creator\">Nikolai
Lobachevsky</metadata>" +
"<write element=\"p\">main_content</write>" +
diff --git
a/tika-pipes/tika-pipes-integration-tests/src/test/resources/META-INF/services/org.junit.jupiter.api.extension.Extension
b/tika-pipes/tika-pipes-integration-tests/src/test/resources/META-INF/services/org.junit.jupiter.api.extension.Extension
new file mode 100644
index 0000000000..80f9a28048
--- /dev/null
+++
b/tika-pipes/tika-pipes-integration-tests/src/test/resources/META-INF/services/org.junit.jupiter.api.extension.Extension
@@ -0,0 +1,15 @@
+# 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.
+org.apache.tika.pipes.core.FontCacheWarmer
diff --git
a/tika-pipes/tika-pipes-integration-tests/src/test/resources/configs/tika-config-passback.json
b/tika-pipes/tika-pipes-integration-tests/src/test/resources/configs/tika-config-passback.json
index 3857d65e30..41847b4ee4 100644
---
a/tika-pipes/tika-pipes-integration-tests/src/test/resources/configs/tika-config-passback.json
+++
b/tika-pipes/tika-pipes-integration-tests/src/test/resources/configs/tika-config-passback.json
@@ -47,7 +47,7 @@
"parse-context": {
"mock-digester-factory": {},
"timeout-limits": {
- "progressTimeoutMillis": 5000
+ "progressTimeoutMillis": 60000
}
},
"plugin-roots": "PLUGINS_PATHS"
diff --git
a/tika-pipes/tika-pipes-integration-tests/src/test/resources/configs/tika-config-truncate.json
b/tika-pipes/tika-pipes-integration-tests/src/test/resources/configs/tika-config-truncate.json
deleted file mode 100644
index 9bc7efa5e8..0000000000
---
a/tika-pipes/tika-pipes-integration-tests/src/test/resources/configs/tika-config-truncate.json
+++ /dev/null
@@ -1,61 +0,0 @@
-{
- "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": {},
- "sax-output-config": {
- "writeFileNameToContent": false
- },
- "runpack-extractor-factory": {
- "maxEmbeddedBytesForExtraction": 10
- },
- "timeout-limits": {
- "progressTimeoutMillis": 5000
- }
- },
- "plugin-roots": "PLUGINS_PATHS"
-}
diff --git
a/tika-pipes/tika-pipes-integration-tests/src/test/resources/configs/tika-config-uppercasing.json
b/tika-pipes/tika-pipes-integration-tests/src/test/resources/configs/tika-config-uppercasing.json
index 7732becfc8..fec18edc38 100644
---
a/tika-pipes/tika-pipes-integration-tests/src/test/resources/configs/tika-config-uppercasing.json
+++
b/tika-pipes/tika-pipes-integration-tests/src/test/resources/configs/tika-config-uppercasing.json
@@ -44,7 +44,7 @@
"parse-context": {
"mock-digester-factory": {},
"timeout-limits": {
- "progressTimeoutMillis": 5000
+ "progressTimeoutMillis": 60000
}
},
"plugin-roots": "PLUGINS_PATHS"
diff --git
a/tika-pipes/tika-pipes-integration-tests/src/test/resources/configs/tika-config-write-limiter.json
b/tika-pipes/tika-pipes-integration-tests/src/test/resources/configs/tika-config-write-limiter.json
index e5b28c369e..31280b14b2 100644
---
a/tika-pipes/tika-pipes-integration-tests/src/test/resources/configs/tika-config-write-limiter.json
+++
b/tika-pipes/tika-pipes-integration-tests/src/test/resources/configs/tika-config-write-limiter.json
@@ -54,7 +54,7 @@
"maxValuesPerField": 5
},
"timeout-limits": {
- "progressTimeoutMillis": 5000
+ "progressTimeoutMillis": 60000
}
},
"plugin-roots": "PLUGINS_PATHS"
diff --git
a/tika-pipes/tika-pipes-integration-tests/src/test/resources/junit-platform.properties
b/tika-pipes/tika-pipes-integration-tests/src/test/resources/junit-platform.properties
new file mode 100644
index 0000000000..54fd4618a1
--- /dev/null
+++
b/tika-pipes/tika-pipes-integration-tests/src/test/resources/junit-platform.properties
@@ -0,0 +1,17 @@
+#
+# 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.
+#
+junit.jupiter.extensions.autodetection.enabled = true
diff --git
a/tika-server/tika-server-core/src/test/java/org/apache/tika/server/core/TikaServerIntegrationTest.java
b/tika-server/tika-server-core/src/test/java/org/apache/tika/server/core/TikaServerIntegrationTest.java
index d934988fb8..062248b78b 100644
---
a/tika-server/tika-server-core/src/test/java/org/apache/tika/server/core/TikaServerIntegrationTest.java
+++
b/tika-server/tika-server-core/src/test/java/org/apache/tika/server/core/TikaServerIntegrationTest.java
@@ -198,7 +198,9 @@ public class TikaServerIntegrationTest extends
IntegrationTestBase {
@Test
@Timeout(60000)
public void testTimeout() throws Exception {
- // With pipes-based parsing, timeout in a child process should NOT
crash the server
+ // With pipes-based parsing, timeout in a child process should NOT
crash the server.
+ // TEST_HEAVY_HANG relies on tika-config-server-pipes-basic.json's
short
+ // progressTimeoutMillis to be detected: keep that short.
startProcess(new String[]{"-config",
getConfig("tika-config-server-pipes-basic.json")});
awaitServerStartup();