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 9bf38562e4 TIKA-4952: send forked pipes servers their config over
stdin (#3298)
9bf38562e4 is described below
commit 9bf38562e45a2901f52a3172fe6175a65249ce03
Author: Tim Allison <[email protected]>
AuthorDate: Tue Oct 6 12:27:47 2026 -0400
TIKA-4952: send forked pipes servers their config over stdin (#3298)
---
CHANGES.txt | 9 +-
docs/modules/ROOT/pages/configuration/index.adoc | 5 +-
.../apache/tika/pipes/grpc/TikaGrpcServerImpl.java | 2 +-
.../tika/pipes/core/PerClientServerManager.java | 69 +++++++++---
.../org/apache/tika/pipes/core/PipesClient.java | 13 ++-
.../org/apache/tika/pipes/core/PipesParser.java | 31 ++++--
.../tika/pipes/core/SharedServerManager.java | 40 ++++---
.../tika/pipes/core/async/AsyncProcessor.java | 17 +--
.../tika/pipes/core/config/ConfigMerger.java | 56 +++++++---
.../tika/pipes/core/protocol/ForkBootstrap.java | 116 ++++++++++++++++++++
.../apache/tika/pipes/core/server/PipesServer.java | 65 +++++-------
.../core/PerClientServerManagerSizingTest.java | 9 +-
.../pipes/core/protocol/ForkBootstrapTest.java | 118 +++++++++++++++++++++
.../apache/tika/pipes/fork/PipesForkParser.java | 26 ++---
.../apache/tika/pipes/core/PipesForkTokenTest.java | 92 ++++++++++++++++
.../tika/pipes/core/PipesParserWarmClientTest.java | 2 +-
.../pipes/core/SharedServerChaosMonkeyTest.java | 4 +-
.../tika/pipes/core/SharedServerModeTest.java | 30 +++---
.../apache/tika/config/loader/TikaJsonConfig.java | 21 +++-
.../org/apache/tika/config/loader/TikaLoader.java | 15 ++-
.../apache/tika/server/core/TikaServerProcess.java | 19 ++--
.../org/apache/tika/server/core/CXFTestBase.java | 2 +-
.../org/apache/tika/server/core/TikaPipesTest.java | 2 +-
.../apache/tika/server/standard/TikaPipesTest.java | 2 +-
24 files changed, 601 insertions(+), 164 deletions(-)
diff --git a/CHANGES.txt b/CHANGES.txt
index 13cf6d8b60..8cf92ee9a8 100644
--- a/CHANGES.txt
+++ b/CHANGES.txt
@@ -1,9 +1,10 @@
Release 4.2.0 - unreleased
- * GeoGebraParser reads structure.json and inline text content with
jackson-core's
- streaming parser; tika-parser-miscoffice-module and
tika-parser-cad-module no
- longer depend on jackson-databind, and OSGi deployments of
tika-bundle-standard
- need only jackson-core (TIKA-4956).
+ * Forked pipes servers get their config and auth token over stdin; nothing
is written
+ to the temp dir. Server-manager and PipesServer signatures change
accordingly (TIKA-4952).
+
+ * GeoGebraParser switched to stream parsing to drop jackson-databind as
+ a dependency (TIKA-4956).
* tika-bundle-standard no longer embeds dependencies that are OSGi bundles
themselves (commons-*, pdfbox, fontbox, bouncycastle, jsoup, asm, xz,
xmpcore,
diff --git a/docs/modules/ROOT/pages/configuration/index.adoc
b/docs/modules/ROOT/pages/configuration/index.adoc
index 2867cbf523..0421974bc5 100644
--- a/docs/modules/ROOT/pages/configuration/index.adoc
+++ b/docs/modules/ROOT/pages/configuration/index.adoc
@@ -281,8 +281,9 @@ set fails the load with the variable's name and the JSON
path, never a value, so
stops the server at startup rather than reaching a service as an empty key. A
literal value works as
before.
-The forked parse workers of tika-server and Pipes inherit the parent's
environment and read the same
-config file, so a reference resolves there too. Presets written inline in the
config file resolve
+The forked parse workers of tika-server and Pipes never read the config file
or the environment:
+the parent resolves every reference once and hands each fork the resolved
config over its stdin,
+so a value that itself contains `${env:...}` is not expanded a second time.
Presets written inline in the config file resolve
at startup with the rest of it; catalog presets and per-request configuration
(`/config`
endpoints, pipes tuples) never do: a request cannot read the server's
environment.
diff --git
a/tika-grpc/src/main/java/org/apache/tika/pipes/grpc/TikaGrpcServerImpl.java
b/tika-grpc/src/main/java/org/apache/tika/pipes/grpc/TikaGrpcServerImpl.java
index d8682b1c61..7323a01c54 100644
--- a/tika-grpc/src/main/java/org/apache/tika/pipes/grpc/TikaGrpcServerImpl.java
+++ b/tika-grpc/src/main/java/org/apache/tika/pipes/grpc/TikaGrpcServerImpl.java
@@ -126,7 +126,7 @@ class TikaGrpcServerImpl extends TikaGrpc.TikaImplBase {
tikaGrpcConfig = TikaGrpcConfig.load(tikaJsonConfig);
// PipesClient is single-threaded; the pool admits pipes.numClients at
a time.
- pipesParser = PipesParser.load(tikaJsonConfig, pipesConfig,
configPath);
+ pipesParser = PipesParser.load(tikaJsonConfig, pipesConfig);
try {
if (pluginRootsOverride != null &&
!pluginRootsOverride.trim().isEmpty()) {
diff --git
a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/PerClientServerManager.java
b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/PerClientServerManager.java
index 24f53d2822..98a9b870f0 100644
---
a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/PerClientServerManager.java
+++
b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/PerClientServerManager.java
@@ -17,6 +17,7 @@
package org.apache.tika.pipes.core;
import java.io.IOException;
+import java.io.OutputStream;
import java.net.InetAddress;
import java.net.InetSocketAddress;
import java.net.ServerSocket;
@@ -35,6 +36,7 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.tika.config.TikaExtras;
+import org.apache.tika.pipes.core.protocol.ForkBootstrap;
import org.apache.tika.pipes.core.server.PipesServer;
import org.apache.tika.utils.ProcessUtils;
@@ -46,13 +48,16 @@ import org.apache.tika.utils.ProcessUtils;
* <p>
* Connection model: The client creates a ServerSocket and the server connects
TO it.
* This is the reverse of typical client-server patterns but allows the client
to
- * control the port assignment.
+ * control the port assignment. The config and a per-process auth token go to
the
+ * server on its stdin ({@link ForkBootstrap}); the server presents the token
as the
+ * first bytes of its connection.
*/
public class PerClientServerManager implements ServerManager {
private static final Logger LOG =
LoggerFactory.getLogger(PerClientServerManager.class);
private static final long WAIT_ON_DESTROY_MS = 10000;
public static final int SOCKET_CONNECT_TIMEOUT_MS = 60000;
+ private static final int TOKEN_READ_TIMEOUT_MS = 10000;
/** Cores reserved for the parent JVM when auto-sizing forked JVMs'
* -XX:ActiveProcessorCount. The parent has client-side serialization,
* response deserialization, and heartbeat bookkeeping; if it's
CPU-starved
@@ -184,12 +189,13 @@ public class PerClientServerManager implements
ServerManager {
private final PipesConfig pipesConfig;
- private final Path tikaConfigPath;
+ private final byte[] tikaConfigJson;
private final int clientId;
private volatile Process process;
private volatile ServerSocket serverSocket;
private volatile Path tmpDir;
+ private volatile byte[] token;
private volatile int port = -1;
private long filesProcessed = 0;
private volatile long generation;
@@ -199,9 +205,13 @@ public class PerClientServerManager implements
ServerManager {
// process after the manager has been torn down (which would leak the
child).
private volatile boolean closed = false;
- public PerClientServerManager(PipesConfig pipesConfig, Path
tikaConfigPath, int clientId) {
+ /**
+ * @param tikaConfigJson the parent's resolved config, from {@link
ForkBootstrap#toBytes};
+ * handed to every fork this manager starts
+ */
+ public PerClientServerManager(PipesConfig pipesConfig, byte[]
tikaConfigJson, int clientId) {
this.pipesConfig = pipesConfig;
- this.tikaConfigPath = tikaConfigPath;
+ this.tikaConfigJson = tikaConfigJson;
this.clientId = clientId;
// Emit CPU-sizing diagnostics once per PipesParser (only on the first
client).
if (clientId == 0) {
@@ -398,6 +408,7 @@ public class PerClientServerManager implements
ServerManager {
// Capture the socket up front: shutdown() may null the field
concurrently, but this
// request keeps using (and detects the close on) the instance it
started with.
ServerSocket ss = serverSocket;
+ byte[] expectedToken = token;
if (ss == null) {
throw new IllegalStateException("Server not started. Call
ensureRunning() first.");
}
@@ -409,6 +420,18 @@ public class PerClientServerManager implements
ServerManager {
while (true) {
try {
Socket socket = ss.accept();
+ if (!presentsToken(socket, expectedToken)) {
+ LOG.warn("clientId={}: rejected a connection that did not
present the fork's token",
+ clientId);
+ try {
+ socket.close();
+ } catch (IOException closeEx) {
+ LOG.debug("clientId={}: error closing rejected
connection", clientId, closeEx);
+ }
+ // Strangers must not be able to hold the deadline open.
+ checkConnectDeadline(startTime);
+ continue;
+ }
socket.setSoTimeout(socketTimeoutMillis);
socket.setTcpNoDelay(true);
LOG.debug("clientId={}: accepted connection from server",
clientId);
@@ -438,18 +461,31 @@ public class PerClientServerManager implements
ServerManager {
throw new IOException(
"Server process died before connecting (exit code
" + exitValue + ") - will retry");
}
- // Check if we've exceeded the overall timeout
- long elapsed = System.currentTimeMillis() - startTime;
- if (elapsed > SOCKET_CONNECT_TIMEOUT_MS) {
- LOG.error("clientId={}: Timed out waiting for server to
connect after {}ms", clientId, elapsed);
- throw new ServerInitializationException(
- "Server did not connect within " +
SOCKET_CONNECT_TIMEOUT_MS + "ms");
- }
+ checkConnectDeadline(startTime);
// Continue polling
}
}
}
+ private void checkConnectDeadline(long startTime) throws
ServerInitializationException {
+ long elapsed = System.currentTimeMillis() - startTime;
+ if (elapsed > SOCKET_CONNECT_TIMEOUT_MS) {
+ LOG.error("clientId={}: Timed out waiting for server to connect
after {}ms", clientId, elapsed);
+ throw new ServerInitializationException(
+ "Server did not connect within " +
SOCKET_CONNECT_TIMEOUT_MS + "ms");
+ }
+ }
+
+ private static boolean presentsToken(Socket socket, byte[] expectedToken) {
+ try {
+ // The fork writes its token immediately; a silent stranger only
costs this wait.
+ socket.setSoTimeout(TOKEN_READ_TIMEOUT_MS);
+ return ForkBootstrap.checkToken(socket.getInputStream(),
expectedToken);
+ } catch (IOException e) {
+ return false;
+ }
+ }
+
@SuppressWarnings("deprecation")
private synchronized void startServer() throws IOException,
InterruptedException, TimeoutException, ServerInitializationException {
if (closed) {
@@ -498,6 +534,8 @@ public class PerClientServerManager implements
ServerManager {
pb.redirectOutput(ProcessBuilder.Redirect.DISCARD);
pb.redirectError(ProcessBuilder.Redirect.DISCARD);
}
+ // stdin carries the ForkBootstrap
+ pb.redirectInput(ProcessBuilder.Redirect.PIPE);
try {
process = pb.start();
@@ -513,6 +551,14 @@ public class PerClientServerManager implements
ServerManager {
throw new ServerInitializationException(msg, e);
}
+ token = ForkBootstrap.newToken();
+ try (OutputStream stdin = process.getOutputStream()) {
+ ForkBootstrap.write(stdin, token, tikaConfigJson);
+ } catch (IOException e) {
+ // The fork died before reading it; connect() sees the exit and
retries.
+ LOG.warn("clientId={}: couldn't send the bootstrap to the server
process", clientId, e);
+ }
+
// Server is started, but we don't wait for connection here.
// The connection is established in connect() method.
LOG.debug("clientId={}: server process started, waiting for connection
in connect()", clientId);
@@ -733,7 +779,6 @@ public class PerClientServerManager implements
ServerManager {
commandLine.add("org.apache.tika.pipes.core.server.PipesServer");
commandLine.add(Integer.toString(port));
- commandLine.add(tikaConfigPath.toAbsolutePath().toString());
LOG.debug("clientId={}: commandline: {}", clientId, commandLine);
return commandLine.toArray(new String[0]);
diff --git
a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/PipesClient.java
b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/PipesClient.java
index aeff9c8d64..54f7e993a5 100644
---
a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/PipesClient.java
+++
b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/PipesClient.java
@@ -39,12 +39,15 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.tika.config.TimeoutLimits;
+import org.apache.tika.config.loader.TikaJsonConfig;
+import org.apache.tika.exception.TikaConfigException;
import org.apache.tika.metadata.Metadata;
import org.apache.tika.metadata.TikaCoreProperties;
import org.apache.tika.pipes.api.FetchEmitTuple;
import org.apache.tika.pipes.api.PipesResult;
import org.apache.tika.pipes.api.emitter.EmitKey;
import org.apache.tika.pipes.core.emitter.EmitDataImpl;
+import org.apache.tika.pipes.core.protocol.ForkBootstrap;
import org.apache.tika.pipes.core.protocol.PayloadLimitExceededException;
import org.apache.tika.pipes.core.protocol.PipesMessage;
import org.apache.tika.pipes.core.protocol.PipesMessageType;
@@ -132,13 +135,17 @@ public class PipesClient implements Closeable {
* lazily on first use and shut down when this client is closed.
*
* @param pipesConfig the pipes configuration
- * @param tikaConfigPath path to the tika config file
+ * @param tikaConfigPath path to the tika config file; read once, here
+ * @throws IOException if the config can't be read
+ * @throws TikaConfigException if the config is invalid
*/
- public PipesClient(PipesConfig pipesConfig, java.nio.file.Path
tikaConfigPath) {
+ public PipesClient(PipesConfig pipesConfig, java.nio.file.Path
tikaConfigPath)
+ throws IOException, TikaConfigException {
this.pipesConfig = pipesConfig;
this.maxIpcPayloadBytes = pipesConfig.getMaxIpcPayloadBytes();
this.pipesClientId = CLIENT_COUNTER.getAndIncrement();
- this.serverManager = new PerClientServerManager(pipesConfig,
tikaConfigPath, pipesClientId);
+ byte[] tikaConfigJson =
ForkBootstrap.toBytes(TikaJsonConfig.load(tikaConfigPath));
+ this.serverManager = new PerClientServerManager(pipesConfig,
tikaConfigJson, pipesClientId);
this.ownsServerManager = true;
}
diff --git
a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/PipesParser.java
b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/PipesParser.java
index b123517ac5..3401916a50 100644
---
a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/PipesParser.java
+++
b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/PipesParser.java
@@ -35,6 +35,7 @@ import org.apache.tika.config.loader.TikaJsonConfig;
import org.apache.tika.exception.TikaConfigException;
import org.apache.tika.pipes.api.FetchEmitTuple;
import org.apache.tika.pipes.api.PipesResult;
+import org.apache.tika.pipes.core.protocol.ForkBootstrap;
import org.apache.tika.plugins.TikaPluginManager;
public class PipesParser implements Closeable {
@@ -59,7 +60,7 @@ public class PipesParser implements Closeable {
public static PipesParser load(Path tikaConfigPath) throws IOException,
TikaConfigException {
TikaJsonConfig tikaJsonConfig = TikaJsonConfig.load(tikaConfigPath);
PipesConfig pipesConfig = PipesConfig.load(tikaJsonConfig);
- return load(tikaJsonConfig, pipesConfig, tikaConfigPath);
+ return load(tikaJsonConfig, pipesConfig);
}
/**
@@ -68,20 +69,29 @@ public class PipesParser implements Closeable {
* Use this method when you need to modify the PipesConfig before creating
* the parser (e.g., to override emit strategy).
*
- * @param tikaJsonConfig the pre-loaded JSON configuration
+ * @param tikaJsonConfig the pre-loaded JSON configuration; forked servers
receive this
+ * config, env references already resolved, over
their stdin
* @param pipesConfig the pipes configuration (may be modified by caller)
- * @param tikaConfigPath path to the config file (passed to child
processes)
* @return a new PipesParser instance
- * @throws IOException if plugin extraction fails
+ * @throws IOException if plugin extraction or config serialization fails
*/
+ public static PipesParser load(TikaJsonConfig tikaJsonConfig, PipesConfig
pipesConfig)
+ throws IOException {
+ TikaPluginManager.preExtractPlugins(tikaJsonConfig);
+ return new PipesParser(pipesConfig,
ForkBootstrap.toBytes(tikaJsonConfig));
+ }
+
+ /**
+ * @deprecated forked servers no longer read the config file; use
+ * {@link #load(TikaJsonConfig, PipesConfig)}. {@code tikaConfigPath} is
ignored.
+ */
+ @Deprecated
public static PipesParser load(TikaJsonConfig tikaJsonConfig, PipesConfig
pipesConfig,
Path tikaConfigPath) throws IOException {
- TikaPluginManager.preExtractPlugins(tikaJsonConfig);
- return new PipesParser(pipesConfig, tikaConfigPath);
+ return load(tikaJsonConfig, pipesConfig);
}
private final PipesConfig pipesConfig;
- private final Path tikaConfigPath;
private final List<PipesClient> clients = new ArrayList<>();
private final List<ServerManager> serverManagers = new ArrayList<>();
// LIFO: the most-recently-returned client is borrowed next, so light
traffic
@@ -90,9 +100,8 @@ public class PipesParser implements Closeable {
private final LinkedBlockingDeque<PipesClient> clientQueue;
private final boolean isSharedMode;
- private PipesParser(PipesConfig pipesConfig, Path tikaConfigPath) {
+ private PipesParser(PipesConfig pipesConfig, byte[] tikaConfigJson) {
this.pipesConfig = pipesConfig;
- this.tikaConfigPath = tikaConfigPath;
this.isSharedMode = pipesConfig.isUseSharedServer();
this.clientQueue = new
LinkedBlockingDeque<>(pipesConfig.getNumClients());
@@ -100,7 +109,7 @@ public class PipesParser implements Closeable {
// Shared mode: one ServerManager for all clients
LOG.info("Using shared server mode with {} clients",
pipesConfig.getNumClients());
SharedServerManager sharedManager = new SharedServerManager(
- pipesConfig, tikaConfigPath, pipesConfig.getNumClients());
+ pipesConfig, tikaConfigJson, pipesConfig.getNumClients());
serverManagers.add(sharedManager);
for (int i = 0; i < pipesConfig.getNumClients(); i++) {
@@ -113,7 +122,7 @@ public class PipesParser implements Closeable {
LOG.info("Using per-client server mode with {} clients",
pipesConfig.getNumClients());
for (int i = 0; i < pipesConfig.getNumClients(); i++) {
PerClientServerManager serverManager = new
PerClientServerManager(
- pipesConfig, tikaConfigPath, i);
+ pipesConfig, tikaConfigJson, i);
serverManagers.add(serverManager);
PipesClient client = new PipesClient(pipesConfig,
serverManager);
diff --git
a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/SharedServerManager.java
b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/SharedServerManager.java
index c0a8088d2c..8374d837b4 100644
---
a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/SharedServerManager.java
+++
b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/SharedServerManager.java
@@ -19,15 +19,14 @@ package org.apache.tika.pipes.core;
import java.io.BufferedReader;
import java.io.IOException;
import java.io.InputStreamReader;
+import java.io.OutputStream;
import java.net.InetAddress;
import java.net.InetSocketAddress;
import java.net.Socket;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
-import java.security.SecureRandom;
import java.util.ArrayList;
-import java.util.HexFormat;
import java.util.List;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
@@ -38,6 +37,7 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.tika.config.TikaExtras;
+import org.apache.tika.pipes.core.protocol.ForkBootstrap;
import org.apache.tika.pipes.core.server.PipesServer;
import org.apache.tika.utils.ProcessUtils;
@@ -53,12 +53,10 @@ import org.apache.tika.utils.ProcessUtils;
* {@link #ensureRunning()} is synchronized to prevent multiple clients from
attempting
* to restart the server simultaneously.
* <p>
- * <b>Security:</b> The server port and a 32-byte auth token are passed to the
child
- * process via environment variables (not command-line args), so they are not
visible
- * in {@code /proc/<pid>/cmdline}. Each client connection must present the
token before
- * the server will accept it. This prevents CVE-style abuse from untrusted
local
- * processes. Note: if a malicious actor has same-uid access to your host and
can read
- * {@code /proc/<pid>/environ}, that is beyond Tika's security model.
+ * <b>Security:</b> A 32-byte auth token goes to the child with its config on
stdin
+ * ({@link ForkBootstrap}), so it is not visible in {@code
/proc/<pid>/cmdline} or
+ * {@code /proc/<pid>/environ}. Each client connection must present the token
before
+ * the server will accept it.
*
* @see PipesConfig#setUseSharedServer(boolean)
*/
@@ -70,7 +68,7 @@ public class SharedServerManager implements ServerManager {
public static final int SOCKET_CONNECT_TIMEOUT_MS = 60000;
private final PipesConfig pipesConfig;
- private final Path tikaConfigPath;
+ private final byte[] tikaConfigJson;
private final int numConnections;
private final Object lock = new Object();
@@ -89,12 +87,12 @@ public class SharedServerManager implements ServerManager {
* Creates a SharedServerManager.
*
* @param pipesConfig the pipes configuration
- * @param tikaConfigPath path to the tika config file
+ * @param tikaConfigJson the parent's resolved config, from {@link
ForkBootstrap#toBytes}
* @param numConnections number of concurrent connections the server
should support
*/
- public SharedServerManager(PipesConfig pipesConfig, Path tikaConfigPath,
int numConnections) {
+ public SharedServerManager(PipesConfig pipesConfig, byte[] tikaConfigJson,
int numConnections) {
this.pipesConfig = pipesConfig;
- this.tikaConfigPath = tikaConfigPath;
+ this.tikaConfigJson = tikaConfigJson;
this.numConnections = numConnections;
}
@@ -302,9 +300,7 @@ public class SharedServerManager implements ServerManager {
shutdownUnsafe();
}
- // Generate auth token for this server instance
- byte[] token = new byte[PipesServer.AUTH_TOKEN_LENGTH_BYTES];
- new SecureRandom().nextBytes(token);
+ byte[] token = ForkBootstrap.newToken();
currentToken = token;
LOG.info("\n\n" +
@@ -336,14 +332,10 @@ public class SharedServerManager implements ServerManager
{
tmpDir =
pipesConfig.createTempDirectory(PipesServer.SHARED_TEMP_DIR_PREFIX);
ProcessBuilder pb = new ProcessBuilder(getCommandline());
- // Pass port and auth token via environment variables so they are not
- // visible in /proc/<pid>/cmdline. The token is only readable via
- // /proc/<pid>/environ which requires same-uid access.
// Pass port=0 so the server binds to any available ephemeral port.
// The actual port is read back from the READY:{port} stdout signal,
// eliminating the TOCTOU race between probing a free port and binding
it.
pb.environment().put("TIKA_PIPES_PORT", "0");
- pb.environment().put("TIKA_PIPES_AUTH_TOKEN",
HexFormat.of().formatHex(token));
// Tell the child our PID so it can watch ProcessHandle.onExit() and
// self-terminate promptly if we die. See
PipesServer.watchParentProcess.
pb.environment().put(PipesServer.PARENT_PID_ENV,
@@ -382,6 +374,13 @@ public class SharedServerManager implements ServerManager {
throw new ServerInitializationException(msg, e);
}
+ try (OutputStream stdin = process.getOutputStream()) {
+ ForkBootstrap.write(stdin, token, tikaConfigJson);
+ } catch (IOException e) {
+ // The server died before reading it; waitForServerReady() reports
the exit.
+ LOG.warn("Couldn't send the bootstrap to the shared server
process", e);
+ }
+
// Wait for the server to signal it's ready and report the port it
actually bound to
serverPort = waitForServerReady();
LOG.info("Shared server started successfully on port {}", serverPort);
@@ -558,10 +557,9 @@ public class SharedServerManager implements ServerManager {
commandLine.add("-Djava.io.tmpdir=" + tmpDir.toAbsolutePath());
commandLine.add("org.apache.tika.pipes.core.server.PipesServer");
- // Shared mode arguments: port and auth token are passed via env vars
+ // Port is passed via env var; token and config via stdin
commandLine.add("--shared");
commandLine.add(Integer.toString(numConnections));
- commandLine.add(tikaConfigPath.toAbsolutePath().toString());
LOG.debug("Shared server commandline: {}", commandLine);
return commandLine.toArray(new String[0]);
diff --git
a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/async/AsyncProcessor.java
b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/async/AsyncProcessor.java
index 1a6b472bbc..6d2ab6d27a 100644
---
a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/async/AsyncProcessor.java
+++
b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/async/AsyncProcessor.java
@@ -56,6 +56,7 @@ import org.apache.tika.pipes.core.ServerManager;
import org.apache.tika.pipes.core.SharedServerManager;
import org.apache.tika.pipes.core.emitter.EmitDataImpl;
import org.apache.tika.pipes.core.emitter.EmitterManager;
+import org.apache.tika.pipes.core.protocol.ForkBootstrap;
import org.apache.tika.pipes.core.reporter.ReporterManager;
import org.apache.tika.plugins.TikaPluginManager;
@@ -76,7 +77,6 @@ public class AsyncProcessor implements Closeable {
private final ExecutorCompletionService<Integer> executorCompletionService;
private final ExecutorService executorService;
private final PipesConfig asyncConfig;
- private final Path tikaConfigPath;
private final PipesReporter pipesReporter;
private final List<ServerManager> serverManagers = new ArrayList<>();
private final AtomicLong totalProcessed = new AtomicLong(0);
@@ -120,15 +120,16 @@ public class AsyncProcessor implements Closeable {
throws TikaException, IOException {
TikaJsonConfig tikaJsonConfig = TikaJsonConfig.load(tikaConfigPath);
TikaPluginManager.preExtractPlugins(tikaJsonConfig);
- return new AsyncProcessor(tikaConfigPath, pipesIterator,
tikaJsonConfig);
+ return new AsyncProcessor(pipesIterator, tikaJsonConfig);
}
- private AsyncProcessor(Path tikaConfigPath, PipesIterator pipesIterator,
- TikaJsonConfig tikaJsonConfig) throws TikaException, IOException {
+ private AsyncProcessor(PipesIterator pipesIterator, TikaJsonConfig
tikaJsonConfig)
+ throws TikaException, IOException {
TikaPluginManager tikaPluginManager =
TikaPluginManager.load(tikaJsonConfig);
- MetadataFilter metadataFilter =
TikaLoader.load(tikaConfigPath).loadMetadataFilters();
+ MetadataFilter metadataFilter = TikaLoader.load(tikaJsonConfig,
+
Thread.currentThread().getContextClassLoader()).loadMetadataFilters();
this.asyncConfig = PipesConfig.load(tikaJsonConfig);
- this.tikaConfigPath = tikaConfigPath;
+ byte[] tikaConfigJson = ForkBootstrap.toBytes(tikaJsonConfig);
this.pipesReporter = ReporterManager.load(tikaPluginManager,
tikaJsonConfig);
LOG.debug("loaded reporter {}", pipesReporter.getClass());
this.fetchEmitTuples = new
ArrayBlockingQueue<>(asyncConfig.getQueueSize());
@@ -166,7 +167,7 @@ public class AsyncProcessor implements Closeable {
if (isSharedMode) {
LOG.info("Using shared server mode with {} workers",
asyncConfig.getNumClients());
SharedServerManager sharedManager = new SharedServerManager(
- asyncConfig, tikaConfigPath,
asyncConfig.getNumClients());
+ asyncConfig, tikaConfigJson,
asyncConfig.getNumClients());
serverManagers.add(sharedManager);
for (int i = 0; i < asyncConfig.getNumClients(); i++) {
@@ -178,7 +179,7 @@ public class AsyncProcessor implements Closeable {
LOG.info("Using per-client server mode with {} workers",
asyncConfig.getNumClients());
for (int i = 0; i < asyncConfig.getNumClients(); i++) {
PerClientServerManager serverManager = new
PerClientServerManager(
- asyncConfig, tikaConfigPath, i);
+ asyncConfig, tikaConfigJson, i);
serverManagers.add(serverManager);
executorCompletionService.submit(
diff --git
a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/config/ConfigMerger.java
b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/config/ConfigMerger.java
index f10ddd7a28..69a85940a2 100644
---
a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/config/ConfigMerger.java
+++
b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/config/ConfigMerger.java
@@ -16,6 +16,7 @@
*/
package org.apache.tika.pipes.core.config;
+import java.io.ByteArrayInputStream;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
@@ -33,7 +34,9 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.tika.config.TimeoutLimits;
+import org.apache.tika.config.loader.TikaJsonConfig;
import org.apache.tika.config.loader.TikaObjectMapperFactory;
+import org.apache.tika.exception.TikaConfigException;
import org.apache.tika.pipes.api.ComponentIds;
/**
@@ -50,10 +53,10 @@ import org.apache.tika.pipes.api.ComponentIds;
* <ul>
* <li>Uses UUID-based names for internal fetchers/emitters to avoid
conflicts with
* user-configured components</li>
- * <li>Returns a MergeResult containing the config path and generated names
so callers
- * can use them</li>
+ * <li>Returns the merged config and generated names so callers can use
them</li>
* <li>Preserves existing config sections when merging</li>
- * <li>Creates temp files that are marked for deletion on JVM exit</li>
+ * <li>{@link #merge} stays in memory; {@link #mergeOrCreate} writes a temp
file that is
+ * marked for deletion on JVM exit</li>
* </ul>
* <p>
* Example usage:
@@ -65,8 +68,9 @@ import org.apache.tika.pipes.api.ComponentIds;
* .setEmitStrategy(EmitStrategy.PASSBACK_ALL)
* .build();
*
- * MergeResult result = ConfigMerger.mergeOrCreate(existingConfigPath,
overrides);
- * // Use result.configPath() for PipesParser.load()
+ * MergedConfig merged = ConfigMerger.merge(existingConfigPath, overrides);
+ * TikaJsonConfig tikaJsonConfig = merged.load();
+ * PipesParser.load(tikaJsonConfig, PipesConfig.load(tikaJsonConfig));
* </pre>
*/
public class ConfigMerger {
@@ -90,6 +94,24 @@ public class ConfigMerger {
*/
public static MergeResult mergeOrCreate(Path existingConfig,
ConfigOverrides overrides)
throws IOException {
+ MergedConfig merged = merge(existingConfig, overrides);
+ Path tempConfig = Files.createTempFile("tika-config-merged-", ".json");
+ Files.write(tempConfig, merged.json());
+ tempConfig.toFile().deleteOnExit();
+ LOG.debug("Created merged config: {}", tempConfig);
+ return new MergeResult(tempConfig, merged.fetcherId(),
merged.emitterId());
+ }
+
+ /**
+ * Like {@link #mergeOrCreate} but in memory: nothing is written to disk.
+ *
+ * @param existingConfig path to existing config (may be null)
+ * @param overrides the overrides to apply
+ * @return the merged config JSON, env references unresolved, and the
generated ids
+ * @throws IOException if the existing config can't be read
+ */
+ public static MergedConfig merge(Path existingConfig, ConfigOverrides
overrides)
+ throws IOException {
// The shared config mapper: same comment and strictness rules as the
main loader.
ObjectMapper mapper = TikaObjectMapperFactory.getMapper();
@@ -214,18 +236,13 @@ public class ConfigMerger {
tl.getProgressTimeoutMillis());
}
- // Write merged config to temp file
- Path tempConfig = Files.createTempFile("tika-config-merged-", ".json");
-
mapper.writerWithDefaultPrettyPrinter().writeValue(tempConfig.toFile(), root);
- tempConfig.toFile().deleteOnExit();
-
- LOG.debug("Created merged config: {}", tempConfig);
+ byte[] json =
mapper.writerWithDefaultPrettyPrinter().writeValueAsBytes(root);
// Return the first generated fetcher/emitter ID (or null if none)
String primaryFetcherId = generatedFetcherIds.isEmpty() ? null :
generatedFetcherIds.get(0);
String primaryEmitterId = generatedEmitterIds.isEmpty() ? null :
generatedEmitterIds.get(0);
- return new MergeResult(tempConfig, primaryFetcherId, primaryEmitterId);
+ return new MergedConfig(json, primaryFetcherId, primaryEmitterId);
}
/**
@@ -308,4 +325,19 @@ public class ConfigMerger {
*/
public record MergeResult(Path configPath, String fetcherId, String
emitterId) {
}
+
+ /**
+ * Result of an in-memory config merge.
+ *
+ * @param json the merged configuration; load it with
+ * {@link TikaJsonConfig#load(java.io.InputStream)} so env
references resolve
+ * @param fetcherId the primary generated fetcher ID (may be null if no
fetchers were added)
+ * @param emitterId the primary generated emitter ID (may be null if no
emitters were added)
+ */
+ public record MergedConfig(byte[] json, String fetcherId, String
emitterId) {
+
+ public TikaJsonConfig load() throws TikaConfigException {
+ return TikaJsonConfig.load(new ByteArrayInputStream(json));
+ }
+ }
}
diff --git
a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/protocol/ForkBootstrap.java
b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/protocol/ForkBootstrap.java
new file mode 100644
index 0000000000..7c284c2389
--- /dev/null
+++
b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/protocol/ForkBootstrap.java
@@ -0,0 +1,116 @@
+/*
+ * 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.protocol;
+
+import java.io.ByteArrayInputStream;
+import java.io.DataInputStream;
+import java.io.DataOutputStream;
+import java.io.EOFException;
+import java.io.IOException;
+import java.io.InputStream;
+import java.io.OutputStream;
+import java.security.MessageDigest;
+import java.security.SecureRandom;
+
+import org.apache.tika.config.loader.TikaJsonConfig;
+import org.apache.tika.config.loader.TikaObjectMapperFactory;
+import org.apache.tika.exception.TikaConfigException;
+
+/**
+ * What the parent hands a forked PipesServer on its stdin before anything
else: an auth
+ * token and the parent's config. The fork presents the token as the first
bytes of every
+ * socket connection so the other end can reject a stranger that connected
first.
+ * <p>
+ * Frame: {@value #TOKEN_LENGTH_BYTES}-byte token, int config length, config
JSON.
+ * The config's {@code ${env:...}} references are resolved by the parent,
never by the fork.
+ */
+public final class ForkBootstrap {
+
+ public static final int TOKEN_LENGTH_BYTES = 32;
+
+ // Fixed: the configurable IPC limits live inside the config this frame
carries.
+ static final int MAX_CONFIG_BYTES = 64 * 1024 * 1024;
+
+ private static final SecureRandom RANDOM = new SecureRandom();
+
+ private final byte[] token;
+ private final byte[] config;
+
+ private ForkBootstrap(byte[] token, byte[] config) {
+ this.token = token;
+ this.config = config;
+ }
+
+ public static byte[] newToken() {
+ byte[] token = new byte[TOKEN_LENGTH_BYTES];
+ RANDOM.nextBytes(token);
+ return token;
+ }
+
+ /**
+ * Serializes the parent's loaded config. Its env references are already
resolved, so the
+ * bytes may carry secrets; they travel only over the fork's stdin.
+ */
+ public static byte[] toBytes(TikaJsonConfig tikaJsonConfig) throws
IOException {
+ return
TikaObjectMapperFactory.getMapper().writeValueAsBytes(tikaJsonConfig.getRootNode());
+ }
+
+ public static void write(OutputStream out, byte[] token, byte[] config)
throws IOException {
+ DataOutputStream dos = new DataOutputStream(out);
+ dos.write(token);
+ dos.writeInt(config.length);
+ dos.write(config);
+ dos.flush();
+ }
+
+ public static ForkBootstrap read(InputStream in) throws IOException {
+ DataInputStream dis = new DataInputStream(in);
+ byte[] token = new byte[TOKEN_LENGTH_BYTES];
+ dis.readFully(token);
+ int length = dis.readInt();
+ if (length < 0 || length > MAX_CONFIG_BYTES) {
+ throw new IOException("bootstrap config length " + length + " is
outside [0, "
+ + MAX_CONFIG_BYTES + "]");
+ }
+ byte[] config = new byte[length];
+ dis.readFully(config);
+ return new ForkBootstrap(token, config);
+ }
+
+ /**
+ * Reads {@value #TOKEN_LENGTH_BYTES} bytes and compares them to {@code
expected} in
+ * constant time. False on a short read or a mismatch.
+ */
+ public static boolean checkToken(InputStream in, byte[] expected) throws
IOException {
+ byte[] presented = new byte[TOKEN_LENGTH_BYTES];
+ try {
+ new DataInputStream(in).readFully(presented);
+ } catch (EOFException e) {
+ return false;
+ }
+ return MessageDigest.isEqual(expected, presented);
+ }
+
+ public byte[] getToken() {
+ return token;
+ }
+
+ /** Loads the config without resolving env references: the parent already
did. */
+ public TikaJsonConfig loadConfig() throws TikaConfigException {
+ return TikaJsonConfig.load(new ByteArrayInputStream(config), false);
+ }
+}
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 3dbe27a8fd..97c20dc91b 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
@@ -21,6 +21,7 @@ import java.io.BufferedOutputStream;
import java.io.DataInputStream;
import java.io.DataOutputStream;
import java.io.IOException;
+import java.io.InputStream;
import java.net.InetAddress;
import java.net.InetSocketAddress;
import java.net.Socket;
@@ -28,8 +29,6 @@ import java.net.SocketTimeoutException;
import java.nio.charset.StandardCharsets;
import java.nio.file.Path;
import java.nio.file.Paths;
-import java.security.MessageDigest;
-import java.util.HexFormat;
import java.util.Locale;
import java.util.Optional;
import java.util.concurrent.ArrayBlockingQueue;
@@ -74,6 +73,7 @@ import org.apache.tika.pipes.core.config.ConfigStore;
import org.apache.tika.pipes.core.config.ConfigStoreFactory;
import org.apache.tika.pipes.core.emitter.EmitterManager;
import org.apache.tika.pipes.core.fetcher.FetcherManager;
+import org.apache.tika.pipes.core.protocol.ForkBootstrap;
import org.apache.tika.pipes.core.protocol.PipesMessage;
import org.apache.tika.pipes.core.protocol.PipesMessageType;
import org.apache.tika.pipes.core.protocol.ShutDownReceivedException;
@@ -89,6 +89,9 @@ import org.apache.tika.utils.ExceptionUtils;
* This server is forked from the PipesClient. This class isolates
* parsing from the client to protect the primary JVM.
* <p>
+ * The config and auth token arrive on stdin as a {@link ForkBootstrap};
nothing is
+ * read from disk.
+ * <p>
* When configuring logging for this class, make absolutely certain
* not to write to STDOUT. This class uses STDOUT to communicate with
* the PipesClient.
@@ -156,8 +159,6 @@ public class PipesServer implements AutoCloseable {
}
}
- public static final int AUTH_TOKEN_LENGTH_BYTES = 32;
-
/** Env var the parent manager sets so the child can watch the parent's
* process handle and exit promptly if the parent dies. */
public static final String PARENT_PID_ENV = "TIKA_PIPES_PARENT_PID";
@@ -214,7 +215,7 @@ public class PipesServer implements AutoCloseable {
private long tIntermediateWaitNanos = -1;
private PipesWorker tLastWorker;
- public static PipesServer load(int port, Path tikaConfigPath) throws
Exception {
+ public static PipesServer load(int port, ForkBootstrap bootstrap) throws
Exception {
String pipesClientId = System.getProperty("pipesClientId",
"unknown");
LOG.debug("connecting to client on port={}", port);
Socket socket = new Socket();
@@ -223,9 +224,12 @@ public class PipesServer implements AutoCloseable {
DataInputStream dis = new DataInputStream(new
BufferedInputStream(socket.getInputStream()));
DataOutputStream dos = new DataOutputStream(new
BufferedOutputStream(socket.getOutputStream()));
+ dos.write(bootstrap.getToken());
+ dos.flush();
ParseContext configContext = null;
try {
- TikaLoader tikaLoader = TikaLoader.load(tikaConfigPath);
+ TikaLoader tikaLoader = TikaLoader.load(bootstrap.loadConfig(),
+ Thread.currentThread().getContextClassLoader());
//load before constructing so the catch block below can redact per
config
configContext = tikaLoader.loadParseContext();
PipesServer pipesServer =
@@ -281,29 +285,26 @@ public class PipesServer implements AutoCloseable {
// below don't strand us if the parent has already died.
watchParentProcess();
- // Check for shared mode: --shared <numConnections> <tikaConfigPath>
+ ForkBootstrap bootstrap;
+ try (InputStream stdin = System.in) {
+ bootstrap = ForkBootstrap.read(stdin);
+ }
+
+ // Check for shared mode: --shared <numConnections>
if (args.length > 0 && "--shared".equals(args[0])) {
String portEnv = System.getenv("TIKA_PIPES_PORT");
if (portEnv == null || portEnv.isEmpty()) {
throw new IllegalStateException("TIKA_PIPES_PORT environment
variable is not set");
}
int port = Integer.parseInt(portEnv);
- String tokenHex = System.getenv("TIKA_PIPES_AUTH_TOKEN");
- if (tokenHex == null || tokenHex.isEmpty()) {
- throw new IllegalStateException("TIKA_PIPES_AUTH_TOKEN
environment variable is not set");
- }
- byte[] expectedToken = HexFormat.of().parseHex(tokenHex);
int numConnections = Integer.parseInt(args[1]);
- Path tikaConfig = Paths.get(args[2]);
LOG.info("Starting shared PipesServer with {} connections",
numConnections);
- runSharedMode(port, numConnections, tikaConfig, expectedToken);
+ runSharedMode(port, numConnections, bootstrap);
} else {
- // Per-client mode: <port> <tikaConfigPath>
+ // Per-client mode: <port>
int port = Integer.parseInt(args[0]);
- Path tikaConfig = Paths.get(args[1]);
- String pipesClientId = System.getProperty("pipesClientId",
"unknown");
LOG.debug("starting pipes server on port={}", port);
- try (PipesServer server = PipesServer.load(port, tikaConfig)) {
+ try (PipesServer server = PipesServer.load(port, bootstrap)) {
server.mainLoop();
} catch (Throwable t) {
LOG.error("crashed", t);
@@ -339,16 +340,13 @@ public class PipesServer implements AutoCloseable {
/**
* Runs the server in shared mode, accepting multiple client connections.
* <p>
- * Each incoming connection must present a valid auth token (32 bytes)
before
- * being accepted. This prevents unauthorized local processes from
connecting.
- * Note: if a malicious actor has access to your localhost and can read
- * /proc/<pid>/environ, that is beyond Tika's security model. This
auth
- * token exists to prevent CVE-style abuse from untrusted local processes
that
- * cannot read the server process's environment.
+ * Each incoming connection must present the bootstrap's auth token before
being
+ * accepted, so a stray local process that finds the port is turned away.
*/
- private static void runSharedMode(int port, int numConnections, Path
tikaConfigPath,
- byte[] expectedToken) throws Exception {
- TikaLoader tikaLoader = TikaLoader.load(tikaConfigPath);
+ private static void runSharedMode(int port, int numConnections,
ForkBootstrap bootstrap)
+ throws Exception {
+ TikaLoader tikaLoader = TikaLoader.load(bootstrap.loadConfig(),
+ Thread.currentThread().getContextClassLoader());
PipesConfig pipesConfig = PipesConfig.load(tikaLoader.getConfig());
validateHeartbeatInterval(pipesConfig);
@@ -377,18 +375,7 @@ public class PipesServer implements AutoCloseable {
clientSocket.setTcpNoDelay(true);
// Validate auth token before creating handler
- byte[] clientToken = new byte[AUTH_TOKEN_LENGTH_BYTES];
- int bytesRead = 0;
- while (bytesRead < AUTH_TOKEN_LENGTH_BYTES) {
- int r = clientSocket.getInputStream().read(
- clientToken, bytesRead,
AUTH_TOKEN_LENGTH_BYTES - bytesRead);
- if (r == -1) {
- break;
- }
- bytesRead += r;
- }
- if (bytesRead < AUTH_TOKEN_LENGTH_BYTES ||
- !MessageDigest.isEqual(expectedToken,
clientToken)) {
+ if
(!ForkBootstrap.checkToken(clientSocket.getInputStream(),
bootstrap.getToken())) {
LOG.warn("Rejected connection with invalid auth
token");
try {
clientSocket.close();
diff --git
a/tika-pipes/tika-pipes-core/src/test/java/org/apache/tika/pipes/core/PerClientServerManagerSizingTest.java
b/tika-pipes/tika-pipes-core/src/test/java/org/apache/tika/pipes/core/PerClientServerManagerSizingTest.java
index eb3cdb3167..d0e20d0614 100644
---
a/tika-pipes/tika-pipes-core/src/test/java/org/apache/tika/pipes/core/PerClientServerManagerSizingTest.java
+++
b/tika-pipes/tika-pipes-core/src/test/java/org/apache/tika/pipes/core/PerClientServerManagerSizingTest.java
@@ -43,7 +43,7 @@ public class PerClientServerManagerSizingTest {
pipesConfig.setNumClients(numClients);
pipesConfig.setForkedJvmArgs(new
ArrayList<>(Arrays.asList(forkedJvmArgs)));
PerClientServerManager manager =
- new PerClientServerManager(pipesConfig,
tmp.resolve("tika-config.json"), 0);
+ new PerClientServerManager(pipesConfig, new byte[0], 0);
return Arrays.asList(manager.getCommandline(tmp));
}
@@ -51,6 +51,13 @@ public class PerClientServerManagerSizingTest {
return args.stream().filter(a -> a.startsWith(prefix)).toList();
}
+ /** The config travels on stdin; only the port follows the main class. */
+ @Test
+ public void noConfigPathOnCommandLine() throws Exception {
+ List<String> args = commandLine(1);
+ assertEquals("org.apache.tika.pipes.core.server.PipesServer",
args.get(args.size() - 2));
+ }
+
/** A lone fork is capped below the fork budget; the parent claims the JVM
default on top. */
@Test
public void heapInjectedForSingleClient() throws Exception {
diff --git
a/tika-pipes/tika-pipes-core/src/test/java/org/apache/tika/pipes/core/protocol/ForkBootstrapTest.java
b/tika-pipes/tika-pipes-core/src/test/java/org/apache/tika/pipes/core/protocol/ForkBootstrapTest.java
new file mode 100644
index 0000000000..35e25d4500
--- /dev/null
+++
b/tika-pipes/tika-pipes-core/src/test/java/org/apache/tika/pipes/core/protocol/ForkBootstrapTest.java
@@ -0,0 +1,118 @@
+/*
+ * 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.protocol;
+
+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.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.io.ByteArrayInputStream;
+import java.io.ByteArrayOutputStream;
+import java.io.DataOutputStream;
+import java.io.EOFException;
+import java.io.IOException;
+import java.nio.charset.StandardCharsets;
+import java.util.Arrays;
+
+import org.junit.jupiter.api.Test;
+
+import org.apache.tika.config.loader.TikaJsonConfig;
+
+public class ForkBootstrapTest {
+
+ private static String basePath(TikaJsonConfig config) {
+ return
config.getRootNode().get("fetchers").get("f").get("file-system-fetcher")
+ .get("basePath").textValue();
+ }
+
+ private static TikaJsonConfig parentLoad(String json) throws Exception {
+ return TikaJsonConfig.load(new
ByteArrayInputStream(json.getBytes(StandardCharsets.UTF_8)));
+ }
+
+ @Test
+ public void testRoundTrip() throws Exception {
+ byte[] token = ForkBootstrap.newToken();
+ byte[] config = "{\"fetchers\":{}}".getBytes(StandardCharsets.UTF_8);
+ ByteArrayOutputStream bos = new ByteArrayOutputStream();
+ ForkBootstrap.write(bos, token, config);
+
+ ForkBootstrap read = ForkBootstrap.read(new
ByteArrayInputStream(bos.toByteArray()));
+ assertArrayEquals(token, read.getToken());
+ assertNotNull(read.loadConfig().getRootNode().get("fetchers"));
+ }
+
+ @Test
+ public void testParentResolvesEnv() throws Exception {
+ String path = System.getenv("PATH");
+ assertNotNull(path, "test needs PATH in the environment");
+ TikaJsonConfig parent = parentLoad(
+
"{\"fetchers\":{\"f\":{\"file-system-fetcher\":{\"basePath\":\"${env:PATH}\"}}}}");
+ ForkBootstrap fork = sendThrough(ForkBootstrap.toBytes(parent));
+ assertEquals(path, basePath(fork.loadConfig()));
+ }
+
+ @Test
+ public void testForkDoesNotResolveEnvAgain() throws Exception {
+ // What a parent-resolved secret whose value contains a reference
looks like on the wire.
+ // Resolving it in the fork would fail on the unset variable or expand
it a second time.
+ String literal = "s3cret-${env:TIKA_FORK_BOOTSTRAP_TEST_UNSET}";
+ byte[] resolved =
("{\"fetchers\":{\"f\":{\"file-system-fetcher\":{\"basePath\":\""
+ + literal + "\"}}}}").getBytes(StandardCharsets.UTF_8);
+ assertEquals(literal, basePath(sendThrough(resolved).loadConfig()));
+ }
+
+ @Test
+ public void testConfigLengthIsBounded() throws Exception {
+ ByteArrayOutputStream bos = new ByteArrayOutputStream();
+ DataOutputStream dos = new DataOutputStream(bos);
+ dos.write(new byte[ForkBootstrap.TOKEN_LENGTH_BYTES]);
+ dos.writeInt(ForkBootstrap.MAX_CONFIG_BYTES + 1);
+ assertThrows(IOException.class,
+ () -> ForkBootstrap.read(new
ByteArrayInputStream(bos.toByteArray())));
+ }
+
+ @Test
+ public void testTruncatedBootstrapFails() throws Exception {
+ ByteArrayOutputStream bos = new ByteArrayOutputStream();
+ ForkBootstrap.write(bos, ForkBootstrap.newToken(), new byte[100]);
+ byte[] truncated = Arrays.copyOf(bos.toByteArray(), bos.size() - 1);
+ assertThrows(EOFException.class,
+ () -> ForkBootstrap.read(new ByteArrayInputStream(truncated)));
+ }
+
+ @Test
+ public void testCheckToken() throws Exception {
+ byte[] token = ForkBootstrap.newToken();
+ assertTrue(ForkBootstrap.checkToken(new ByteArrayInputStream(token),
token));
+
+ byte[] wrong = token.clone();
+ wrong[0] ^= 1;
+ assertFalse(ForkBootstrap.checkToken(new ByteArrayInputStream(wrong),
token));
+
+ byte[] shortRead = Arrays.copyOf(token, token.length - 1);
+ assertFalse(ForkBootstrap.checkToken(new
ByteArrayInputStream(shortRead), token));
+ }
+
+ private static ForkBootstrap sendThrough(byte[] config) throws IOException
{
+ ByteArrayOutputStream bos = new ByteArrayOutputStream();
+ ForkBootstrap.write(bos, ForkBootstrap.newToken(), config);
+ return ForkBootstrap.read(new ByteArrayInputStream(bos.toByteArray()));
+ }
+}
diff --git
a/tika-pipes/tika-pipes-fork-parser/src/main/java/org/apache/tika/pipes/fork/PipesForkParser.java
b/tika-pipes/tika-pipes-fork-parser/src/main/java/org/apache/tika/pipes/fork/PipesForkParser.java
index 1719601228..9e9e893570 100644
---
a/tika-pipes/tika-pipes-fork-parser/src/main/java/org/apache/tika/pipes/fork/PipesForkParser.java
+++
b/tika-pipes/tika-pipes-fork-parser/src/main/java/org/apache/tika/pipes/fork/PipesForkParser.java
@@ -24,6 +24,7 @@ import java.util.Map;
import java.util.UUID;
import org.apache.tika.config.EmbeddedLimits;
+import org.apache.tika.config.loader.TikaJsonConfig;
import org.apache.tika.exception.TikaConfigException;
import org.apache.tika.exception.TikaException;
import org.apache.tika.io.TikaInputStream;
@@ -121,13 +122,12 @@ public class PipesForkParser implements Closeable {
private final PipesForkParserConfig config;
private final PipesParser pipesParser;
- private final Path tikaConfigPath;
private final String internalFetcherId;
/**
* Creates a new PipesForkParser with default configuration.
*
- * @throws IOException if the temporary config file cannot be created
+ * @throws IOException if the user config cannot be read
* @throws TikaConfigException if configuration is invalid
*/
public PipesForkParser() throws IOException, TikaConfigException {
@@ -138,17 +138,17 @@ public class PipesForkParser implements Closeable {
* Creates a new PipesForkParser with the specified configuration.
*
* @param config the configuration for this parser
- * @throws IOException if the temporary config file cannot be created
+ * @throws IOException if the user config cannot be read
* @throws TikaConfigException if configuration is invalid
*/
public PipesForkParser(PipesForkParserConfig config) throws IOException,
TikaConfigException {
this.config = config;
// Jackson-deserialized configs are checked on binding; setter-built
ones are checked here.
config.getPipesConfig().checkPayloadLimits();
- ConfigMerger.MergeResult mergeResult = createTikaConfigFile();
- this.tikaConfigPath = mergeResult.configPath();
- this.internalFetcherId = mergeResult.fetcherId();
- this.pipesParser = PipesParser.load(tikaConfigPath);
+ ConfigMerger.MergedConfig merged = createTikaConfig();
+ this.internalFetcherId = merged.fetcherId();
+ TikaJsonConfig tikaJsonConfig = merged.load();
+ this.pipesParser = PipesParser.load(tikaJsonConfig,
PipesConfig.load(tikaJsonConfig));
}
/**
@@ -402,23 +402,19 @@ public class PipesForkParser implements Closeable {
@Override
public void close() throws IOException {
pipesParser.close();
- // Clean up temp config file
- if (tikaConfigPath != null) {
- Files.deleteIfExists(tikaConfigPath);
- }
}
/**
- * Creates a temporary tika-config.json file for the forked process.
+ * Builds the tika-config for the forked processes, in memory.
* <p>
* Uses ConfigMerger to:
* - Add a FileSystemFetcher with UUID-based name for absolute path access
* - Set PASSBACK_ALL emit strategy (no emitter, return results to client)
* - Merge with user config if provided
*
- * @return MergeResult containing the config path and generated fetcher ID
+ * @return the merged config and generated fetcher ID
*/
- private ConfigMerger.MergeResult createTikaConfigFile() throws IOException
{
+ private ConfigMerger.MergedConfig createTikaConfig() throws IOException {
PipesConfig pc = config.getPipesConfig();
// Build configuration overrides
@@ -454,7 +450,7 @@ public class PipesForkParser implements Closeable {
ConfigOverrides overrides = builder.build();
// Merge with user config if provided, otherwise create new
- return ConfigMerger.mergeOrCreate(config.getUserConfigPath(),
overrides);
+ return ConfigMerger.merge(config.getUserConfigPath(), overrides);
}
diff --git
a/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/PipesForkTokenTest.java
b/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/PipesForkTokenTest.java
new file mode 100644
index 0000000000..cc4c08f071
--- /dev/null
+++
b/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/PipesForkTokenTest.java
@@ -0,0 +1,92 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.tika.pipes.core;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+import java.io.IOException;
+import java.io.InputStream;
+import java.net.InetAddress;
+import java.net.InetSocketAddress;
+import java.net.Socket;
+import java.nio.file.Path;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CompletionException;
+import java.util.concurrent.TimeUnit;
+
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+import org.apache.tika.config.loader.TikaJsonConfig;
+import org.apache.tika.metadata.Metadata;
+import org.apache.tika.parser.ParseContext;
+import org.apache.tika.pipes.api.FetchEmitTuple;
+import org.apache.tika.pipes.api.PipesResult;
+import org.apache.tika.pipes.api.emitter.EmitKey;
+import org.apache.tika.pipes.api.fetcher.FetchKey;
+import org.apache.tika.pipes.core.protocol.ForkBootstrap;
+
+/**
+ * A per-client fork presents the token it got on stdin; the parent turns away
anything else.
+ */
+public class PipesForkTokenTest {
+
+ private static final String TEST_DOC = "testOverlappingText.pdf";
+
+ /** A local process that reaches the parent's port before the fork is
turned away. */
+ @Test
+ public void perClientRejectsConnectionWithoutToken(@TempDir Path tmp)
throws Exception {
+ Path tikaConfigPath = PluginsTestHelper.getFileSystemFetcherConfig(
+ tmp, tmp.resolve("input"), tmp.resolve("output"));
+ PluginsTestHelper.copyTestFilesToTmpInput(tmp, TEST_DOC);
+ TikaJsonConfig tikaJsonConfig = TikaJsonConfig.load(tikaConfigPath);
+ PipesConfig pipesConfig = PipesConfig.load(tikaJsonConfig);
+ try (PerClientServerManager manager = new PerClientServerManager(
+ pipesConfig, ForkBootstrap.toBytes(tikaJsonConfig), 0);
+ PipesClient client = new PipesClient(pipesConfig, manager)) {
+ // The port is bound before the fork is launched; connecting as
soon as it appears
+ // puts the stranger first in the backlog regardless of how fast
the fork's JVM starts.
+ CompletableFuture<Socket> strangerFuture =
CompletableFuture.supplyAsync(() -> {
+ try {
+ while (manager.getPort() <= 0) {
+ Thread.sleep(1);
+ }
+ Socket stranger = new Socket();
+ stranger.connect(new
InetSocketAddress(InetAddress.getLoopbackAddress(),
+ manager.getPort()));
+ stranger.getOutputStream().write(new
byte[ForkBootstrap.TOKEN_LENGTH_BYTES]);
+ stranger.getOutputStream().flush();
+ return stranger;
+ } catch (IOException | InterruptedException e) {
+ throw new CompletionException(e);
+ }
+ });
+ manager.ensureRunning();
+ try (Socket stranger = strangerFuture.get(30, TimeUnit.SECONDS)) {
+ PipesResult result = client.process(new
FetchEmitTuple(TEST_DOC,
+ new FetchKey("fsf", TEST_DOC), new EmitKey(), new
Metadata(),
+ new ParseContext(),
FetchEmitTuple.ON_PARSE_EXCEPTION.SKIP));
+ assertEquals(PipesResult.RESULT_STATUS.PARSE_SUCCESS,
result.status());
+
+ stranger.setSoTimeout(30000);
+ try (InputStream in = stranger.getInputStream()) {
+ assertEquals(-1, in.read(), "the parent must close the
stranger's connection");
+ }
+ }
+ }
+ }
+}
diff --git
a/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/PipesParserWarmClientTest.java
b/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/PipesParserWarmClientTest.java
index 486fcc43ea..f6671b2b63 100644
---
a/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/PipesParserWarmClientTest.java
+++
b/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/PipesParserWarmClientTest.java
@@ -49,7 +49,7 @@ public class PipesParserWarmClientTest {
PipesConfig pipesConfig = PipesConfig.load(tikaJsonConfig);
pipesConfig.setNumClients(3);
- try (PipesParser parser = PipesParser.load(tikaJsonConfig,
pipesConfig, tikaConfigPath)) {
+ try (PipesParser parser = PipesParser.load(tikaJsonConfig,
pipesConfig)) {
for (int i = 0; i < 4; i++) {
PipesResult result = parser.parse(
new FetchEmitTuple(testDoc + "-" + i, new
FetchKey("fsf", testDoc),
diff --git
a/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/SharedServerChaosMonkeyTest.java
b/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/SharedServerChaosMonkeyTest.java
index 324520d106..52f7bb44f2 100644
---
a/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/SharedServerChaosMonkeyTest.java
+++
b/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/SharedServerChaosMonkeyTest.java
@@ -144,7 +144,7 @@ public class SharedServerChaosMonkeyTest {
AtomicInteger observedCrash = new AtomicInteger(0);
AtomicInteger collateralDamage = new AtomicInteger(0); // OK files
that failed due to concurrent crash
- try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig, tikaConfigPath)) {
+ try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig)) {
assertTrue(pipesParser.isSharedMode(), "Should be in shared mode");
ExecutorService executor = Executors.newFixedThreadPool(8);
@@ -281,7 +281,7 @@ public class SharedServerChaosMonkeyTest {
TikaJsonConfig tikaJsonConfig = TikaJsonConfig.load(tikaConfigPath);
PipesConfig pipesConfig = PipesConfig.load(tikaJsonConfig);
- try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig, tikaConfigPath)) {
+ try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig)) {
// Process chaos files (some will crash)
for (int i = 0; i < 20; i++) {
pipesParser.parse(new FetchEmitTuple(
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 b7f0e15bc6..5913ce93e8 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
@@ -104,7 +104,7 @@ public class SharedServerModeTest {
// Verify shared mode is enabled
assertTrue(pipesConfig.isUseSharedServer(), "Shared server mode should
be enabled");
- try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig, tikaConfigPath)) {
+ try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig)) {
assertTrue(pipesParser.isSharedMode(), "PipesParser should be in
shared mode");
PipesResult result = pipesParser.parse(new FetchEmitTuple(
@@ -137,7 +137,7 @@ public class SharedServerModeTest {
TikaJsonConfig tikaJsonConfig = TikaJsonConfig.load(tikaConfigPath);
PipesConfig pipesConfig = PipesConfig.load(tikaJsonConfig);
- try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig, tikaConfigPath)) {
+ try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig)) {
ExecutorService executor = Executors.newFixedThreadPool(8);
List<Future<PipesResult>> futures = new ArrayList<>();
@@ -183,7 +183,7 @@ public class SharedServerModeTest {
TikaJsonConfig tikaJsonConfig = TikaJsonConfig.load(tikaConfigPath);
PipesConfig pipesConfig = PipesConfig.load(tikaJsonConfig);
- try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig, tikaConfigPath)) {
+ try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig)) {
// Process files sequentially with the same PipesParser
for (int i = 1; i <= 3; i++) {
String fileName = "file" + i + ".xml";
@@ -213,7 +213,7 @@ public class SharedServerModeTest {
// Create and close parser multiple times to verify graceful
shutdown/restart
for (int iteration = 0; iteration < 3; iteration++) {
- try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig, tikaConfigPath)) {
+ try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig)) {
PipesResult result = pipesParser.parse(new FetchEmitTuple(
"test.xml",
new FetchKey(FETCHER_NAME, "test.xml"),
@@ -257,7 +257,7 @@ public class SharedServerModeTest {
TikaJsonConfig tikaJsonConfig = TikaJsonConfig.load(tikaConfigPath);
PipesConfig pipesConfig = PipesConfig.load(tikaJsonConfig);
- PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig, tikaConfigPath);
+ PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig);
ExecutorService executor = Executors.newSingleThreadExecutor();
try {
// Warmup so the shared server is fully started; without this,
short
@@ -316,7 +316,7 @@ public class SharedServerModeTest {
// Verify shared mode is NOT enabled (default)
assertTrue(!pipesConfig.isUseSharedServer(), "Shared server mode
should be disabled by default");
- try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig, tikaConfigPath)) {
+ try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig)) {
assertTrue(!pipesParser.isSharedMode(), "PipesParser should NOT be
in shared mode");
PipesResult result = pipesParser.parse(new FetchEmitTuple(
@@ -348,7 +348,7 @@ public class SharedServerModeTest {
PipesConfig perClientPipesConfig = PipesConfig.load(perClientConfig);
Metadata perClientMetadata;
- try (PipesParser parser = PipesParser.load(perClientConfig,
perClientPipesConfig, perClientConfigPath)) {
+ try (PipesParser parser = PipesParser.load(perClientConfig,
perClientPipesConfig)) {
PipesResult result = parser.parse(new FetchEmitTuple(
testFile,
new FetchKey(FETCHER_NAME, testFile),
@@ -367,7 +367,7 @@ public class SharedServerModeTest {
PipesConfig sharedPipesConfig = PipesConfig.load(sharedConfig);
Metadata sharedMetadata;
- try (PipesParser parser = PipesParser.load(sharedConfig,
sharedPipesConfig, sharedConfigPath)) {
+ try (PipesParser parser = PipesParser.load(sharedConfig,
sharedPipesConfig)) {
PipesResult result = parser.parse(new FetchEmitTuple(
testFile,
new FetchKey(FETCHER_NAME, testFile),
@@ -406,7 +406,7 @@ public class SharedServerModeTest {
TikaJsonConfig tikaJsonConfig = TikaJsonConfig.load(tikaConfigPath);
PipesConfig pipesConfig = PipesConfig.load(tikaJsonConfig);
- try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig, tikaConfigPath)) {
+ try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig)) {
// First, trigger OOM
PipesResult oomResult = pipesParser.parse(new FetchEmitTuple(
"oom.xml",
@@ -451,7 +451,7 @@ public class SharedServerModeTest {
TikaJsonConfig tikaJsonConfig = TikaJsonConfig.load(tikaConfigPath);
PipesConfig pipesConfig = PipesConfig.load(tikaJsonConfig);
- try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig, tikaConfigPath)) {
+ try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig)) {
// First, make a request to ensure the server is started (lazy
initialization)
PipesResult warmupResult = pipesParser.parse(new FetchEmitTuple(
"warmup.xml",
@@ -511,7 +511,7 @@ public class SharedServerModeTest {
TikaJsonConfig tikaJsonConfig = TikaJsonConfig.load(tikaConfigPath);
PipesConfig pipesConfig = PipesConfig.load(tikaJsonConfig);
- try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig, tikaConfigPath)) {
+ try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig)) {
for (int i = 0; i < 3; i++) {
// Trigger OOM
PipesResult oomResult = pipesParser.parse(new FetchEmitTuple(
@@ -557,7 +557,7 @@ public class SharedServerModeTest {
TikaJsonConfig tikaJsonConfig = TikaJsonConfig.load(tikaConfigPath);
PipesConfig pipesConfig = PipesConfig.load(tikaJsonConfig);
- try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig, tikaConfigPath)) {
+ try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig)) {
ExecutorService executor = Executors.newFixedThreadPool(6);
List<Future<PipesResult>> futures = new ArrayList<>();
@@ -640,7 +640,7 @@ public class SharedServerModeTest {
TikaJsonConfig tikaJsonConfig = TikaJsonConfig.load(tikaConfigPath);
PipesConfig pipesConfig = PipesConfig.load(tikaJsonConfig);
- try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig, tikaConfigPath)) {
+ try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig)) {
// Trigger timeout
PipesResult timeoutResult = pipesParser.parse(new FetchEmitTuple(
"timeout.xml",
@@ -690,7 +690,7 @@ public class SharedServerModeTest {
TikaJsonConfig tikaJsonConfig = TikaJsonConfig.load(tikaConfigPath);
PipesConfig pipesConfig = PipesConfig.load(tikaJsonConfig);
- try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig, tikaConfigPath)) {
+ try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig)) {
ExecutorService executor = Executors.newFixedThreadPool(6);
// Phase 1: Submit 5 slow requests + 1 OOM concurrently
@@ -835,7 +835,7 @@ public class SharedServerModeTest {
TikaJsonConfig tikaJsonConfig = TikaJsonConfig.load(tikaConfigPath);
PipesConfig pipesConfig = PipesConfig.load(tikaJsonConfig);
- try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig, tikaConfigPath)) {
+ try (PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig)) {
ExecutorService executor = Executors.newFixedThreadPool(6);
try {
List<Future<PipesResult>> futures = new ArrayList<>();
diff --git
a/tika-serialization/src/main/java/org/apache/tika/config/loader/TikaJsonConfig.java
b/tika-serialization/src/main/java/org/apache/tika/config/loader/TikaJsonConfig.java
index 7c363054a6..4349294b9f 100644
---
a/tika-serialization/src/main/java/org/apache/tika/config/loader/TikaJsonConfig.java
+++
b/tika-serialization/src/main/java/org/apache/tika/config/loader/TikaJsonConfig.java
@@ -170,10 +170,27 @@ public class TikaJsonConfig {
* @throws TikaConfigException if loading or parsing fails
*/
public static TikaJsonConfig load(InputStream inputStream) throws
TikaConfigException {
+ return load(inputStream, true);
+ }
+
+ /**
+ * Loads configuration from an input stream, resolving {@code ${env:NAME}}
references only
+ * when {@code resolveEnv} is true. Pass false for JSON whose references
were already
+ * resolved, so that a resolved value containing {@code ${env:...}} is not
expanded again.
+ *
+ * @param inputStream the input stream containing JSON configuration
+ * @param resolveEnv whether to resolve environment variable references
+ * @return the parsed configuration
+ * @throws TikaConfigException if loading or parsing fails
+ */
+ public static TikaJsonConfig load(InputStream inputStream, boolean
resolveEnv)
+ throws TikaConfigException {
try {
JsonNode rootNode = OBJECT_MAPPER.readTree(inputStream);
- // startup config only: a request's JSON never reaches this method
- EnvInterpolator.resolve(rootNode, System::getenv);
+ if (resolveEnv) {
+ // startup config only: a request's JSON never reaches this
method
+ EnvInterpolator.resolve(rootNode, System::getenv);
+ }
TikaJsonConfig tikaJsonConfig = new TikaJsonConfig(rootNode);
tikaJsonConfig.validateKeys();
return tikaJsonConfig;
diff --git
a/tika-serialization/src/main/java/org/apache/tika/config/loader/TikaLoader.java
b/tika-serialization/src/main/java/org/apache/tika/config/loader/TikaLoader.java
index 4e96125080..ef9b422016 100644
---
a/tika-serialization/src/main/java/org/apache/tika/config/loader/TikaLoader.java
+++
b/tika-serialization/src/main/java/org/apache/tika/config/loader/TikaLoader.java
@@ -235,7 +235,20 @@ public class TikaLoader {
*/
public static TikaLoader load(Path configPath, ClassLoader classLoader)
throws TikaConfigException, IOException {
- TikaJsonConfig config = TikaJsonConfig.load(configPath);
+ return load(TikaJsonConfig.load(configPath), classLoader);
+ }
+
+ /**
+ * Creates a Tika loader from an already-loaded configuration.
+ * Global settings are automatically loaded and applied during
initialization.
+ *
+ * @param config the configuration
+ * @param classLoader the class loader to use for loading components
+ * @return the Tika loader
+ * @throws TikaConfigException if loading global settings fails
+ */
+ public static TikaLoader load(TikaJsonConfig config, ClassLoader
classLoader)
+ throws TikaConfigException, IOException {
TikaLoader loader = new TikaLoader(config, classLoader);
loader.init();
return loader;
diff --git
a/tika-server/tika-server-core/src/main/java/org/apache/tika/server/core/TikaServerProcess.java
b/tika-server/tika-server-core/src/main/java/org/apache/tika/server/core/TikaServerProcess.java
index 01302d2c1c..be314cd92c 100644
---
a/tika-server/tika-server-core/src/main/java/org/apache/tika/server/core/TikaServerProcess.java
+++
b/tika-server/tika-server-core/src/main/java/org/apache/tika/server/core/TikaServerProcess.java
@@ -700,9 +700,8 @@ public class TikaServerProcess {
// Create or merge config with server components using ConfigMerger
Path existingConfigPath = tikaServerConfig.hasConfigFile() ?
tikaServerConfig.getConfigPath() : null;
- Path configPath = createServerConfig(existingConfigPath,
+ TikaJsonConfig tikaJsonConfig = createServerConfig(existingConfigPath,
inputTempDirectory, unpackTempDirectory);
- TikaJsonConfig tikaJsonConfig = TikaJsonConfig.load(configPath);
// Load or create PipesConfig with defaults
PipesConfig pipesConfig = tikaJsonConfig.deserialize("pipes",
PipesConfig.class);
@@ -714,7 +713,7 @@ public class TikaServerProcess {
pipesConfig.setEmitStrategy(new
EmitStrategyConfig(EmitStrategy.PASSBACK_ALL));
// Create PipesParser
- PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig, configPath);
+ PipesParser pipesParser = PipesParser.load(tikaJsonConfig,
pipesConfig);
// Create and return the helper
PipesParsingHelper helper = new PipesParsingHelper(pipesParser,
pipesConfig,
@@ -740,11 +739,12 @@ public class TikaServerProcess {
* @param existingConfigPath the existing config file path (may be null)
* @param inputTempDirectory the temp directory for input files
* @param unpackTempDirectory the temp directory for unpack output files
(may be null)
- * @return the config path to use
+ * @return the merged config
*/
- private static Path createServerConfig(Path existingConfigPath,
- Path inputTempDirectory,
- Path unpackTempDirectory) throws
IOException {
+ private static TikaJsonConfig createServerConfig(Path existingConfigPath,
+ Path inputTempDirectory,
+ Path unpackTempDirectory)
+ throws IOException, TikaConfigException {
LOG.info("Configuring {} with basePath={}",
PipesParsingHelper.DEFAULT_FETCHER_ID, inputTempDirectory);
// Build configuration overrides
@@ -779,10 +779,7 @@ public class TikaServerProcess {
ConfigOverrides overrides = builder.build();
// Merge with existing config or create new
- ConfigMerger.MergeResult result =
ConfigMerger.mergeOrCreate(existingConfigPath, overrides);
-
- LOG.debug("Created server config: {}", result.configPath());
- return result.configPath();
+ return ConfigMerger.merge(existingConfigPath, overrides).load();
}
private static class ServerDetails {
diff --git
a/tika-server/tika-server-core/src/test/java/org/apache/tika/server/core/CXFTestBase.java
b/tika-server/tika-server-core/src/test/java/org/apache/tika/server/core/CXFTestBase.java
index e21913a1f9..6e2748e2ae 100644
---
a/tika-server/tika-server-core/src/test/java/org/apache/tika/server/core/CXFTestBase.java
+++
b/tika-server/tika-server-core/src/test/java/org/apache/tika/server/core/CXFTestBase.java
@@ -213,7 +213,7 @@ public abstract class CXFTestBase {
pipesConfig = new PipesConfig();
}
pipesConfig.setEmitStrategy(new
EmitStrategyConfig(EmitStrategy.PASSBACK_ALL));
- this.pipesParser = PipesParser.load(tikaJsonConfig, pipesConfig,
this.pipesConfigPath);
+ this.pipesParser = PipesParser.load(tikaJsonConfig, pipesConfig);
PipesParsingHelper pipesParsingHelper = new
PipesParsingHelper(this.pipesParser, pipesConfig,
inputTempDirectory, getUnpackEmitterBasePath());
diff --git
a/tika-server/tika-server-core/src/test/java/org/apache/tika/server/core/TikaPipesTest.java
b/tika-server/tika-server-core/src/test/java/org/apache/tika/server/core/TikaPipesTest.java
index 88991250b2..8a6657358d 100644
---
a/tika-server/tika-server-core/src/test/java/org/apache/tika/server/core/TikaPipesTest.java
+++
b/tika-server/tika-server-core/src/test/java/org/apache/tika/server/core/TikaPipesTest.java
@@ -151,7 +151,7 @@ public class TikaPipesTest extends CXFTestBase {
TikaJsonConfig tikaJsonConfig =
TikaJsonConfig.load(tikaConfigPath);
PipesConfig pipesConfig = PipesConfig.load(tikaJsonConfig);
pipesConfig.setEmitStrategy(new
EmitStrategyConfig(EmitStrategy.EMIT_ALL));
- pipesParser = PipesParser.load(tikaJsonConfig, pipesConfig,
tikaConfigPath);
+ pipesParser = PipesParser.load(tikaJsonConfig, pipesConfig);
pipesResource = new PipesResource(pipesParser, pipesConfig);
rCoreProviders.add(new SingletonResourceProvider(pipesResource));
} catch (IOException | TikaConfigException e) {
diff --git
a/tika-server/tika-server-standard/src/test/java/org/apache/tika/server/standard/TikaPipesTest.java
b/tika-server/tika-server-standard/src/test/java/org/apache/tika/server/standard/TikaPipesTest.java
index d536ad12e3..a3cf4bb072 100644
---
a/tika-server/tika-server-standard/src/test/java/org/apache/tika/server/standard/TikaPipesTest.java
+++
b/tika-server/tika-server-standard/src/test/java/org/apache/tika/server/standard/TikaPipesTest.java
@@ -144,7 +144,7 @@ public class TikaPipesTest extends CXFTestBase {
TikaJsonConfig tikaJsonConfig =
TikaJsonConfig.load(tikaConfigPath);
PipesConfig pipesConfig = PipesConfig.load(tikaJsonConfig);
pipesConfig.setEmitStrategy(new
EmitStrategyConfig(EmitStrategy.EMIT_ALL));
- pipesParser = PipesParser.load(tikaJsonConfig, pipesConfig,
tikaConfigPath);
+ pipesParser = PipesParser.load(tikaJsonConfig, pipesConfig);
pipesResource = new PipesResource(pipesParser, pipesConfig);
rCoreProviders.add(new SingletonResourceProvider(pipesResource));
} catch (IOException | TikaConfigException e) {