From 53d01c5025e5f7219b9202dd4ef329408e5f5e65 Mon Sep 17 00:00:00 2001 From: Alexander Koshevoy Date: Fri, 10 Aug 2018 13:51:20 +0300 Subject: [PATCH] PY-18029 Fix irrelevant exceptions on closing Python Console Formalize Python Console frontend client/server and Python Console backend process initialization procedures. --- .../console/PydevConsoleCommunication.java | 280 ++++-------------- .../PydevConsoleCommunicationClient.kt | 173 +++++++++++ .../PydevConsoleCommunicationServer.kt | 224 ++++++++++++++ .../console/PydevConsoleRunnerImpl.java | 14 +- .../transport/server/ServerClosedException.kt | 4 + .../console/transport/server/TNettyServer.kt | 2 +- .../transport/server/TNettyServerTransport.kt | 29 +- .../env/python/console/PyConsoleTask.java | 4 +- 8 files changed, 500 insertions(+), 230 deletions(-) create mode 100644 python/src/com/jetbrains/python/console/PydevConsoleCommunicationClient.kt create mode 100644 python/src/com/jetbrains/python/console/PydevConsoleCommunicationServer.kt create mode 100644 python/src/com/jetbrains/python/console/transport/server/ServerClosedException.kt diff --git a/python/src/com/jetbrains/python/console/PydevConsoleCommunication.java b/python/src/com/jetbrains/python/console/PydevConsoleCommunication.java index d2366f62c178..8d44e7f0cd12 100644 --- a/python/src/com/jetbrains/python/console/PydevConsoleCommunication.java +++ b/python/src/com/jetbrains/python/console/PydevConsoleCommunication.java @@ -1,7 +1,6 @@ // 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; @@ -26,16 +25,10 @@ import com.jetbrains.python.console.protocol.*; import com.jetbrains.python.console.pydev.AbstractConsoleCommunication; import com.jetbrains.python.console.pydev.InterpreterResponse; import com.jetbrains.python.console.pydev.PydevCompletionVariant; -import com.jetbrains.python.console.transport.client.TNettyClientTransport; -import com.jetbrains.python.console.transport.server.TNettyServer; -import com.jetbrains.python.console.transport.server.TNettyServerTransport; import com.jetbrains.python.debugger.*; import com.jetbrains.python.debugger.containerview.PyViewNumericContainerAction; import com.jetbrains.python.debugger.pydev.GetVariableCommand; import org.apache.thrift.TException; -import org.apache.thrift.protocol.TBinaryProtocol; -import org.apache.thrift.transport.TServerTransport; -import org.apache.thrift.transport.TTransport; import org.jetbrains.annotations.NotNull; import org.jetbrains.annotations.Nullable; @@ -43,7 +36,10 @@ import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.Map; -import java.util.concurrent.*; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; import static com.jetbrains.python.console.PydevConsoleCommunicationUtil.*; @@ -53,17 +49,7 @@ import static com.jetbrains.python.console.PydevConsoleCommunicationUtil.*; * * @author Fabio */ -public class PydevConsoleCommunication extends AbstractConsoleCommunication implements PyFrameAccessor { - /** - * Thrift RPC client for sending messages to the server. - */ - private PythonConsoleBackendServiceDisposable myClient; - - /** - * This is the server responsible for giving input to a raw_input() requested. - */ - @Nullable private TNettyServer myServer; - +public abstract class PydevConsoleCommunication extends AbstractConsoleCommunication implements PyFrameAccessor { private static final Logger LOG = Logger.getInstance(PydevConsoleCommunication.class); /** @@ -97,17 +83,6 @@ 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. */ @@ -115,99 +90,6 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl super(project); } - public void startServer(int port) throws InterruptedException { - PythonConsoleFrontendHandler serverHandler = new PythonConsoleFrontendHandler(); - PythonConsoleFrontendService.Processor serverProcessor = - new PythonConsoleFrontendService.Processor<>(serverHandler); - //noinspection IOResourceOpenedButNotSafelyClosed - TNettyServerTransport serverTransport = new TNettyServerTransport(port); - TNettyServer server = new TNettyServer(serverTransport, serverProcessor); - - ApplicationManager.getApplication().executeOnPooledThread(() -> server.serve()); - - ApplicationManager.getApplication().executeOnPooledThread(() -> { - TTransport clientTransport = serverTransport.getReverseTransport(); - TBinaryProtocol clientProtocol = new TBinaryProtocol(clientTransport); - PythonConsoleBackendService.Client client = new PythonConsoleBackendService.Client(clientProtocol); - - this.myServer = server; - this.myInitialPythonConsoleClientFuture.set(client); - - PyDebugValueExecutionService executionService = PyDebugValueExecutionService.getInstance(myProject); - executionService.sessionStarted(this); - addFrameListener(new PyFrameListener() { - @Override - public void frameChanged() { - executionService.cancelSubmittedTasks(PydevConsoleCommunication.this); - } - }); - }); - - serverTransport.waitForBind(); - } - - public void startClient(@NotNull String host, int port, @NotNull Process pythonConsoleProcess) { - ApplicationManager.getApplication().executeOnPooledThread(() -> { - TNettyClientTransport clientTransport = new TNettyClientTransport(host, port); - clientTransport.open(); - - TBinaryProtocol clientProtocol = new TBinaryProtocol(clientTransport); - PythonConsoleBackendService.Client client = new PythonConsoleBackendService.Client(clientProtocol); - - TServerTransport serverTransport = clientTransport.getServerTransport(); - - PythonConsoleFrontendHandler serverHandler = new PythonConsoleFrontendHandler(); - PythonConsoleFrontendService.Processor serverProcessor = - new PythonConsoleFrontendService.Processor<>(serverHandler); - - TNettyServer server = new TNettyServer(serverTransport, serverProcessor); - - ApplicationManager.getApplication().executeOnPooledThread(() -> server.serve()); - - this.myServer = server; - this.myInitialPythonConsoleClientFuture.set(client); - this.myPythonConsoleProcessFuture.set(pythonConsoleProcess); - - PyDebugValueExecutionService executionService = PyDebugValueExecutionService.getInstance(myProject); - executionService.sessionStarted(this); - addFrameListener(new PyFrameListener() { - @Override - public void frameChanged() { - executionService.cancelSubmittedTasks(PydevConsoleCommunication.this); - } - }); - }); - } - - @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, Python Console process-aware and disposable * {@link PythonConsoleBackendService.Iface}. Requests to the returned @@ -218,57 +100,39 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl * * @return thread safe and related Python Console process-aware * {@link PythonConsoleBackendService.Iface} + * @throws CommunicationClosedException if transport is closed */ @NotNull - private PythonConsoleBackendServiceDisposable getPythonConsoleBackendClient() { - if (myClient == null) { - myClient = PythonConsoleClientUtil.synchronizedPythonConsoleClient(PydevConsoleCommunication.class.getClassLoader(), - getInitialPythonConsoleBackendClient(), getPythonConsoleProcess()); - } - return myClient; - } + protected abstract PythonConsoleBackendServiceDisposable getPythonConsoleBackendClient(); + /** + * Sends handshake message to Python Console backend. Returns + * {@code true} if Python Console backend replies with PyCharm string. + * Returns {@code false} if Python Console backend replies with unexpected + * message or Python Console process is finished or Python Console is closed. + * + * @return whether handshake with Python Console backend succeeded + * @throws RuntimeException if transport (protocol) error occurs + */ public boolean handshake() { - if (!myClose) { + if (!isCommunicationClosed()) { try { return "PyCharm".equals(getPythonConsoleBackendClient().handshake()); } - catch (PyConsoleProcessFinishedException | TException e) { + catch (CommunicationClosedException | PyConsoleProcessFinishedException e) { + return false; + } + catch (TException e) { throw new RuntimeException(e); } } return false; } - private void sendCloseMessageToScript() { - if (!myClose) { - myClose = true; - new Task.Backgroundable(myProject, "Close Console Communication", true) { - @Override - public void run(@NotNull ProgressIndicator indicator) { - try { - getPythonConsoleBackendClient().close(); - } - catch (Exception e) { - //Ok, we can ignore this one on close. - } - finally { - getPythonConsoleBackendClient().dispose(); - } - } - }.queue(); - } - } - /** * Stops the communication with the client (passes message for it to quit). */ public void close() { - if (myClose) { - return; - } - myClose = true; - PyDebugValueExecutionService.getInstance(myProject).sessionStopped(this); myCallbackHashMap.clear(); @@ -276,32 +140,14 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl @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. - } - finally { - getPythonConsoleBackendClient().dispose(); - } - 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(); - } + closeCommunication().get(); } - finally { - if (myServer != null) { - myServer.stop(); - myServer = null; - } + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + catch (ExecutionException e) { + // it could help us diagnose some intricate cases + LOG.debug(e); } } }.queue(); @@ -310,32 +156,33 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl /** * Stops the communication with the client (passes message for it to quit). * - * @return {@link Future} that allows to wait for Python console server - * thread {@link WebServer#listener} to die + * @return {@link Future} that allows to wait for Python Console transport + * thread(s) to finish its execution */ @NotNull - public Future closeAsync() { - sendCloseMessageToScript(); + public Future closeAsync() { PyDebugValueExecutionService.getInstance(myProject).sessionStopped(this); myCallbackHashMap.clear(); - return ApplicationManager.getApplication().executeOnPooledThread(() -> { - try { - getPythonConsoleProcess().waitFor(5L, TimeUnit.SECONDS); - } - catch (InterruptedException e) { - Thread.currentThread().interrupt(); - } - finally { - if (myServer != null) { - Future stopFuture = myServer.stop(); - myServer = null; - stopFuture.get(); - } - } - return null; - }); + + return closeCommunication(); } + /** + * Closes the communication with Python Console backend gracefully. Returns + * {@link Future} that allows to wait for communication resources + * (corresponding {@link java.util.concurrent.ExecutorService} and threads) + * to be finished. + *

