PY-18029 Fix irrelevant exceptions on closing Python Console

Formalize Python Console frontend client/server and Python Console backend process initialization procedures.
This commit is contained in:
Alexander Koshevoy
2018-08-22 23:16:40 +03:00
parent 3c714802b4
commit 53d01c5025
8 changed files with 500 additions and 230 deletions
@@ -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<PythonConsoleBackendService.Client> myInitialPythonConsoleClientFuture = SettableFuture.create();
@NotNull
private final SettableFuture<Process> 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<PythonConsoleFrontendService.Iface> 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<PythonConsoleFrontendService.Iface> 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 <i>handshake</i> message to Python Console backend. Returns
* {@code true} if Python Console backend replies with <i>PyCharm</i> 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 <i>handshake</i> 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<Void> 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<Void> 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.
* <p>
* 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<DebugValue> 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<DebugValue> 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<PyDebugValue>` but this is no longer true
}
catch (PyConsoleProcessFinishedException | TException e) {
catch (CommunicationClosedException | PyConsoleProcessFinishedException | TException e) {
for (PyAsyncValue<String> 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<DebugValue> 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<Void> 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 {
}
}
@@ -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<PythonConsoleFrontendService.Iface>(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 }
}
@@ -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<PythonConsoleFrontendService.Iface>(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)
}
}
@@ -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<String, String> envs = generalCommandLine.getEnvironment();
EncodingEnvironmentUtil.setLocaleEnvironmentIfMac(envs, generalCommandLine.getCharset());
Future<Void> 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());
}
}
@@ -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()
@@ -31,7 +31,7 @@ class TNettyServer private constructor(transport: TServerTransport, processor: T
server.serve()
}
fun stop(): Future<Void?> {
fun stop(): Future<*> {
server.stop()
return object : Future<Void?> {
@@ -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].
@@ -140,8 +140,8 @@ public class PyConsoleTask extends PyExecutionFixtureTestTask {
}
@NotNull
private Future<Void> disposeConsoleAsync() {
Future<Void> shutdownFuture;
private Future<?> disposeConsoleAsync() {
Future<?> shutdownFuture;
if (myCommunication != null) {
shutdownFuture = UIUtil.invokeAndWaitIfNeeded(() -> {
try {