diff --git a/python/src/com/jetbrains/python/console/PydevConsoleCommunication.java b/python/src/com/jetbrains/python/console/PydevConsoleCommunication.java index 0040eac1d64a..7767701ea425 100644 --- a/python/src/com/jetbrains/python/console/PydevConsoleCommunication.java +++ b/python/src/com/jetbrains/python/console/PydevConsoleCommunication.java @@ -1,6 +1,7 @@ // Copyright 2000-2018 JetBrains s.r.o. Use of this source code is governed by the Apache 2.0 license that can be found in the LICENSE file. package com.jetbrains.python.console; +import com.google.common.util.concurrent.SettableFuture; import com.intellij.openapi.application.ApplicationManager; import com.intellij.openapi.diagnostic.Logger; import com.intellij.openapi.fileEditor.FileEditorManager; @@ -41,10 +42,7 @@ import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.Map; -import java.util.concurrent.CompletableFuture; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.Future; -import java.util.concurrent.TimeUnit; +import java.util.concurrent.*; import java.util.stream.Collectors; import static com.jetbrains.python.console.PydevConsoleCommunicationUtil.*; @@ -98,6 +96,17 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl @Nullable private XCompositeNode myCurrentRootNode; + @NotNull + private final SettableFuture myInitialPythonConsoleClientFuture = SettableFuture.create(); + + @NotNull + private final SettableFuture myPythonConsoleProcessFuture = SettableFuture.create(); + + /** + * Indicates that {@link #close()} was executed. + */ + private boolean myClose = false; + /** * Initializes the bidirectional RPC communication. */ @@ -121,7 +130,7 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl PythonConsoleBackendService.Client client = new PythonConsoleBackendService.Client(clientProtocol); this.myServer = server; - this.myClient = PythonConsoleClientUtil.synchronizedPythonConsoleClient(PydevConsoleCommunication.class.getClassLoader(), client); + this.myInitialPythonConsoleClientFuture.set(client); PyDebugValueExecutionService executionService = PyDebugValueExecutionService.getInstance(myProject); executionService.sessionStarted(this); @@ -155,8 +164,8 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl ApplicationManager.getApplication().executeOnPooledThread(() -> server.serve()); this.myServer = server; - this.myClient = PythonConsoleClientUtil.synchronizedPythonConsoleClient(PydevConsoleCommunication.class.getClassLoader(), client, - pythonConsoleProcess); + this.myInitialPythonConsoleClientFuture.set(client); + this.myPythonConsoleProcessFuture.set(pythonConsoleProcess); PyDebugValueExecutionService executionService = PyDebugValueExecutionService.getInstance(myProject); executionService.sessionStarted(this); @@ -169,10 +178,58 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl }); } + @NotNull + private Process getPythonConsoleProcess() { + try { + return myPythonConsoleProcessFuture.get(); + } + catch (InterruptedException | ExecutionException e) { + throw new IllegalStateException(e); + } + } + + public void setPythonConsoleProcess(@NotNull Process pythonConsoleProcess) { + myPythonConsoleProcessFuture.set(pythonConsoleProcess); + } + + /** + * Returns initial non-thread safe {@link PythonConsoleBackendService.Client}. + * + * @return {@link PythonConsoleBackendService.Client} + */ + @NotNull + private PythonConsoleBackendService.Client getInitialPythonConsoleBackendClient() { + try { + return myInitialPythonConsoleClientFuture.get(); + } + catch (InterruptedException | ExecutionException e) { + throw new IllegalStateException(e); + } + } + + /** + * Returns thread safe {@link PythonConsoleBackendService.Iface}: requests to + * the returned {@link PythonConsoleBackendService.Iface} will be processed + * sequentially. Also if Python Console process is detected to be finished + * the current request will be interrupted with + * {@link PyConsoleProcessFinishedException}. + * + * @return thread safe and related Python Console process-aware + * {@link PythonConsoleBackendService.Iface} + */ + @NotNull + private PythonConsoleBackendService.Iface getPythonConsoleBackendClient() { + if (myClient == null) { + myClient = PythonConsoleClientUtil.synchronizedPythonConsoleClient(PydevConsoleCommunication.class.getClassLoader(), + getInitialPythonConsoleBackendClient(), getPythonConsoleProcess()); + } + return myClient; + } + public boolean handshake() { - if (myClient != null) { + if (!myClose) { try { - return "PyCharm".equals(myClient.handshake()); + return "PyCharm".equals(getPythonConsoleBackendClient().handshake()); } catch (PyConsoleProcessFinishedException | TException e) { throw new RuntimeException(e); @@ -182,17 +239,17 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl } private void sendCloseMessageToScript() { - if (this.myClient != null) { + if (!myClose) { + myClose = true; new Task.Backgroundable(myProject, "Close Console Communication", true) { @Override public void run(@NotNull ProgressIndicator indicator) { try { - PydevConsoleCommunication.this.myClient.close(); + getPythonConsoleBackendClient().close(); } catch (Exception e) { //Ok, we can ignore this one on close. } - PydevConsoleCommunication.this.myClient = null; } }.queue(); } @@ -201,15 +258,45 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl /** * Stops the communication with the client (passes message for it to quit). */ - public synchronized void close() { - sendCloseMessageToScript(); + public void close() { + if (myClose) { + return; + } + myClose = true; + PyDebugValueExecutionService.getInstance(myProject).sessionStopped(this); myCallbackHashMap.clear(); - if (myServer != null) { - myServer.stop(); - myServer = null; - } + new Task.Backgroundable(myProject, "Close Console Communication", true) { + @Override + public void run(@NotNull ProgressIndicator indicator) { + try { + indicator.setText2("Sending close message to Python Console..."); + try { + getPythonConsoleBackendClient().close(); + } + catch (Exception e) { + //Ok, we can ignore this one on close. + } + indicator.setText2("Waiting for Python Console process to finish..."); + try { + do { + indicator.checkCanceled(); + } + while (!getPythonConsoleProcess().waitFor(500, TimeUnit.MILLISECONDS)); + } + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } + finally { + if (myServer != null) { + myServer.stop(); + myServer = null; + } + } + } + }.queue(); } /** @@ -219,16 +306,25 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl * thread {@link WebServer#listener} to die */ @NotNull - public synchronized Future closeAsync() { + public Future closeAsync() { sendCloseMessageToScript(); PyDebugValueExecutionService.getInstance(myProject).sessionStopped(this); myCallbackHashMap.clear(); - - if (myServer != null) { - return myServer.stop(); - } - - return completedFuture(); + return ApplicationManager.getApplication().executeOnPooledThread(() -> { + try { + getPythonConsoleProcess().waitFor(5L, TimeUnit.SECONDS); + } + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + finally { + if (myServer != null) { + myServer.stop(); + myServer = null; + } + } + return null; + }); } /** @@ -309,10 +405,10 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl boolean more; try { if (command.isSingleLine()) { - more = myClient.execLine(command.getText()); + more = getPythonConsoleBackendClient().execLine(command.getText()); } else { - more = myClient.execMultipleLines(command.getText()); + more = getPythonConsoleBackendClient().execMultipleLines(command.getText()); } } catch (TException e) { @@ -339,7 +435,7 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl if (waitingForInput) { return Collections.emptyList(); } - List fromServer = myClient.getCompletions(text, actTok); + List fromServer = getPythonConsoleBackendClient().getCompletions(text, actTok); return fromServer.stream().map(option -> toPydevCompletionVariant(option)).collect(Collectors.toList()); } @@ -362,7 +458,7 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl return "Unable to get description: waiting for input."; } - ThrowableComputable doGetDesc = () -> myClient.getDescription(text); + ThrowableComputable doGetDesc = () -> getPythonConsoleBackendClient().getDescription(text); if (ApplicationManager.getApplication().isDispatchThread()) { return ProgressManager.getInstance().runProcessWithProgressSynchronously(doGetDesc, "Getting Description", true, myProject); } @@ -499,7 +595,7 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl @Override public void interrupt() { try { - myClient.interrupt(); + getPythonConsoleBackendClient().interrupt(); } catch (PyConsoleProcessFinishedException | TException e) { LOG.error(e); @@ -518,9 +614,9 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl @Override public PyDebugValue evaluate(String expression, boolean execute, boolean doTrunc) throws PyDebuggerException { - if (myClient != null) { + if (!myClose) { try { - List debugValues = myClient.evaluate(expression); + List debugValues = getPythonConsoleBackendClient().evaluate(expression); return createPyDebugValue(debugValues.iterator().next(), this); } catch (Exception e) { @@ -533,9 +629,9 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl @Nullable @Override public XValueChildrenList loadFrame() throws PyDebuggerException { - if (myClient != null) { + if (!myClose) { try { - List frame = myClient.getFrame(); + List frame = getPythonConsoleBackendClient().getFrame(); return parseVars(frame, null, this); } catch (PyConsoleProcessFinishedException | TException e) { @@ -553,30 +649,28 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl @Override public void loadAsyncVariablesValues(@NotNull List> pyAsyncValues) { PyDebugValueExecutionService.getInstance(myProject).submitTask(this, () -> { - if (myClient != null) { - try { - List evaluationExpressions = new ArrayList<>(); - for (PyAsyncValue asyncValue : pyAsyncValues) { - evaluationExpressions.add(GetVariableCommand.composeName(asyncValue.getDebugValue())); - } - final int seq = getNextFullValueSeq(); - myCallbackHashMap.put(seq, pyAsyncValues); - - myClient.loadFullValue(seq, evaluationExpressions); - - // previously `loadFullValue()` might return `List` but this is no longer true + try { + List evaluationExpressions = new ArrayList<>(); + for (PyAsyncValue asyncValue : pyAsyncValues) { + evaluationExpressions.add(GetVariableCommand.composeName(asyncValue.getDebugValue())); } - catch (PyConsoleProcessFinishedException | TException e) { - for (PyAsyncValue asyncValue : pyAsyncValues) { - PyDebugValue value = asyncValue.getDebugValue(); - XValueNode node = value.getLastNode(); - if (node != null && !node.isObsolete()) { - if (e.getMessage().startsWith("Timeout") || e.getMessage().startsWith("Console already exited")) { - value.updateNodeValueAfterLoading(node, " ", "", PyVariableViewSettings.LOADING_TIMED_OUT); - } - else { - LOG.error(e); - } + final int seq = getNextFullValueSeq(); + myCallbackHashMap.put(seq, pyAsyncValues); + + getPythonConsoleBackendClient().loadFullValue(seq, evaluationExpressions); + + // previously `loadFullValue()` might return `List` but this is no longer true + } + catch (PyConsoleProcessFinishedException | TException e) { + for (PyAsyncValue asyncValue : pyAsyncValues) { + PyDebugValue value = asyncValue.getDebugValue(); + XValueNode node = value.getLastNode(); + if (node != null && !node.isObsolete()) { + if (e.getMessage().startsWith("Timeout") || e.getMessage().startsWith("Console already exited")) { + value.updateNodeValueAfterLoading(node, " ", "", PyVariableViewSettings.LOADING_TIMED_OUT); + } + else { + LOG.error(e); } } } @@ -586,9 +680,9 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl @Override public XValueChildrenList loadVariable(PyDebugValue var) throws PyDebuggerException { - if (myClient != null) { + if (!myClose) { try { - List ret = myClient.getVariable(GetVariableCommand.composeName(var)); + List ret = getPythonConsoleBackendClient().getVariable(GetVariableCommand.composeName(var)); return parseVars(ret, var, this); } catch (PyConsoleProcessFinishedException e) { @@ -614,11 +708,11 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl @Override public void changeVariable(PyDebugValue variable, String value) throws PyDebuggerException { - if (myClient != null) { + if (!myClose) { try { // NOTE: The actual change is being scheduled in the exec_queue in main thread // This method is async now - myClient.changeVariable(variable.getEvaluationExpression(), value); + getPythonConsoleBackendClient().changeVariable(variable.getEvaluationExpression(), value); } catch (PyConsoleProcessFinishedException | TException e) { throw new PyDebuggerException("Get change variable", e); @@ -635,9 +729,9 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl @Override public ArrayChunk getArrayItems(PyDebugValue var, int rowOffset, int colOffset, int rows, int cols, String format) throws PyDebuggerException { - if (myClient != null) { + if (!myClose) { try { - GetArrayResponse ret = myClient.getArray(var.getName(), rowOffset, colOffset, rows, cols, format); + GetArrayResponse ret = getPythonConsoleBackendClient().getArray(var.getName(), rowOffset, colOffset, rows, cols, format); return createArrayChunk(ret, this); } catch (Exception e) { @@ -674,7 +768,7 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl } try { // though `connectToDebugger` returns "connect complete" string, let us just ignore it - myClient.connectToDebugger(localPort, dbgOpts, extraEnvs); + getPythonConsoleBackendClient().connectToDebugger(localPort, dbgOpts, extraEnvs); } catch (PyConsoleProcessFinishedException | TException e) { throw new PyDebuggerException("pydevconsole failed to execute connectToDebugger", e); diff --git a/python/src/com/jetbrains/python/console/PydevConsoleRunnerImpl.java b/python/src/com/jetbrains/python/console/PydevConsoleRunnerImpl.java index 44701e2064f1..0dd972556248 100644 --- a/python/src/com/jetbrains/python/console/PydevConsoleRunnerImpl.java +++ b/python/src/com/jetbrains/python/console/PydevConsoleRunnerImpl.java @@ -413,7 +413,9 @@ public class PydevConsoleRunnerImpl implements PydevConsoleRunner { throw new ExecutionException(e.getMessage(), e); } - return new CommandLineProcess(generalCommandLine.createProcess(), generalCommandLine.getCommandLineString()); + Process process = generalCommandLine.createProcess(); + myPydevConsoleCommunication.setPythonConsoleProcess(process); + return new CommandLineProcess(process, generalCommandLine.getCommandLineString()); } }