diff --git a/python/python-console/src/com/jetbrains/python/console/thrift/TCumulativeTransport.kt b/python/python-console/src/com/jetbrains/python/console/thrift/TCumulativeTransport.kt index 38f83983ed17..e8660e412a29 100644 --- a/python/python-console/src/com/jetbrains/python/console/thrift/TCumulativeTransport.kt +++ b/python/python-console/src/com/jetbrains/python/console/thrift/TCumulativeTransport.kt @@ -4,6 +4,8 @@ import com.intellij.openapi.diagnostic.Logger import io.netty.buffer.ByteBuf import io.netty.buffer.Unpooled import org.apache.thrift.transport.TTransport +import org.apache.thrift.transport.TTransportException +import java.io.IOException import java.io.PipedInputStream import java.io.PipedOutputStream @@ -38,7 +40,15 @@ abstract class TCumulativeTransport : TTransport() { pipedInputStream.close() } - final override fun read(buf: ByteArray, off: Int, len: Int): Int = pipedInputStream.read(buf, off, len) + @Throws(TTransportException::class) + final override fun read(buf: ByteArray, off: Int, len: Int): Int { + try { + return pipedInputStream.read(buf, off, len) + } + catch (e: IOException) { + throw TTransportException(TTransportException.UNKNOWN, e) + } + } companion object { val LOG = Logger.getInstance(TCumulativeTransport::class.java) diff --git a/python/src/com/jetbrains/python/console/PyConsoleProcessFinishedException.kt b/python/src/com/jetbrains/python/console/PyConsoleProcessFinishedException.kt new file mode 100644 index 000000000000..a622d20c0fc0 --- /dev/null +++ b/python/src/com/jetbrains/python/console/PyConsoleProcessFinishedException.kt @@ -0,0 +1,12 @@ +// 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 + +/** + * The [PyConsoleProcessFinishedException] is thrown when IDE is waiting for + * the reply from Python console and the Python console process is discovered + * to be finished. + * + * @see [synchronizedPythonConsoleClient] + */ +class PyConsoleProcessFinishedException(exitValue: Int) + : RuntimeException("Console already exited with value: $exitValue while waiting for an answer.") \ No newline at end of file diff --git a/python/src/com/jetbrains/python/console/PydevConsoleCommunication.java b/python/src/com/jetbrains/python/console/PydevConsoleCommunication.java index eda73394973b..1fcf89125576 100644 --- a/python/src/com/jetbrains/python/console/PydevConsoleCommunication.java +++ b/python/src/com/jetbrains/python/console/PydevConsoleCommunication.java @@ -174,7 +174,7 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl try { return "PyCharm".equals(myClient.handshake()); } - catch (TException e) { + catch (PyConsoleProcessFinishedException | TException e) { throw new RuntimeException(e); } } @@ -501,7 +501,7 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl try { myClient.interrupt(); } - catch (TException e) { + catch (PyConsoleProcessFinishedException | TException e) { LOG.error(e); } } @@ -520,7 +520,6 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl public PyDebugValue evaluate(String expression, boolean execute, boolean doTrunc) throws PyDebuggerException { if (myClient != null) { try { - // @alexander todo add specific exception to the method (previously processed by `checkError()`) List debugValues = myClient.evaluate(expression); return createPyDebugValue(debugValues.iterator().next(), this); } @@ -536,11 +535,10 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl public XValueChildrenList loadFrame() throws PyDebuggerException { if (myClient != null) { try { - // @alexander todo add specific exception to the method (previously processed by `checkError()`) List frame = myClient.getFrame(); return parseVars(frame, null, this); } - catch (TException e) { + catch (PyConsoleProcessFinishedException | TException e) { throw new PyDebuggerException("Get frame from console failed", e); } } @@ -564,12 +562,11 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl final int seq = getNextFullValueSeq(); myCallbackHashMap.put(seq, pyAsyncValues); - // @alexander todo add specific exception to the method (previously processed by `checkError()`) myClient.loadFullValue(seq, evaluationExpressions); // previously `loadFullValue()` might return `List` but this is no longer true } - catch (TException e) { + catch (PyConsoleProcessFinishedException | TException e) { for (PyAsyncValue asyncValue : pyAsyncValues) { PyDebugValue value = asyncValue.getDebugValue(); XValueNode node = value.getLastNode(); @@ -591,10 +588,12 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl public XValueChildrenList loadVariable(PyDebugValue var) throws PyDebuggerException { if (myClient != null) { try { - // @alexander todo add specific exception to the method (previously processed by `checkError()`) List ret = myClient.getVariable(GetVariableCommand.composeName(var)); return parseVars(ret, var, this); } + catch (PyConsoleProcessFinishedException e) { + throw new PyDebuggerException(e.getLocalizedMessage(), e); + } catch (TException e) { throw new PyDebuggerException("Get variable from console failed", e); } @@ -619,10 +618,9 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl try { // NOTE: The actual change is being scheduled in the exec_queue in main thread // This method is async now - // @alexander todo add specific exception to the method (previously processed by `checkError()`) myClient.changeVariable(variable.getEvaluationExpression(), value); } - catch (TException e) { + catch (PyConsoleProcessFinishedException | TException e) { throw new PyDebuggerException("Get change variable", e); } } @@ -678,7 +676,7 @@ public class PydevConsoleCommunication extends AbstractConsoleCommunication impl // though `connectToDebugger` returns "connect complete" string, let us just ignore it myClient.connectToDebugger(localPort, dbgOpts, extraEnvs); } - catch (TException e) { + catch (PyConsoleProcessFinishedException | TException e) { throw new PyDebuggerException("pydevconsole failed to execute connectToDebugger", e); } } diff --git a/python/src/com/jetbrains/python/console/PythonConsoleClientUtil.kt b/python/src/com/jetbrains/python/console/PythonConsoleClientUtil.kt index 8167f3f165c7..1b416eb4e77c 100644 --- a/python/src/com/jetbrains/python/console/PythonConsoleClientUtil.kt +++ b/python/src/com/jetbrains/python/console/PythonConsoleClientUtil.kt @@ -3,38 +3,34 @@ package com.jetbrains.python.console -import com.intellij.openapi.application.ApplicationManager +import com.intellij.util.ConcurrencyUtil import java.lang.reflect.InvocationHandler +import java.lang.reflect.InvocationTargetException +import java.lang.reflect.Method import java.lang.reflect.Proxy -import java.util.concurrent.Callable -import java.util.concurrent.TimeUnit -import java.util.concurrent.TimeoutException -import java.util.concurrent.locks.ReentrantLock +import java.util.concurrent.* + +private const val PYTHON_CONSOLE_COMMAND_THREAD_FACTORY_NAME: String = "Python Console Command Executor" @JvmOverloads fun synchronizedPythonConsoleClient(loader: ClassLoader, delegate: PythonConsoleBackendService.Iface, pythonConsoleProcess: Process? = null): PythonConsoleBackendService.Iface { - val lock = ReentrantLock() + val executorService = newSingleThreadPythonConsoleCommandExecutor() return Proxy.newProxyInstance(loader, arrayOf>(PythonConsoleBackendService.Iface::class.java), InvocationHandler { _, method, args -> - val future = ApplicationManager.getApplication().executeOnPooledThread(Callable { - lock.lock() - try { - if (args == null) { - return@Callable method.invoke(delegate) - } - else { - return@Callable method.invoke(delegate, *args) - } - } - finally { - lock.unlock() - } + // we evaluate the original method in the other thread in order to control it + val future = executorService.submit(Callable { + return@Callable invokeOriginalMethod(args, method, delegate) }) if (pythonConsoleProcess == null) { - return@InvocationHandler future.get() + try { + return@InvocationHandler future.get() + } + catch (e: ExecutionException) { + throw e.cause ?: e + } } while (true) { @@ -44,10 +40,40 @@ fun synchronizedPythonConsoleClient(loader: ClassLoader, catch (e: TimeoutException) { if (!pythonConsoleProcess.isAlive) { val exitValue = pythonConsoleProcess.exitValue() - throw RuntimeException( - "Console already exited with value: $exitValue while waiting for an answer.\n") + throw PyConsoleProcessFinishedException(exitValue) } + // continue waiting for the end of the operation execution + } + catch (e: ExecutionException) { + if (!pythonConsoleProcess.isAlive) { + val exitValue = pythonConsoleProcess.exitValue() + throw PyConsoleProcessFinishedException(exitValue) + } + throw e.cause ?: e } } }) as PythonConsoleBackendService.Iface } + +/** + * Creates new single thread executor with [ThreadPoolExecutor.keepAliveTime] + * equals to 60 sec. + */ +private fun newSingleThreadPythonConsoleCommandExecutor(): ExecutorService { + val threadFactory = ConcurrencyUtil.newNamedThreadFactory(PYTHON_CONSOLE_COMMAND_THREAD_FACTORY_NAME) + return ThreadPoolExecutor(0, 1, 60L, TimeUnit.SECONDS, LinkedBlockingQueue(), threadFactory) +} + +private fun invokeOriginalMethod(args: Array?, method: Method, delegate: Any): Any? { + return try { + if (args != null) { + method.invoke(delegate, *args) + } + else { + method.invoke(delegate) + } + } + catch (e: InvocationTargetException) { + throw e.cause ?: e + } +} \ No newline at end of file