This is an automated email from the ASF dual-hosted git repository. tballison pushed a commit to branch TIKA-4835-spill-primitives in repository https://gitbox.apache.org/repos/asf/tika.git
commit 5aebae53d8106720449f624ab952514f89b17ff4 Author: tallison <[email protected]> AuthorDate: Wed Aug 26 09:26:21 2026 -0400 TIKA-4835 -- tika-core primitives for spill-less parsing --- CHANGES.txt | 13 ++ .../apache/tika/digest/BufferingDigestSink.java | 84 ++++++++++ .../org/apache/tika/digest/CompositeDigester.java | 59 +++++++ .../java/org/apache/tika/digest/DigestHelper.java | 17 +-- .../main/java/org/apache/tika/digest/Digester.java | 19 +++ .../apache/tika/digest/InputStreamDigester.java | 32 ++++ .../java/org/apache/tika/io/CachingSource.java | 27 +--- .../apache/tika/io/MemorySeekableByteChannel.java | 6 + .../org/apache/tika/io/TemporaryResources.java | 46 ++++-- .../java/org/apache/tika/io/TikaInputStream.java | 19 +++ .../org/apache/tika/digest/DigestSinkTest.java | 135 ++++++++++++++++ .../apache/tika/io/InMemoryContentViewTest.java | 170 +++++++++++++++++++++ .../org/apache/tika/io/StreamCacheBudgetTest.java | 39 +++++ .../org/apache/tika/io/TemporaryResourcesTest.java | 54 +++++++ .../apache/tika/pipes/core/server/PipesServer.java | 11 +- .../core/server/CacheMemoryBudgetSeedingTest.java | 6 +- 16 files changed, 681 insertions(+), 56 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 9b7edca602..402927e239 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,5 +1,18 @@ 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) now defaults to a quarter + of the fork's heap, so raising -Xmx raises it; it was a fixed 256MB. It is + one pool per forked JVM shared by all of its threads. + -Dtika.pipes.cacheMemoryBudgetBytes in forkedJvmArgs still overrides it + (below the quarter-heap ceiling; <=0 disables). TikaInputStream.hasFile() + now also reports content the stream cache spilled on its own, not only + content a getPath() call put on disk. Digester gains digestSink(), an + OutputStream that digests as it is written; DigestHelper uses it for + translated embedded streams, which no longer touch a temp file. + TemporaryResources and CachingSource now close every remaining resource + when one close() throws an unchecked exception (TIKA-4835). + * Pipes plugins no longer bundle their own Jackson: jackson-core, -databind and -annotations are provided by the host (tika-serialization) and the plugins parent pom now bans bundling them, so a mapper can cross the 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..fbde5d39f4 --- /dev/null +++ b/tika-core/src/main/java/org/apache/tika/digest/BufferingDigestSink.java @@ -0,0 +1,84 @@ +/* + * 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.InputStream; +import java.io.OutputStream; + +import org.apache.commons.io.output.DeferredFileOutputStream; + +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 OutputStream { + + static final int MEMORY_THRESHOLD = 1024 * 1024; + + private final Digester digester; + private final Metadata metadata; + private final ParseContext context; + private final TemporaryResources tmp = new TemporaryResources(); + private final DeferredFileOutputStream buffer; + private boolean closed; + + BufferingDigestSink(Digester digester, Metadata metadata, ParseContext context) + throws IOException { + this.digester = digester; + this.metadata = metadata; + this.context = context; + this.buffer = DeferredFileOutputStream.builder() + .setThreshold(MEMORY_THRESHOLD) + .setOutputFile(tmp.createTempFile().toFile()) + .get(); + } + + @Override + public void write(int b) throws IOException { + buffer.write(b); + } + + @Override + public void write(byte[] b, int off, int len) throws IOException { + buffer.write(b, off, len); + } + + @Override + public void close() throws IOException { + if (closed) { + return; + } + closed = true; + try { + buffer.close(); + // toInputStream() serves memory without copying and the file otherwise + try (InputStream in = buffer.toInputStream(); + TikaInputStream tis = TikaInputStream.get(in, new TemporaryResources(), null)) { + digester.digest(tis, metadata, context); + } + } finally { + tmp.close(); + } + } +} 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..71636c8a44 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 @@ -17,6 +17,7 @@ package org.apache.tika.digest; import java.io.IOException; +import java.io.OutputStream; import org.apache.tika.io.TikaInputStream; import org.apache.tika.metadata.Metadata; @@ -37,4 +38,62 @@ public class CompositeDigester implements Digester { digester.digest(tis, m, parseContext); } } + + /** Fans each write out to every child's sink; close closes them all. */ + @Override + public OutputStream digestSink(Metadata m, ParseContext parseContext) throws IOException { + OutputStream[] sinks = new OutputStream[digesters.length]; + try { + for (int i = 0; i < digesters.length; i++) { + sinks[i] = digesters[i].digestSink(m, parseContext); + } + } catch (IOException | RuntimeException e) { + closeAll(sinks, e); + throw e; + } + return new OutputStream() { + @Override + public void write(int b) throws IOException { + for (OutputStream sink : sinks) { + sink.write(b); + } + } + + @Override + public void write(byte[] b, int off, int len) throws IOException { + for (OutputStream sink : sinks) { + sink.write(b, off, len); + } + } + + @Override + public void close() throws IOException { + closeAll(sinks, null); + } + }; + } + + // Closes every non-null sink even if one throws; the first failure is what propagates. + private static void closeAll(OutputStream[] sinks, Throwable pending) throws IOException { + IOException first = null; + for (OutputStream sink : sinks) { + if (sink == null) { + continue; + } + try { + sink.close(); + } catch (IOException e) { + if (pending != null) { + pending.addSuppressed(e); + } else if (first == null) { + first = e; + } else { + first.addSuppressed(e); + } + } + } + if (first != null) { + throw first; + } + } } 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..4505b4e9cd 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 @@ -18,13 +18,10 @@ 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.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; @@ -82,17 +79,13 @@ public class DigestHelper { Digester digester = digesterFactory.build(); // The translator consumes `tis` (e.g. OLE2), so enableRewind() before and rewind() - // after -- otherwise the caller would see an exhausted stream. + // after -- otherwise the caller would see an exhausted stream. The translated bytes + // go straight into the digest: the translator only writes, so nothing can mark, + // reset or skip between it and the digest. 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 (OutputStream sink = digester.digestSink(metadata, context)) { + EMBEDDED_STREAM_TRANSLATOR.translate(tis, metadata, sink); } finally { tis.rewind(); } 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..d39654c305 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 @@ -17,6 +17,7 @@ package org.apache.tika.digest; import java.io.IOException; +import java.io.OutputStream; import org.apache.tika.io.TikaInputStream; import org.apache.tika.metadata.Metadata; @@ -42,4 +43,22 @@ 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 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. + * <p> + * The default buffers what is written -- in memory below a fixed threshold, in a temp + * file above it -- and calls {@link #digest} on close, so an implementation 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 to set the values for on close + * @param parseContext ParseContext + * @return a sink; the caller must close it, and the values are set only on close + */ + default OutputStream 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..4ad7d9bbf7 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 @@ -17,6 +17,7 @@ package org.apache.tika.digest; import java.io.IOException; +import java.io.OutputStream; import java.security.MessageDigest; import java.security.NoSuchAlgorithmException; import java.security.Provider; @@ -137,4 +138,35 @@ public class InputStreamDigester implements Digester { tis.rewind(); } + /** Streams: every write goes straight into the MessageDigest; close sets the values. */ + @Override + public OutputStream digestSink(Metadata metadata, ParseContext parseContext) { + MessageDigest messageDigest = newMessageDigest(); + return new OutputStream() { + private long total; + private boolean closed; + + @Override + public void write(int b) { + messageDigest.update((byte) b); + total++; + } + + @Override + public void write(byte[] b, int off, int len) { + messageDigest.update(b, off, len); + total += len; + } + + @Override + public void close() { + if (closed) { + return; + } + closed = true; + setContentLength(total, metadata); + metadata.set(metadataProperty, encoder.encode(messageDigest.digest())); + } + }; + } } 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..3035a534ab 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; @@ -260,7 +259,10 @@ class CachingSource extends InputStream implements TikaInputSource { @Override public boolean hasPath() { - return spilledPath != null; + // A threshold spill inside the cache puts the content on disk without going through + // getPath(). Reporting only the getPath() case sends callers that branch on + // hasFile() down the read-it-into-memory path for content already on disk. + return spilledPath != null || (cachingStream != null && cachingStream.isFileBacked()); } @Override @@ -326,24 +328,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/MemorySeekableByteChannel.java b/tika-core/src/main/java/org/apache/tika/io/MemorySeekableByteChannel.java index 71d71d93ed..6187dda620 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,12 @@ 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(); + return ByteBuffer.wrap(data, 0, length).asReadOnlyBuffer(); + } + @Override public int read(ByteBuffer dst) throws IOException { ensureOpen(); 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..f2793f2d32 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,42 @@ 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: a tracked resource + * whose close() throws a RuntimeException must not leave the ones after it open. The + * first failure propagates with the rest attached as suppressed. + */ + 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; + closeable.close(); + } catch (Throwable t) { + if (first == null) { + first = t; } else { - exception.addSuppressed(e); + 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/TikaInputStream.java b/tika-core/src/main/java/org/apache/tika/io/TikaInputStream.java index 6de71cfa8e..a45d80a15c 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; @@ -552,6 +553,24 @@ 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 (no release, no reuse), and the content is fully + * drained before any channel is handed out (no growth). Keep the channel open for as + * long as the view is in use, then close it. + */ + 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 "; 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..f8615bde40 --- /dev/null +++ b/tika-core/src/test/java/org/apache/tika/digest/DigestSinkTest.java @@ -0,0 +1,135 @@ +/* + * 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.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; + +import java.io.IOException; +import java.io.OutputStream; +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 sink must produce exactly what the pull-style digest produces for the same bytes, + * through every Digester shape: streaming, composite, and the buffering default. + */ +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 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 (OutputStream 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; + } + } + return m; + } + + @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)); + } + + /** A third-party Digester that only implements digest() gets the buffering default. */ + @Test + public void testDefaultBuffersForPullOnlyDigester() throws Exception { + InputStreamDigester inner = new InputStreamDigester("SHA-256", SHA_KEY, HEX); + Digester pullOnly = (tis, m, ctx) -> inner.digest(tis, m, ctx); + // past the buffering threshold so the temp-file branch runs too + byte[] data = data(BufferingDigestSink.MEMORY_THRESHOLD + 12_345); + Metadata expected = viaStream(inner, data); + Metadata actual = viaSink(pullOnly, data); + assertEquals(expected.get(SHA_KEY), actual.get(SHA_KEY)); + assertEquals(Integer.toString(data.length), actual.get(HttpHeaders.CONTENT_LENGTH)); + } + + @Test + public void testDefaultBuffersSmallInMemory() throws Exception { + InputStreamDigester inner = new InputStreamDigester("MD5", MD5_KEY, HEX); + Digester pullOnly = (tis, m, ctx) -> inner.digest(tis, m, ctx); + byte[] data = data(1234); + assertEquals(viaStream(inner, data).get(MD5_KEY), viaSink(pullOnly, data).get(MD5_KEY)); + } + + @Test + public void testValuesSetOnlyOnClose() throws Exception { + Digester d = new InputStreamDigester("MD5", MD5_KEY, HEX); + Metadata m = new Metadata(); + OutputStream sink = d.digestSink(m, new ParseContext()); + sink.write(data(100), 0, 100); + assertNull(m.get(MD5_KEY), "digest must not be visible before close"); + sink.close(); + assertNotNull(m.get(MD5_KEY)); + sink.close(); // idempotent + } +} 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..4279e2b51b --- /dev/null +++ b/tika-core/src/test/java/org/apache/tika/io/InMemoryContentViewTest.java @@ -0,0 +1,170 @@ +/* + * 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.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.Metadata; + +/** + * {@link TikaInputStream#inMemoryContent} hands a library the cache's own array. These pin + * the contract: the view is exact, read-only, and stays valid for the life of the channel + * even when the stream is spilled or advanced underneath it. + */ +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()); + 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()); + // reading through the view does not move the channel either + view.get(new byte[100]); + assertEquals(4_000, channel.position()); + } + } + + /** The whole point: a library may hold the view while something else spills the stream. */ + @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"); + // and the stream itself is still usable from the file + tis.rewind(); + assertArrayEquals(data, tis.readAllBytes()); + } + } + } + + @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 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)); + } + } + + /** Budget accounting: the view's pin holds the reservation, closing the channel releases it. */ + @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"); + } + } +} 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..b47f0c4055 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,14 @@ 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 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 +158,40 @@ 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, or + * callers that branch on it read already-spilled content back into memory. + */ + @Test + public void testHasFileAfterInternalCacheSpill() throws Exception { + // no budget => StreamCache's 1MB per-object threshold applies + byte[] big = new byte[3 * 1024 * 1024]; + try (TemporaryResources tmp = new TemporaryResources()) { + TikaInputStream tis = + TikaInputStream.get(new ByteArrayInputStream(big), tmp, new Metadata()); + assertFalse(tis.hasFile(), "nothing read yet, nothing spilled"); + tis.enableRewind(null); + // drain through the cache; past the threshold it spills on its own + try (java.nio.channels.SeekableByteChannel c = tis.getSeekableByteChannel()) { + assertEquals(big.length, c.size()); + } + assertTrue(tis.hasFile(), "cache spilled to disk; hasFile() must report it"); + assertEquals(big.length, java.nio.file.Files.size(tis.getPath())); + } + } + + @Test + public void testHasFileStaysFalseWhenCacheKeepsContentInMemory() throws Exception { + byte[] small = new byte[1024]; + try (TemporaryResources tmp = new TemporaryResources()) { + TikaInputStream tis = + TikaInputStream.get(new ByteArrayInputStream(small), tmp, new Metadata()); + tis.enableRewind(null); + try (java.nio.channels.SeekableByteChannel c = tis.getSeekableByteChannel()) { + assertEquals(small.length, c.size()); + } + assertFalse(tis.hasFile(), "content fit in memory; no file exists"); + } + } } 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..a3219901f7 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,18 @@ */ 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.ByteArrayInputStream; +import java.io.Closeable; import java.io.IOException; +import java.io.InputStream; import java.nio.file.Files; import java.nio.file.Path; +import java.util.concurrent.atomic.AtomicInteger; import org.junit.jupiter.api.Test; @@ -37,4 +44,51 @@ 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"); + } + + @Test + public void testCachingSourceCloseSurvivesUncheckedThrow() throws IOException { + AtomicInteger sourceClosed = new AtomicInteger(); + InputStream source = new ByteArrayInputStream(new byte[100]) { + @Override + public void close() { + sourceClosed.incrementAndGet(); + throw new IllegalStateException("source close"); + } + }; + try (TemporaryResources tmp = new TemporaryResources()) { + TikaInputStream tis = TikaInputStream.get(source, tmp, null); + tis.enableRewind(null); + // spill so a file stream exists ahead of the (throwing) source in the close order + tis.getPath(); + assertThrows(IllegalStateException.class, tis::close); + assertEquals(1, sourceClosed.get()); + } + } } 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..e7a90267e6 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,15 +94,17 @@ 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. Default and ceiling: a quarter of + // the fork's max heap, so raising -Xmx raises it. NOTE: read in THIS (forked server) JVM -- + // set -Dtika.pipes.cacheMemoryBudgetBytes via the config's forkedJvmArgs, not on the parent + // JVM; <=0 disables (falls back to the 1MB per-object default). static final String CACHE_MEMORY_BUDGET_BYTES_PROP = "tika.pipes.cacheMemoryBudgetBytes"; static final CacheMemoryBudget CACHE_MEMORY_BUDGET = initCacheMemoryBudget(); private static CacheMemoryBudget initCacheMemoryBudget() { - long bytes = 256L * 1024 * 1024; + long clamp = Runtime.getRuntime().maxMemory() / 4; + long bytes = clamp; String val = System.getProperty(CACHE_MEMORY_BUDGET_BYTES_PROP); if (val != null) { try { @@ -116,7 +118,6 @@ public class PipesServer implements AutoCloseable { LOG.info("Cache memory budget disabled; per-object 1MB spill threshold applies"); 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); 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..d1b8298a27 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 @@ -29,11 +29,9 @@ 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
