From 1aa489856fbb0736b59ad7ffc82fa634f462ee29 Mon Sep 17 00:00:00 2001 From: dev-donghwan Date: Wed, 30 Sep 2026 16:49:08 +0900 Subject: [PATCH 1/4] [ZEPPELIN-6721] Keep a group created while another group with the same id is closing --- .../interpreter/InterpreterSetting.java | 18 ++- .../interpreter/ManagedInterpreterGroup.java | 2 +- .../interpreter/InterpreterSettingTest.java | 129 ++++++++++++++++++ 3 files changed, 147 insertions(+), 2 deletions(-) diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java index f26bc54f0e7..a26593e50b5 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java @@ -476,6 +476,22 @@ void removeInterpreterGroup(String groupId) { } } + /** + * Removes the given group only if it is still the one registered under its id. A group with the + * same id can be created while this one is closing, and that group must stay registered. + * Compares references, because {@link InterpreterGroup#equals} compares ids only. + */ + void removeInterpreterGroup(ManagedInterpreterGroup interpreterGroup) { + try { + interpreterGroupWriteLock.lock(); + if (this.interpreterGroups.get(interpreterGroup.getId()) == interpreterGroup) { + this.interpreterGroups.remove(interpreterGroup.getId()); + } + } finally { + interpreterGroupWriteLock.unlock(); + } + } + public ManagedInterpreterGroup getInterpreterGroup(String user, String noteId) { return getInterpreterGroup(getExecutionContext(user, noteId)); } @@ -534,7 +550,7 @@ public void closeInterpreters(ExecutionContext executionContext) { String sessionId = getInterpreterSessionId(executionContext); interpreterGroup.close(sessionId); if (interpreterGroup.isEmpty()) { - interpreterGroups.remove(interpreterGroup.getId()); + removeInterpreterGroup(interpreterGroup); } } } diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/ManagedInterpreterGroup.java b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/ManagedInterpreterGroup.java index 3a8b14ee81e..539e085e88f 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/ManagedInterpreterGroup.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/ManagedInterpreterGroup.java @@ -150,7 +150,7 @@ public void close(String sessionId) { synchronized (this) { if (sessions.isEmpty() && interpreterSetting != null) { LOGGER.info("Remove this InterpreterGroup: {} as all the sessions are closed", id); - interpreterSetting.removeInterpreterGroup(id); + interpreterSetting.removeInterpreterGroup(this); if (remoteInterpreterProcess != null) { LOGGER.info("Kill RemoteInterpreterProcess"); remoteInterpreterProcess.stop(); diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/InterpreterSettingTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/InterpreterSettingTest.java index 7d4051687bf..c9d27b1afb4 100644 --- a/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/InterpreterSettingTest.java +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/InterpreterSettingTest.java @@ -20,21 +20,31 @@ import com.google.common.collect.Lists; import org.apache.zeppelin.dep.Dependency; import org.apache.zeppelin.dep.DependencyResolver; +import org.apache.zeppelin.interpreter.remote.RemoteInterpreterProcess; import org.apache.zeppelin.user.AuthenticationInfo; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.lang.reflect.Field; import java.util.ArrayList; import java.util.HashMap; import java.util.List; +import java.util.Properties; +import java.util.concurrent.CyclicBarrier; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; 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.assertNotSame; import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertSame; import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.when; @@ -600,4 +610,123 @@ void testLoadDependency() throws InterruptedException { assertNull(interpreterSetting.getErrorReason()); } + + @Test + void testCloseInterpretersKeepsGroupCreatedWhileProcessStops() throws Exception { + InterpreterSetting interpreterSetting = createEchoInterpreterSetting(InterpreterOption.SHARED); + interpreterSetting.getDefaultInterpreter("user1", note1Id); + ManagedInterpreterGroup closingGroup = + interpreterSetting.getInterpreterGroup("user1", note1Id); + + // A paragraph that runs while the old process is stopping creates a new group with the same id. + AtomicReference newGroup = new AtomicReference<>(); + setInterpreterProcess(closingGroup, processThatRunsOnStop(() -> + newGroup.set(interpreterSetting.getOrCreateInterpreterGroup("user1", note2Id)))); + + interpreterSetting.closeInterpreters("user1", note1Id); + + assertNotSame(closingGroup, newGroup.get()); + assertSame(newGroup.get(), interpreterSetting.getInterpreterGroup("user1", note2Id)); + } + + @Test + @Timeout(30) + void testConcurrentCloseOfLastSessionsKeepsGroupCreatedWhileProcessStops() throws Exception { + InterpreterSetting interpreterSetting = createEchoInterpreterSetting(InterpreterOption.SCOPED); + ManagedInterpreterGroup closingGroup = + interpreterSetting.getOrCreateInterpreterGroup("user1", note1Id); + + // Both closers remove their session before either of them checks whether the group is empty. + CyclicBarrier sessionsRemoved = new CyclicBarrier(2); + closingGroup.sessions.put(note1Id, Lists.newArrayList(new BarrierInterpreter(sessionsRemoved))); + closingGroup.sessions.put(note2Id, Lists.newArrayList(new BarrierInterpreter(sessionsRemoved))); + + AtomicReference newGroup = new AtomicReference<>(); + setInterpreterProcess(closingGroup, processThatRunsOnStop(() -> + newGroup.set(interpreterSetting.getOrCreateInterpreterGroup("user1", note1Id)))); + + Thread closer1 = new Thread(() -> interpreterSetting.closeInterpreters("user1", note1Id)); + Thread closer2 = new Thread(() -> interpreterSetting.closeInterpreters("user1", note2Id)); + closer1.start(); + closer2.start(); + closer1.join(); + closer2.join(); + + assertNotSame(closingGroup, newGroup.get()); + assertSame(newGroup.get(), interpreterSetting.getInterpreterGroup("user1", note1Id)); + } + + private InterpreterSetting createEchoInterpreterSetting(String perNote) { + InterpreterOption interpreterOption = new InterpreterOption(); + interpreterOption.setPerNote(perNote); + InterpreterInfo interpreterInfo = new InterpreterInfo(EchoInterpreter.class.getName(), + "echo", true, new HashMap(), new HashMap()); + return new InterpreterSetting.Builder() + .setId("id") + .setName("test") + .setGroup("test") + .setInterpreterInfos(Lists.newArrayList(interpreterInfo)) + .setOption(interpreterOption) + .setIntepreterSettingManager(interpreterSettingManager) + .setConf(zConf) + .create(); + } + + private static RemoteInterpreterProcess processThatRunsOnStop(Runnable onStop) { + RemoteInterpreterProcess process = mock(RemoteInterpreterProcess.class); + doAnswer(invocation -> { + onStop.run(); + return null; + }).when(process).stop(); + return process; + } + + private static void setInterpreterProcess(ManagedInterpreterGroup interpreterGroup, + RemoteInterpreterProcess process) throws Exception { + Field field = ManagedInterpreterGroup.class.getDeclaredField("remoteInterpreterProcess"); + field.setAccessible(true); + field.set(interpreterGroup, process); + } + + private static class BarrierInterpreter extends Interpreter { + + private final CyclicBarrier barrier; + + BarrierInterpreter(CyclicBarrier barrier) { + super(new Properties()); + this.barrier = barrier; + } + + @Override + public void close() { + try { + barrier.await(10, TimeUnit.SECONDS); + } catch (Exception e) { + throw new RuntimeException(e); + } + } + + @Override + public void open() { + } + + @Override + public InterpreterResult interpret(String st, InterpreterContext context) { + return null; + } + + @Override + public void cancel(InterpreterContext context) { + } + + @Override + public FormType getFormType() { + return FormType.NATIVE; + } + + @Override + public int getProgress(InterpreterContext context) { + return 0; + } + } } From e75c0d07d00bcda443b67323dcf83ba942c06e50 Mon Sep 17 00:00:00 2001 From: dev-donghwan Date: Wed, 30 Sep 2026 16:54:12 +0900 Subject: [PATCH 2/4] [ZEPPELIN-6721] Destroy the interpreter process when its launch times out --- .../interpreter/util/ProcessLauncher.java | 12 ++- .../remote/ExecRemoteInterpreterProcess.java | 7 ++ .../ExecRemoteInterpreterProcessTest.java | 96 +++++++++++++++++++ 3 files changed, 114 insertions(+), 1 deletion(-) create mode 100644 zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcessTest.java diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/util/ProcessLauncher.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/util/ProcessLauncher.java index 266094e18b6..997fbcfcb9c 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/util/ProcessLauncher.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/util/ProcessLauncher.java @@ -161,7 +161,17 @@ public boolean isRunning() { } public void stop() { - if (watchdog != null && isRunning()) { + if (isRunning()) { + destroyProcess(); + } + } + + /** + * Destroys the process whatever its state, e.g. for a launch that is given up before the + * process reports that it is running. + */ + protected void destroyProcess() { + if (watchdog != null) { watchdog.destroyProcess(); watchdog = null; } diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcess.java b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcess.java index 6dd2793adf8..52a7fcc541c 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcess.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcess.java @@ -218,6 +218,13 @@ public void waitForReady(int timeout) { } } + @Override + public void onTimeout() { + super.onTimeout(); + // The process never reported that it is running, so stop() leaves it alive. + destroyProcess(); + } + @Override public void onProcessRunning() { super.onProcessRunning(); diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcessTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcessTest.java new file mode 100644 index 00000000000..c2f0e1bc7ca --- /dev/null +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcessTest.java @@ -0,0 +1,96 @@ +/* + * 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.zeppelin.interpreter.remote; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertThrows; + +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.HashMap; +import java.util.concurrent.TimeUnit; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; +import org.junit.jupiter.api.condition.DisabledOnOs; +import org.junit.jupiter.api.condition.OS; +import org.junit.jupiter.api.io.TempDir; + +/** + * Launches a stand-in interpreter process that never registers with the server, so the launch + * stays in LAUNCHED until it times out or is cancelled. + */ +@DisabledOnOs(OS.WINDOWS) +class ExecRemoteInterpreterProcessTest { + + @TempDir + Path tempDir; + + @Test + @Timeout(60) + void launchTimeoutDestroysTheProcess() throws Exception { + Path pidFile = tempDir.resolve("pid"); + ExecRemoteInterpreterProcess process = createProcess(neverRegisteringRunner(pidFile), 3000); + + assertThrows(IOException.class, () -> process.start("anonymous")); + + assertExits(readPid(pidFile)); + } + + private Path neverRegisteringRunner(Path pidFile) throws IOException { + Path runner = tempDir.resolve("interpreter.sh"); + Files.write(runner, ("#!/bin/sh\n" + + "echo $$ > " + pidFile + "\n" + + "exec sleep 600\n").getBytes(StandardCharsets.UTF_8)); + runner.toFile().setExecutable(true); + return runner; + } + + private ExecRemoteInterpreterProcess createProcess(Path runner, int connectTimeout) { + return new ExecRemoteInterpreterProcess( + 0, "127.0.0.1", ":", tempDir.toString(), tempDir.toString(), new HashMap<>(), + connectTimeout, 10, "test", "test-shared_process", false, runner.toString()); + } + + private static long readPid(Path pidFile) throws Exception { + long deadline = System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(10); + while (System.currentTimeMillis() < deadline) { + if (Files.exists(pidFile)) { + String pid = new String(Files.readAllBytes(pidFile), StandardCharsets.UTF_8).trim(); + if (!pid.isEmpty()) { + return Long.parseLong(pid); + } + } + Thread.sleep(50); + } + throw new IllegalStateException("the runner did not write its pid"); + } + + private static void assertExits(long pid) throws InterruptedException { + long deadline = System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(10); + while (isAlive(pid) && System.currentTimeMillis() < deadline) { + Thread.sleep(50); + } + assertFalse(isAlive(pid), "the launched process " + pid + " is still running"); + } + + private static boolean isAlive(long pid) { + return ProcessHandle.of(pid).map(ProcessHandle::isAlive).orElse(false); + } +} From 36af106a28f85365acfce32b5dfa7f46689933e8 Mon Sep 17 00:00:00 2001 From: dev-donghwan Date: Wed, 30 Sep 2026 18:10:57 +0900 Subject: [PATCH 3/4] [ZEPPELIN-6721] End the launch and destroy the process when a launching interpreter process is stopped --- .../remote/ExecRemoteInterpreterProcess.java | 35 +++++++++++++++++-- .../ExecRemoteInterpreterProcessTest.java | 29 +++++++++++++++ 2 files changed, 62 insertions(+), 2 deletions(-) diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcess.java b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcess.java index 52a7fcc541c..c5768aa7b3c 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcess.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcess.java @@ -40,6 +40,8 @@ public class ExecRemoteInterpreterProcess extends RemoteInterpreterManagedProces private final String interpreterRunner; private InterpreterProcessLauncher interpreterProcessLauncher; + // Guarded by this. Set once stop() is called before the process reports that it is running. + private boolean launchCancelled; public ExecRemoteInterpreterProcess( int intpEventServerPort, @@ -82,8 +84,14 @@ public void start(String userName) throws IOException { cmdLine.addArgument("-g", false); cmdLine.addArgument(getInterpreterSettingName(), false); - interpreterProcessLauncher = new InterpreterProcessLauncher(cmdLine, getEnv()); - interpreterProcessLauncher.launch(); + synchronized (this) { + if (launchCancelled) { + throw new IOException("Interpreter process of interpreter group " + getInterpreterGroupId() + + " is stopped before it is launched"); + } + interpreterProcessLauncher = new InterpreterProcessLauncher(cmdLine, getEnv()); + interpreterProcessLauncher.launch(); + } interpreterProcessLauncher.waitForReady(getConnectTimeout()); if (interpreterProcessLauncher.isLaunchTimeout()) { throw new IOException( @@ -137,10 +145,22 @@ public void stop() { } else { // Shutdown connection super.close(); + cancelLaunch(); LOGGER.warn("Try to stop a not running interpreter process of interpreter group: {}", getInterpreterGroupId()); } } + /** + * A process that is still launching has not registered yet, so the server cannot ask it to shut + * down. Destroy it and end the launch, instead of leaving start() to wait for the connect timeout. + */ + private synchronized void cancelLaunch() { + launchCancelled = true; + if (interpreterProcessLauncher != null) { + interpreterProcessLauncher.cancelLaunch(); + } + } + @VisibleForTesting public String getInterpreterRunner() { return interpreterRunner; @@ -218,6 +238,17 @@ public void waitForReady(int timeout) { } } + public void cancelLaunch() { + synchronized (this) { + destroyProcess(); + if (state == State.LAUNCHED) { + errorMessage = "The launch is cancelled, because the interpreter process is stopped"; + transition(State.TERMINATED); + } + notifyAll(); + } + } + @Override public void onTimeout() { super.onTimeout(); diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcessTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcessTest.java index c2f0e1bc7ca..ffb26294892 100644 --- a/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcessTest.java +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcessTest.java @@ -18,6 +18,7 @@ package org.apache.zeppelin.interpreter.remote; import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; import static org.junit.jupiter.api.Assertions.assertThrows; import java.io.IOException; @@ -25,6 +26,9 @@ import java.nio.file.Files; import java.nio.file.Path; import java.util.HashMap; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; +import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Timeout; @@ -53,6 +57,31 @@ void launchTimeoutDestroysTheProcess() throws Exception { assertExits(readPid(pidFile)); } + @Test + @Timeout(60) + void stopWhileLaunchingEndsTheLaunchAndDestroysTheProcess() throws Exception { + Path pidFile = tempDir.resolve("pid"); + ExecRemoteInterpreterProcess process = + createProcess(neverRegisteringRunner(pidFile), (int) TimeUnit.MINUTES.toMillis(10)); + + CompletableFuture launch = CompletableFuture.runAsync(() -> { + try { + process.start("anonymous"); + } catch (IOException e) { + throw new CompletionException(e); + } + }); + long pid = readPid(pidFile); + + // What ManagedInterpreterGroup.close() does when the group is closed during the launch. + process.stop(); + + ExecutionException launchFailure = + assertThrows(ExecutionException.class, () -> launch.get(10, TimeUnit.SECONDS)); + assertTrue(launchFailure.getCause() instanceof IOException, launchFailure.toString()); + assertExits(pid); + } + private Path neverRegisteringRunner(Path pidFile) throws IOException { Path runner = tempDir.resolve("interpreter.sh"); Files.write(runner, ("#!/bin/sh\n" From 2b4f7152b6236239c2308e1a7dc68d87cf9f43ad Mon Sep 17 00:00:00 2001 From: dev-donghwan Date: Wed, 30 Sep 2026 18:30:44 +0900 Subject: [PATCH 4/4] [ZEPPELIN-6721] Kill an abandoned launch forcibly so that its shutdown hook does not unregister the group id --- .../interpreter/util/ProcessLauncher.java | 39 ++++++++++++++++++- .../remote/ExecRemoteInterpreterProcess.java | 10 ++++- .../ExecRemoteInterpreterProcessTest.java | 13 ++++++- 3 files changed, 56 insertions(+), 6 deletions(-) diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/util/ProcessLauncher.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/util/ProcessLauncher.java index 997fbcfcb9c..9b647eccb61 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/util/ProcessLauncher.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/util/ProcessLauncher.java @@ -50,7 +50,7 @@ public enum State { private CommandLine commandLine; private Map envs; - private ExecuteWatchdog watchdog; + private ProcessWatchdog watchdog; private ProcessLogOutputStream processOutput; protected String errorMessage = null; protected volatile State state = State.NEW; @@ -87,7 +87,7 @@ public void setRedirectedContext(InterpreterContext redirectedContext) { public void launch() { DefaultExecutor executor = new DefaultExecutor(); executor.setStreamHandler(new PumpStreamHandler(processOutput)); - this.watchdog = new ExecuteWatchdog(ExecuteWatchdog.INFINITE_TIMEOUT); + this.watchdog = new ProcessWatchdog(); executor.setWatchdog(watchdog); try { executor.execute(commandLine, envs, this); @@ -177,6 +177,41 @@ protected void destroyProcess() { } } + /** + * Like {@link #destroyProcess()}, but kills the process without running its shutdown hooks. + */ + protected void destroyProcessForcibly() { + if (watchdog != null) { + watchdog.destroyProcessForcibly(); + watchdog = null; + } + } + + /** + * ExecuteWatchdog only destroys the process with Process.destroy(), so keep the process to be + * able to kill it forcibly. + */ + private static class ProcessWatchdog extends ExecuteWatchdog { + + private Process process; + + ProcessWatchdog() { + super(ExecuteWatchdog.INFINITE_TIMEOUT); + } + + @Override + public synchronized void start(Process process) { + this.process = process; + super.start(process); + } + + synchronized void destroyProcessForcibly() { + if (process != null) { + process.destroyForcibly(); + } + } + } + public void stopCatchLaunchOutput() { processOutput.stopCatchLaunchOutput(); } diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcess.java b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcess.java index c5768aa7b3c..36b2f0f224b 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcess.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcess.java @@ -183,6 +183,12 @@ public String getErrorMessage() { : ""; } + /** + * A launch that is given up, on timeout or when the process is stopped while launching, kills + * the process forcibly. It has not been initialized, so it has no open interpreter to close, and + * its shutdown hook would unregister its interpreter group id, which by then can belong to + * another group (ZEPPELIN-6723). + */ private class InterpreterProcessLauncher extends ProcessLauncher { public InterpreterProcessLauncher(CommandLine commandLine, Map envs) { @@ -240,7 +246,7 @@ public void waitForReady(int timeout) { public void cancelLaunch() { synchronized (this) { - destroyProcess(); + destroyProcessForcibly(); if (state == State.LAUNCHED) { errorMessage = "The launch is cancelled, because the interpreter process is stopped"; transition(State.TERMINATED); @@ -253,7 +259,7 @@ public void cancelLaunch() { public void onTimeout() { super.onTimeout(); // The process never reported that it is running, so stop() leaves it alive. - destroyProcess(); + destroyProcessForcibly(); } @Override diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcessTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcessTest.java index ffb26294892..488dd057eac 100644 --- a/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcessTest.java +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/remote/ExecRemoteInterpreterProcessTest.java @@ -38,7 +38,8 @@ /** * Launches a stand-in interpreter process that never registers with the server, so the launch - * stays in LAUNCHED until it times out or is cancelled. + * stays in LAUNCHED until it times out or is cancelled. It records a SIGTERM, which in a real + * interpreter process would run the shutdown hook that unregisters its interpreter group id. */ @DisabledOnOs(OS.WINDOWS) class ExecRemoteInterpreterProcessTest { @@ -55,6 +56,7 @@ void launchTimeoutDestroysTheProcess() throws Exception { assertThrows(IOException.class, () -> process.start("anonymous")); assertExits(readPid(pidFile)); + assertNoShutdownHookRan(); } @Test @@ -80,13 +82,15 @@ void stopWhileLaunchingEndsTheLaunchAndDestroysTheProcess() throws Exception { assertThrows(ExecutionException.class, () -> launch.get(10, TimeUnit.SECONDS)); assertTrue(launchFailure.getCause() instanceof IOException, launchFailure.toString()); assertExits(pid); + assertNoShutdownHookRan(); } private Path neverRegisteringRunner(Path pidFile) throws IOException { Path runner = tempDir.resolve("interpreter.sh"); Files.write(runner, ("#!/bin/sh\n" + + "trap 'echo TERM > " + tempDir.resolve("sigterm") + "; exit 143' TERM\n" + "echo $$ > " + pidFile + "\n" - + "exec sleep 600\n").getBytes(StandardCharsets.UTF_8)); + + "while true; do sleep 1; done\n").getBytes(StandardCharsets.UTF_8)); runner.toFile().setExecutable(true); return runner; } @@ -119,6 +123,11 @@ private static void assertExits(long pid) throws InterruptedException { assertFalse(isAlive(pid), "the launched process " + pid + " is still running"); } + private void assertNoShutdownHookRan() { + assertFalse(Files.exists(tempDir.resolve("sigterm")), + "the process received SIGTERM, so a real interpreter process would run its shutdown hook"); + } + private static boolean isAlive(long pid) { return ProcessHandle.of(pid).map(ProcessHandle::isAlive).orElse(false); }