This is an automated email from the ASF dual-hosted git repository. tballison pushed a commit to branch TIKA-4844 in repository https://gitbox.apache.org/repos/asf/tika.git
commit 6e236e737002a0e0d7300b8fc35ec7c1edda9684 Author: tallison <[email protected]> AuthorDate: Wed Aug 26 16:25:51 2026 -0400 TIKA-4844: scope shared-server restart reports to the process they describe --- CHANGES.txt | 18 ++++++ .../org/apache/tika/pipes/core/PipesClient.java | 22 +++++-- .../org/apache/tika/pipes/core/PipesParser.java | 9 +++ .../org/apache/tika/pipes/core/ServerManager.java | 23 +++++++ .../tika/pipes/core/SharedServerManager.java | 71 ++++++++++++++++++++-- .../pipes/core/SharedServerManagerClosedTest.java | 45 ++++++++++++++ .../tika/pipes/core/SharedServerModeTest.java | 18 ++++++ 7 files changed, 195 insertions(+), 11 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index d7907ba138..9b012ff0ed 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,5 +1,23 @@ Release 4.1.0 - unreleased + * Shared pipes server (useSharedServer: true, not the default): a client whose + in-flight parse was killed by another client's restart could restart the + healthy replacement. ensureRunning holds its lock across the whole fork, so + siblings cannot report a dead worker until after the replacement is up, and + the pending-restart flag carried no process identity -- so a report about + the process that just died was applied to its successor, which was then + destroyed and re-forked. One worker death produced two restarts and a second + round of destroyed in-flight work; under sustained concurrent load it + sustained itself at one spurious restart per round, appearing as periodic + unexplained worker churn and intermittent parse failures that succeed on + retry. Each fork now carries a generation that clients capture when they + connect and hand back with every report, and reports about a superseded + process are dropped. Also fixed in shared mode: ensureRunning could fork a + replacement after shutdown() that nothing owned and nothing would ever + destroy, and an interrupt during process teardown left the process handle + pointing at a killed process and leaked the temp directory. Affects 4.0.0 + and earlier (TIKA-4844). + * Documentation: corrected a batch of pages and javadocs that contradicted the code. Notably: the ES/OpenSearch attachmentStrategy has no default (unset means embedded documents get neither the parent field nor the diff --git a/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 7b2b34a974..c4743e1a65 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 @@ -82,6 +82,13 @@ public class PipesClient implements Closeable { // calling close(). The lock lets one thread atomically claim and null the tuple. private final Object connectionLock = new Object(); private volatile ConnectionTuple connectionTuple; + /** + * The server generation this client's connection belongs to, captured in {@link #reconnect()}. + * Starts at MAX_VALUE so a report made before we have ever connected is never mistaken for a + * stale one: dropping a legitimate report can wedge the pool, while an extra one only costs a + * restart, so the un-connected case must fail toward reporting. + */ + private volatile long connectionGeneration = Long.MAX_VALUE; private int filesProcessed = 0; /** @@ -325,6 +332,9 @@ public class PipesClient implements Closeable { // Get port after ensureRunning - this is the port we'll connect to int port = serverManager.getPort(); + // Captured with the port: every report we make below is about THIS process, and must be + // dropped if a sibling has already replaced it. + connectionGeneration = serverManager.getGeneration(); LOG.debug("pipesClientId={}: connecting to server", pipesClientId); // Connect to server. Use the generous startup timeout as the read SO_TIMEOUT so the @@ -404,7 +414,7 @@ public class PipesClient implements Closeable { LOG.warn("clientId={}: client-side backstop timeout: id={} elapsed={}ms limit={}ms " + "-- server should have self-terminated well before this", pipesClientId, t.getId(), totalElapsed, clientBackstopMillis); - serverManager.markServerForRestart(); + serverManager.markServerForRestart(connectionGeneration); closeConnection(); return buildFatalResult(t.getId(), t.getEmitKey(), TIMEOUT, intermediateResult.get()); } @@ -420,19 +430,19 @@ public class PipesClient implements Closeable { switch (msg.type()) { case OOM: String oomMsg = JsonPipesIpc.fromBytes(msg.payload(), String.class); - serverManager.markServerForRestart(); + serverManager.markServerForRestart(connectionGeneration); closeConnection(); return buildFatalResult(t.getId(), t.getEmitKey(), PipesResult.RESULT_STATUS.OOM, intermediateResult.get(), oomMsg); case TIMEOUT: String timeoutMsg = JsonPipesIpc.fromBytes(msg.payload(), String.class); - serverManager.markServerForRestart(); + serverManager.markServerForRestart(connectionGeneration); closeConnection(); return buildFatalResult(t.getId(), t.getEmitKey(), TIMEOUT, intermediateResult.get(), timeoutMsg); case UNSPECIFIED_CRASH: String crashMsg = JsonPipesIpc.fromBytes(msg.payload(), String.class); - serverManager.markServerForRestart(); + serverManager.markServerForRestart(connectionGeneration); closeConnection(); return buildFatalResult(t.getId(), t.getEmitKey(), UNSPECIFIED_CRASH, intermediateResult.get(), crashMsg); @@ -466,7 +476,7 @@ public class PipesClient implements Closeable { } catch (SocketTimeoutException e) { LOG.warn("clientId={}: Socket timeout exception while waiting for server", pipesClientId, e); // Mark for restart - server is stuck on current request and needs to be restarted - serverManager.markServerForRestart(); + serverManager.markServerForRestart(connectionGeneration); closeConnection(); return buildFatalResult(t.getId(), t.getEmitKey(), TIMEOUT, intermediateResult.get(), ExceptionUtils.getStackTrace(e)); @@ -482,7 +492,7 @@ public class PipesClient implements Closeable { } catch (Exception e) { LOG.warn("clientId={} - crash while waiting for server", pipesClientId, e); // Handle crash and determine status based on exit code - int exitCode = serverManager.handleCrashAndGetExitCode(); + int exitCode = serverManager.handleCrashAndGetExitCode(connectionGeneration); PipesResult.RESULT_STATUS status = UNSPECIFIED_CRASH; if (exitCode == PipesMessageType.OOM.getExitCode().orElse(-1)) { status = PipesResult.RESULT_STATUS.OOM; 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 dd68dca002..a03e24efe5 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 @@ -173,6 +173,15 @@ public class PipesParser implements Closeable { * * @return true if using shared server mode */ + /** + * Total forks performed across this parser's server managers. One more than the number of + * restarts, since the first start of each manager is not a restart. Exposed for tests: it is + * the only observable that distinguishes a justified restart from a spurious one. + */ + public long getGeneration() { + return serverManagers.stream().mapToLong(ServerManager::getGeneration).sum(); + } + public boolean isSharedMode() { return isSharedMode; } diff --git a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/ServerManager.java b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/ServerManager.java index 3a663193e6..b40d8c4757 100644 --- a/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/ServerManager.java +++ b/tika-pipes/tika-pipes-core/src/main/java/org/apache/tika/pipes/core/ServerManager.java @@ -163,4 +163,27 @@ public interface ServerManager extends Closeable { markServerForRestart(); return -1; } + + /** + * The generation of the currently running process: a counter incremented every time this + * manager forks a replacement. A client captures it when it connects and hands it back with + * every report, so a report about a process that has already been replaced can be recognised + * and dropped rather than being applied to its healthy successor. + */ + default long getGeneration() { + return 0; + } + + /** + * As {@link #markServerForRestart()}, but only if {@code generation} is still current. + * Reports about a superseded process are dropped. + */ + default void markServerForRestart(long generation) { + markServerForRestart(); + } + + /** As {@link #handleCrashAndGetExitCode()}, but only if {@code generation} is still current. */ + default int handleCrashAndGetExitCode(long generation) { + return handleCrashAndGetExitCode(); + } } 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 ea6db80039..93b4f03828 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 @@ -76,7 +76,9 @@ public class SharedServerManager implements ServerManager { private final Object lock = new Object(); private final AtomicLong filesProcessed = new AtomicLong(0); private volatile boolean restarting = false; - private volatile boolean pendingRestart = false; // Set when a client reports fatal error or max files reached + private volatile boolean pendingRestart = false; + private volatile long generation; + private volatile boolean closed; // Set when a client reports fatal error or max files reached private volatile byte[] currentToken; private Process process; private Path tmpDir; @@ -141,9 +143,13 @@ public class SharedServerManager implements ServerManager { if (isProcessAlive() && !pendingRestart) { return; } + if (closed) { + throw new IllegalStateException("shared server manager is closed"); + } restarting = true; try { startServer(); + generation++; // supersedes every report still in flight against the old process pendingRestart = false; // Clear the flag after successful restart filesProcessed.set(0); // Reset file counter on restart } finally { @@ -171,6 +177,51 @@ public class SharedServerManager implements ServerManager { } } + @Override + public long getGeneration() { + synchronized (lock) { + return generation; + } + } + + @Override + public void markServerForRestart(long generation) { + synchronized (lock) { + if (isSuperseded(generation)) { + return; + } + LOG.debug("Server marked for restart - will restart on next ensureRunning()"); + pendingRestart = true; + } + } + + @Override + public int handleCrashAndGetExitCode(long generation) { + synchronized (lock) { + if (isSuperseded(generation)) { + return -1; + } + pendingRestart = true; + } + return -1; + } + + /** + * A restart kills the shared JVM out from under every sibling still parsing on it, and + * {@code ensureRunning} holds {@code lock} for the whole fork, so those siblings cannot + * report until after the replacement is up. Their reports describe the process that was + * already destroyed; applying them would destroy the healthy replacement too. Callers + * hold {@code lock}. + */ + private boolean isSuperseded(long reportedGeneration) { + if (reportedGeneration >= generation) { + return false; + } + LOG.debug("dropping stale restart report for generation {}; current generation is {}", + reportedGeneration, generation); + return true; + } + /** * Increments the count of files processed and marks for restart if limit reached. */ @@ -378,6 +429,11 @@ public class SharedServerManager implements ServerManager { @Override public void shutdown() throws InterruptedException { synchronized (lock) { + // Latch here, not in shutdownUnsafe(): startServer() calls that to reap the previous + // process on every restart, so latching there would brick the manager on restart #1. + // A request thread racing us into ensureRunning must fail rather than fork a + // replacement that nothing owns and nothing will ever destroy. + closed = true; shutdownUnsafe(); } } @@ -398,11 +454,16 @@ public class SharedServerManager implements ServerManager { private void destroyProcessUnsafe() throws InterruptedException { if (process != null) { process.destroyForcibly(); - process.waitFor(WAIT_ON_DESTROY_MS, TimeUnit.MILLISECONDS); - if (process.isAlive()) { - LOG.error("Shared server process still alive after {}ms", WAIT_ON_DESTROY_MS); + try { + process.waitFor(WAIT_ON_DESTROY_MS, TimeUnit.MILLISECONDS); + if (process.isAlive()) { + LOG.error("Shared server process still alive after {}ms", WAIT_ON_DESTROY_MS); + } + } finally { + // An interrupt here must not leave the field pointing at a SIGKILLed process: + // startServer() would then try to reap it again and tmpDir would never be deleted. + process = null; } - process = null; } } diff --git a/tika-pipes/tika-pipes-core/src/test/java/org/apache/tika/pipes/core/SharedServerManagerClosedTest.java b/tika-pipes/tika-pipes-core/src/test/java/org/apache/tika/pipes/core/SharedServerManagerClosedTest.java new file mode 100644 index 0000000000..49cdb65492 --- /dev/null +++ b/tika-pipes/tika-pipes-core/src/test/java/org/apache/tika/pipes/core/SharedServerManagerClosedTest.java @@ -0,0 +1,45 @@ +/* + * 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.assertDoesNotThrow; +import static org.junit.jupiter.api.Assertions.assertThrows; + +import org.junit.jupiter.api.Test; + +/** + * A request thread can race a teardown into {@code ensureRunning}. Without a closed latch it + * forks a replacement that nothing owns and nothing will ever destroy; the child only exits when + * the parent JVM does. PerClientServerManager has guarded this since 4.0.0; shared mode did not. + */ +public class SharedServerManagerClosedTest { + + @Test + public void testEnsureRunningAfterShutdownDoesNotFork() throws Exception { + SharedServerManager sm = new SharedServerManager(new PipesConfig(), null, 2); + sm.shutdown(); + assertThrows(IllegalStateException.class, sm::ensureRunning, + "ensureRunning must refuse to fork once the manager has been shut down"); + } + + @Test + public void testShutdownIsIdempotent() throws Exception { + SharedServerManager sm = new SharedServerManager(new PipesConfig(), null, 2); + sm.shutdown(); + assertDoesNotThrow(sm::shutdown, "a second shutdown must be a no-op, not a failure"); + } +} diff --git a/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/SharedServerModeTest.java b/tika-pipes/tika-pipes-integration-tests/src/test/java/org/apache/tika/pipes/core/SharedServerModeTest.java index 59b8b4b3a5..02de36e279 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 @@ -617,6 +617,8 @@ public class SharedServerModeTest { assertTrue(verifyResult.isSuccess(), "After concurrent OOM, server should restart and process new request. Got: " + verifyResult.status()); + + assertOneForkPerDeath(pipesParser, 1); } } @@ -760,9 +762,25 @@ public class SharedServerModeTest { "All 5 phase 2 requests should succeed after server recovery. " + "Phase 1 had: " + phase1OomCount + " OOMs, " + phase1SuccessCount + " successes, " + phase1CrashCount + " crashes"); + + assertOneForkPerDeath(pipesParser, 1); } } + /** + * The shared JVM died {@code deaths} times. Siblings whose in-flight parses were killed by a + * death report a crash <em>result</em>, which is correct -- but none of them is another death, + * so none may cause another fork. Generation counts forks directly, so it cannot be satisfied + * by relabelling a spurious restart: one death, one replacement. + */ + private void assertOneForkPerDeath(PipesParser pipesParser, int deaths) { + long forks = pipesParser.getGeneration(); + assertEquals(deaths + 1, forks, + "expected the initial fork plus exactly " + deaths + " replacement(s); a sibling's " + + "report about an already-replaced process must not restart its healthy " + + "successor (forks=" + forks + ")"); + } + private Path setupInputDir(Path tmp) throws Exception { Path inputDir = tmp.resolve("input"); Files.createDirectories(inputDir);
