Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,7 @@ public enum State {

private CommandLine commandLine;
private Map<String, String> envs;
private ExecuteWatchdog watchdog;
private ProcessWatchdog watchdog;
private ProcessLogOutputStream processOutput;
protected String errorMessage = null;
protected volatile State state = State.NEW;
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -161,12 +161,57 @@ 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;
}
}

/**
* 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();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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));
}
Expand Down Expand Up @@ -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);
}
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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;
Expand All @@ -163,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<String, String> envs) {
Expand Down Expand Up @@ -218,6 +244,24 @@ public void waitForReady(int timeout) {
}
}

public void cancelLaunch() {
synchronized (this) {
destroyProcessForcibly();
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();
// The process never reported that it is running, so stop() leaves it alive.
destroyProcessForcibly();
}

@Override
public void onProcessRunning() {
super.onProcessRunning();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -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<ManagedInterpreterGroup> 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<ManagedInterpreterGroup> 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<String, Object>(), new HashMap<String, Object>());
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;
}
}
}
Loading
Loading