+ * The method is not expected to throw any exception as well as the returned + * {@link Future}. + * + * @return {@link Future} + */ + @NotNull + protected abstract Future closeCommunication(); + + protected abstract boolean isCommunicationClosed(); + /** * Variables that control when we're expecting to give some input to the server or when we're * adding some line to be executed @@ -606,7 +453,7 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl try { getPythonConsoleBackendClient().interrupt(); } - catch (PyConsoleProcessFinishedException | TException e) { + catch (CommunicationClosedException | PyConsoleProcessFinishedException | TException e) { LOG.error(e); } } @@ -623,7 +470,7 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl @Override public PyDebugValue evaluate(String expression, boolean execute, boolean doTrunc) throws PyDebuggerException { - if (!myClose) { + if (!isCommunicationClosed()) { try { List debugValues = getPythonConsoleBackendClient().evaluate(expression); return createPyDebugValue(debugValues.iterator().next(), this); @@ -638,12 +485,12 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl @Nullable @Override public XValueChildrenList loadFrame() throws PyDebuggerException { - if (!myClose) { + if (!isCommunicationClosed()) { try { List frame = getPythonConsoleBackendClient().getFrame(); return parseVars(frame, null, this); } - catch (PyConsoleProcessFinishedException | TException e) { + catch (CommunicationClosedException | PyConsoleProcessFinishedException | TException e) { throw new PyDebuggerException("Get frame from console failed", e); } } @@ -670,7 +517,7 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl // previously `loadFullValue()` might return `List` but this is no longer true } - catch (PyConsoleProcessFinishedException | TException e) { + catch (CommunicationClosedException | PyConsoleProcessFinishedException | TException e) { for (PyAsyncValue asyncValue : pyAsyncValues) { PyDebugValue value = asyncValue.getDebugValue(); XValueNode node = value.getLastNode(); @@ -689,12 +536,12 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl @Override public XValueChildrenList loadVariable(PyDebugValue var) throws PyDebuggerException { - if (!myClose) { + if (!isCommunicationClosed()) { try { List ret = getPythonConsoleBackendClient().getVariable(GetVariableCommand.composeName(var)); return parseVars(ret, var, this); } - catch (PyConsoleProcessFinishedException e) { + catch (CommunicationClosedException | PyConsoleProcessFinishedException e) { throw new PyDebuggerException(e.getLocalizedMessage(), e); } catch (TException e) { @@ -717,13 +564,13 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl @Override public void changeVariable(PyDebugValue variable, String value) throws PyDebuggerException { - if (!myClose) { + if (!isCommunicationClosed()) { try { // NOTE: The actual change is being scheduled in the exec_queue in main thread // This method is async now getPythonConsoleBackendClient().changeVariable(variable.getEvaluationExpression(), value); } - catch (PyConsoleProcessFinishedException | TException e) { + catch (CommunicationClosedException | PyConsoleProcessFinishedException | TException e) { throw new PyDebuggerException("Get change variable", e); } } @@ -738,7 +585,7 @@ 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 (!myClose) { + if (!isCommunicationClosed()) { try { GetArrayResponse ret = getPythonConsoleBackendClient().getArray(var.getName(), rowOffset, colOffset, rows, cols, format); return createArrayChunk(ret, this); @@ -779,7 +626,7 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl // though `connectToDebugger` returns "connect complete" string, let us just ignore it getPythonConsoleBackendClient().connectToDebugger(localPort, dbgOpts, extraEnvs); } - catch (PyConsoleProcessFinishedException | TException e) { + catch (CommunicationClosedException | PyConsoleProcessFinishedException | TException e) { throw new PyDebuggerException("pydevconsole failed to execute connectToDebugger", e); } } @@ -816,8 +663,8 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl } @NotNull - private static Future completedFuture() { - return CompletableFuture.completedFuture(null); + protected final PythonConsoleFrontendService.Iface createPythonConsoleFrontendHandler() { + return new PythonConsoleFrontendHandler(); } private class PythonConsoleFrontendHandler implements PythonConsoleFrontendService.Iface { @@ -870,4 +717,7 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl return execIPythonEditor(path); } } + + protected static class CommunicationClosedException extends RuntimeException { + } } diff --git a/python/src/com/jetbrains/python/console/PydevConsoleCommunicationClient.kt b/python/src/com/jetbrains/python/console/PydevConsoleCommunicationClient.kt new file mode 100644 index 000000000000..0fa42c7caef5 --- /dev/null +++ b/python/src/com/jetbrains/python/console/PydevConsoleCommunicationClient.kt @@ -0,0 +1,173 @@ +// 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.intellij.openapi.application.ApplicationManager +import com.intellij.openapi.progress.ProcessCanceledException +import com.intellij.openapi.progress.ProgressIndicator +import com.intellij.openapi.progress.ProgressIndicatorProvider +import com.intellij.openapi.project.Project +import com.jetbrains.python.console.protocol.PythonConsoleBackendService +import com.jetbrains.python.console.protocol.PythonConsoleFrontendService +import com.jetbrains.python.console.transport.client.TNettyClientTransport +import com.jetbrains.python.console.transport.server.TNettyServer +import com.jetbrains.python.debugger.PyDebugValueExecutionService +import org.apache.thrift.protocol.TBinaryProtocol +import java.util.concurrent.CompletableFuture +import java.util.concurrent.Future +import java.util.concurrent.TimeUnit +import java.util.concurrent.locks.Condition +import java.util.concurrent.locks.Lock +import java.util.concurrent.locks.ReentrantLock +import kotlin.concurrent.withLock + +/** + * This is [PydevConsoleCommunication] where Python Console backend acts like a + * server and IDE acts like a client. + * + * Python Console [Process] is expected to be already started. It is passed as + * [_pythonConsoleProcess] property. + */ +class PydevConsoleCommunicationClient(project: Project, + private val host: String, private val port: Int, + private val _pythonConsoleProcess: Process) : PydevConsoleCommunication(project) { + private var server: TNettyServer? = null + + /** + * Thrift RPC client for sending messages to the server. + * + * Guarded by [stateLock]. + */ + private var client: PythonConsoleBackendServiceDisposable? = null + + private val clientTransport: TNettyClientTransport = TNettyClientTransport(host, port) + + private val stateLock: Lock = ReentrantLock() + private val stateChanged: Condition = stateLock.newCondition() + + /** + * Initial non-thread safe [PythonConsoleBackendService.Client]. + * + * Guarded by [stateLock]. + */ + private var initialPythonConsoleClient: PythonConsoleBackendService.Iface? = null + + /** + * Guarded by [stateLock]. + */ + private var isClosed = false + + /** + * Establishes connection to Python Console backend listening at + * [host]:[port]. + */ + fun connect() { + ApplicationManager.getApplication().executeOnPooledThread { + // TODO handle exception scenario + clientTransport.open() + + val clientProtocol = TBinaryProtocol(clientTransport) + val client = PythonConsoleBackendService.Client(clientProtocol) + + val serverTransport = clientTransport.serverTransport + + val serverHandler = createPythonConsoleFrontendHandler() + val serverProcessor = PythonConsoleFrontendService.Processor(serverHandler) + + val server = TNettyServer(serverTransport, serverProcessor) + + stateLock.withLock { + if (isClosed) throw ProcessCanceledException() + + this.server = server + initialPythonConsoleClient = client + + stateChanged.signalAll() + } + + ApplicationManager.getApplication().executeOnPooledThread { server.serve() } + + val executionService = PyDebugValueExecutionService.getInstance(myProject) + executionService.sessionStarted(this) + addFrameListener { executionService.cancelSubmittedTasks(this@PydevConsoleCommunicationClient) } + } + } + + override fun getPythonConsoleBackendClient(): PythonConsoleBackendServiceDisposable { + stateLock.withLock { + while (!isClosed && _pythonConsoleProcess.isAlive) { + // if `client` is set just return it + client?.let { + return it + } + + val initialPythonConsoleClient = initialPythonConsoleClient + + if (initialPythonConsoleClient != null) { + val newClient = synchronizedPythonConsoleClient(PydevConsoleCommunication::class.java.classLoader, + initialPythonConsoleClient, _pythonConsoleProcess) + client = newClient + return newClient + } + else { + stateChanged.await() + } + } + if (!_pythonConsoleProcess.isAlive) { + throw PyConsoleProcessFinishedException(_pythonConsoleProcess.exitValue()) + } + throw CommunicationClosedException() + } + } + + override fun closeCommunication(): Future<*> { + stateLock.withLock { + try { + isClosed = true + } + finally { + stateChanged.signalAll() + } + } + + // `client` cannot be assigned after `isClosed` is set + + val progressIndicator: ProgressIndicator? = ProgressIndicatorProvider.getInstance().progressIndicator + + // if client exists then try to gracefully `close()` it + try { + client?.apply { + progressIndicator?.text2 = "Sending close message to Python Console..." + + close() + dispose() + } + } + catch (e: Exception) { + // ignore exceptions on `client` shutdown + } + + _pythonConsoleProcess.let { + progressIndicator?.text2 = "Waiting for Python Console process to finish..." + + // TODO move under the future! + try { + do { + progressIndicator?.checkCanceled() + } + while (!it.waitFor(500, TimeUnit.MILLISECONDS)) + } + catch (e: InterruptedException) { + Thread.currentThread().interrupt() + } + } + + + // explicitly close Netty client + clientTransport.close() + + // we know that in this case `server.stop()` would do almost nothing + return server?.stop() ?: CompletableFuture.completedFuture(null) + } + + override fun isCommunicationClosed(): Boolean = stateLock.withLock { isClosed } +} \ No newline at end of file diff --git a/python/src/com/jetbrains/python/console/PydevConsoleCommunicationServer.kt b/python/src/com/jetbrains/python/console/PydevConsoleCommunicationServer.kt new file mode 100644 index 000000000000..14f8fff99a5c --- /dev/null +++ b/python/src/com/jetbrains/python/console/PydevConsoleCommunicationServer.kt @@ -0,0 +1,224 @@ +// 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.intellij.openapi.application.ApplicationManager +import com.intellij.openapi.diagnostic.Logger +import com.intellij.openapi.progress.ProcessCanceledException +import com.intellij.openapi.progress.ProgressIndicator +import com.intellij.openapi.progress.ProgressIndicatorProvider +import com.intellij.openapi.project.Project +import com.jetbrains.python.console.protocol.PythonConsoleBackendService +import com.jetbrains.python.console.protocol.PythonConsoleFrontendService +import com.jetbrains.python.console.transport.server.ServerClosedException +import com.jetbrains.python.console.transport.server.TNettyServer +import com.jetbrains.python.console.transport.server.TNettyServerTransport +import com.jetbrains.python.debugger.PyDebugValueExecutionService +import org.apache.thrift.protocol.TBinaryProtocol +import org.apache.thrift.transport.TTransport +import java.util.concurrent.Future +import java.util.concurrent.TimeUnit +import java.util.concurrent.locks.Condition +import java.util.concurrent.locks.Lock +import java.util.concurrent.locks.ReentrantLock +import kotlin.concurrent.withLock + +class PydevConsoleCommunicationServer(project: Project, port: Int) : PydevConsoleCommunication(project) { + private val serverTransport: TNettyServerTransport + + /** + * This is the server responsible for giving input to a raw_input() requested. + */ + private val server: TNettyServer + + /** + * Thrift RPC client for sending messages to the server. + */ + private var client: PythonConsoleBackendServiceDisposable? = null + + private val stateLock: Lock = ReentrantLock() + private val stateChanged: Condition = stateLock.newCondition() + + /** + * Initial non-thread safe [PythonConsoleBackendService.Client]. + * + * Guarded by [stateLock]. + */ + private var initialPythonConsoleClient: PythonConsoleBackendService.Iface? = null + + /** + * Guarded by [stateLock]. + */ + private var _pythonConsoleProcess: Process? = null + + /** + * Guarded by [stateLock]. + */ + private var isFailedOnBound: Boolean = false + + /** + * Guarded by [stateLock]. + */ + private var isServerBound: Boolean = false + + /** + * Guarded by [stateLock]. + */ + private var isClosed: Boolean = false + + init { + val serverHandler = createPythonConsoleFrontendHandler() + val serverProcessor = PythonConsoleFrontendService.Processor(serverHandler) + //noinspection IOResourceOpenedButNotSafelyClosed + serverTransport = TNettyServerTransport(port) + server = TNettyServer(serverTransport, serverProcessor) + } + + /** + * Must be called once. + */ + fun serve() { + // start server in the separate thread + ApplicationManager.getApplication().executeOnPooledThread { server.serve() } + + ApplicationManager.getApplication().executeOnPooledThread { + // this will wait for the connection of Python Console to the IDE + val clientTransport: TTransport + try { + clientTransport = serverTransport.getReverseTransport() + } + catch (e: ServerClosedException) { + // this is the normal execution flow + throw ProcessCanceledException(e) + } + + val clientProtocol = TBinaryProtocol(clientTransport) + val client = PythonConsoleBackendService.Client(clientProtocol) + + stateLock.withLock { + // early close + if (isClosed) { + server.stop() + + // this is the normal execution flow + throw ProcessCanceledException() + } + + initialPythonConsoleClient = client + + stateChanged.signalAll() + } + + val executionService = PyDebugValueExecutionService.getInstance(myProject) + executionService.sessionStarted(this) + addFrameListener { executionService.cancelSubmittedTasks(this@PydevConsoleCommunicationServer) } + } + + stateLock.withLock { + try { + // waiting on `CountDownLatch.await()` within `stateLock` might be harmful + serverTransport.waitForBind() + isServerBound = true + } + finally { + isFailedOnBound = !isServerBound + + stateChanged.signalAll() + } + } + } + + /** + * The Python Console process is expected to be set after the server is + * bound. + */ + fun setPythonConsoleProcess(pythonConsoleProcess: Process) { + stateLock.withLock { + if (isClosed) { + throw CommunicationClosedException() + } + + if (!isServerBound) { + LOG.warn("Python Console process is set before IDE server is bound, the process may not be able to connect to the server") + } + + _pythonConsoleProcess = pythonConsoleProcess + + stateChanged.signalAll() + } + } + + + override fun getPythonConsoleBackendClient(): PythonConsoleBackendServiceDisposable { + stateLock.withLock { + while (!isClosed && !isFailedOnBound) { + // if `client` is set just return it + client?.let { + return it + } + + val initialPythonConsoleClient = initialPythonConsoleClient + val pythonConsoleProcess = _pythonConsoleProcess + + if (initialPythonConsoleClient != null && pythonConsoleProcess != null) { + val newClient = synchronizedPythonConsoleClient(PydevConsoleCommunication::class.java.classLoader, + initialPythonConsoleClient, pythonConsoleProcess) + client = newClient + return newClient + } + else { + stateChanged.await() + } + } + throw CommunicationClosedException() + } + } + + override fun closeCommunication(): Future<*> { + val progressIndicator: ProgressIndicator? = ProgressIndicatorProvider.getInstance().progressIndicator + + stateLock.withLock { + try { + isClosed = true + } + finally { + stateChanged.signalAll() + } + } + + // if client exists then try to gracefully `close()` it + try { + client?.apply { + progressIndicator?.text2 = "Sending close message to Python Console..." + + close() + dispose() + } + } + catch (e: Exception) { + // ignore exceptions on `client` shutdown + } + + _pythonConsoleProcess?.let { + progressIndicator?.text2 = "Waiting for Python Console process to finish..." + + // TODO move under feature! + try { + do { + progressIndicator?.checkCanceled() + } + while (!it.waitFor(500, TimeUnit.MILLISECONDS)) + } + catch (e: InterruptedException) { + Thread.currentThread().interrupt() + } + } + + return server.stop() + } + + override fun isCommunicationClosed(): Boolean = stateLock.withLock { isClosed } + + companion object { + val LOG: Logger = Logger.getInstance(PydevConsoleCommunicationServer::class.java) + } +} \ No newline at end of file diff --git a/python/src/com/jetbrains/python/console/PydevConsoleRunnerImpl.java b/python/src/com/jetbrains/python/console/PydevConsoleRunnerImpl.java index 0dd972556248..a78efe422fd8 100644 --- a/python/src/com/jetbrains/python/console/PydevConsoleRunnerImpl.java +++ b/python/src/com/jetbrains/python/console/PydevConsoleRunnerImpl.java @@ -94,7 +94,6 @@ import java.util.Collections; import java.util.List; import java.util.Map; import java.util.Scanner; -import java.util.concurrent.Future; import java.util.stream.Collectors; import static com.intellij.execution.runners.AbstractConsoleRunnerWithHistory.registerActionShortcuts; @@ -400,21 +399,18 @@ public class PydevConsoleRunnerImpl implements PydevConsoleRunner { Map envs = generalCommandLine.getEnvironment(); EncodingEnvironmentUtil.setLocaleEnvironmentIfMac(envs, generalCommandLine.getCharset()); - Future connectionFuture; - - myPydevConsoleCommunication = new PydevConsoleCommunication(myProject); - // first of all - start server + PydevConsoleCommunicationServer communicationServer = new PydevConsoleCommunicationServer(myProject, port); + myPydevConsoleCommunication = communicationServer; try { - // todo use process in `PydevConsoleCommunication` - // todo we might want to add a timeout here on start - myPydevConsoleCommunication.startServer(port); + communicationServer.serve(); } catch (Exception e) { + communicationServer.close(); throw new ExecutionException(e.getMessage(), e); } Process process = generalCommandLine.createProcess(); - myPydevConsoleCommunication.setPythonConsoleProcess(process); + communicationServer.setPythonConsoleProcess(process); return new CommandLineProcess(process, generalCommandLine.getCommandLineString()); } } diff --git a/python/src/com/jetbrains/python/console/transport/server/ServerClosedException.kt b/python/src/com/jetbrains/python/console/transport/server/ServerClosedException.kt new file mode 100644 index 000000000000..ca9d5fce2be3 --- /dev/null +++ b/python/src/com/jetbrains/python/console/transport/server/ServerClosedException.kt @@ -0,0 +1,4 @@ +// 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.transport.server + +class ServerClosedException : RuntimeException() \ No newline at end of file diff --git a/python/src/com/jetbrains/python/console/transport/server/TNettyServer.kt b/python/src/com/jetbrains/python/console/transport/server/TNettyServer.kt index e36aa844cf7a..85b5fe47ea5f 100644 --- a/python/src/com/jetbrains/python/console/transport/server/TNettyServer.kt +++ b/python/src/com/jetbrains/python/console/transport/server/TNettyServer.kt @@ -31,7 +31,7 @@ class TNettyServer private constructor(transport: TServerTransport, processor: T server.serve() } - fun stop(): Future { + fun stop(): Future<*> { server.stop() return object : Future { diff --git a/python/src/com/jetbrains/python/console/transport/server/TNettyServerTransport.kt b/python/src/com/jetbrains/python/console/transport/server/TNettyServerTransport.kt index 85839bda381c..413a83edf881 100644 --- a/python/src/com/jetbrains/python/console/transport/server/TNettyServerTransport.kt +++ b/python/src/com/jetbrains/python/console/transport/server/TNettyServerTransport.kt @@ -62,6 +62,7 @@ class TNettyServerTransport(port: Int) : TServerTransport() { nettyServer.close() } + @Throws(InterruptedException::class) fun getReverseTransport(): TTransport = nettyServer.takeReverseTransport() private class NettyServer(val port: Int) { @@ -117,7 +118,7 @@ class TNettyServerTransport(port: Int) : TServerTransport() { ch.pipeline().addLast(DirectedMessageHandler(reverseTransport.outputStream, thriftTransport.outputStream)) - ch.pipeline().addLast(object: ChannelInboundHandlerAdapter() { + ch.pipeline().addLast(object : ChannelInboundHandlerAdapter() { override fun channelInactive(ctx: ChannelHandlerContext) { thriftTransport.close() reverseTransport.close() @@ -156,9 +157,18 @@ class TNettyServerTransport(port: Int) : TServerTransport() { serverBound.countDown() } + /** + * @throws InterruptedException if [CountDownLatch.await] is interrupted + * @throws ServerClosedException if [NettyServer] gets closed + */ @Throws(InterruptedException::class) fun waitForBind() { - serverBound.await() + while (!closed.get()) { + if (serverBound.await(100L, TimeUnit.MILLISECONDS)) { + return + } + } + throw ServerClosedException() } fun accept(): TTransport { @@ -178,7 +188,20 @@ class TNettyServerTransport(port: Int) : TServerTransport() { } } - fun takeReverseTransport(): TTransport = reverseTransportQueue.take() + /** + * @throws InterruptedException if [BlockingQueue.poll] is interrupted + * @throws ServerClosedException if [NettyServer] gets closed + */ + @Throws(InterruptedException::class) + fun takeReverseTransport(): TTransport { + while (!closed.get()) { + val element = reverseTransportQueue.poll(100L, TimeUnit.MILLISECONDS) + if (element != null) { + return element + } + } + throw ServerClosedException() + } /** * Shutdown the server [NioEventLoopGroup]. diff --git a/python/testSrc/com/jetbrains/env/python/console/PyConsoleTask.java b/python/testSrc/com/jetbrains/env/python/console/PyConsoleTask.java index 5654375286c2..e532033e108c 100644 --- a/python/testSrc/com/jetbrains/env/python/console/PyConsoleTask.java +++ b/python/testSrc/com/jetbrains/env/python/console/PyConsoleTask.java @@ -140,8 +140,8 @@ public class PyConsoleTask extends PyExecutionFixtureTestTask { } @NotNull - private Future disposeConsoleAsync() { - Future shutdownFuture; + private Future disposeConsoleAsync() { + Future shutdownFuture; if (myCommunication != null) { shutdownFuture = UIUtil.invokeAndWaitIfNeeded(() -> { try {