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

Reply via email to