diff --git a/python/helpers/pydev/pydev_console/thrift_rpc.py b/python/helpers/pydev/pydev_console/thrift_rpc.py index 193955e3dea0..51d60cc48a37 100644 --- a/python/helpers/pydev/pydev_console/thrift_rpc.py +++ b/python/helpers/pydev/pydev_console/thrift_rpc.py @@ -1,35 +1,35 @@ -from pydev_console.thrift_transport import TBidirectionalClientTransport +import threading + +from pydev_console.thrift_transport import TBidirectionalClientTransport, TSyncClient from thriftpy.protocol import TBinaryProtocolFactory from thriftpy.server import TThreadedServer -from thriftpy.thrift import TProcessor, TClient +from thriftpy.thrift import TProcessor def make_rpc_client(client_service, host, port, proto_factory=TBinaryProtocolFactory()): - """ - - :param client_service: - :param server_service: - :param server_handler: - :param host: connection host - :param port: connection port - :param proto_factory: protocol factory for client - :return: - """ - # instantiate client transport = TBidirectionalClientTransport(host, port) protocol = proto_factory.get_protocol(transport) transport.open() - client = TClient(client_service, protocol) + client = TSyncClient(client_service, protocol) server_transport = transport.get_server_transport() return client, server_transport -def make_rpc_server(server_transport, server_service, server_handler, proto_factory=TBinaryProtocolFactory()): +def start_rpc_server(server_transport, server_service, server_handler, proto_factory=TBinaryProtocolFactory(strict_read=False, + strict_write=False)): # setup server processor = TProcessor(server_service, server_handler) - return TThreadedServer(processor, server_transport, iprot_factory=proto_factory) + # todo as `server.serve()` is excessive we may want to get rid of `server` as `TThreadedServer` + server = TThreadedServer(processor, server_transport, iprot_factory=proto_factory) + + client = server.trans.accept() + t = threading.Thread(target=server.handle, args=(client,)) + # t.setDaemon(self.daemon) + t.start() + + return server diff --git a/python/helpers/pydev/pydev_console/thrift_transport.py b/python/helpers/pydev/pydev_console/thrift_transport.py index 318814579a16..b0b26528169e 100644 --- a/python/helpers/pydev/pydev_console/thrift_transport.py +++ b/python/helpers/pydev/pydev_console/thrift_transport.py @@ -257,3 +257,14 @@ class TReversedServerAcceptedTransport(FramedWriter): def read(self, sz): return self._read_fn(sz) + + +class TSyncClient(TClient): + + def __init__(self, service, iprot, oprot=None): + super(TSyncClient, self).__init__(service, iprot, oprot) + self._lock = threading.RLock() + + def _req(self, _api, *args, **kwargs): + with self._lock: + return super(TSyncClient, self)._req(_api, *args, **kwargs) diff --git a/python/helpers/pydev/pydevconsole.py b/python/helpers/pydev/pydevconsole.py index c855fbd6f00a..1fa0ab1145d3 100644 --- a/python/helpers/pydev/pydevconsole.py +++ b/python/helpers/pydev/pydevconsole.py @@ -2,7 +2,7 @@ Entry point module to start the interactive console. ''' from _pydev_imps._pydev_saved_modules import thread -from pydev_console.thrift_rpc import make_rpc_client, make_rpc_server +from pydev_console.thrift_rpc import make_rpc_client, start_rpc_server start_new_thread = thread.start_new_thread @@ -326,6 +326,28 @@ def start_console_server(host, port, interpreter): if connection_queue is not None: connection_queue.put(False) + +def enable_thrift_logging(): + import logging + + # create logger + logger = logging.getLogger('thriftpy') + logger.setLevel(logging.DEBUG) + + # create console handler and set level to debug + ch = logging.StreamHandler() + ch.setLevel(logging.DEBUG) + + # create formatter + formatter = logging.Formatter('%(asctime)s - %(name)s - %(levelname)s - %(message)s') + + # add formatter to ch + ch.setFormatter(formatter) + + # add ch to logger + logger.addHandler(ch) + + def start_server(host, port, client_port, client_host = None): if not client_host: client_host = host @@ -336,11 +358,13 @@ def start_server(host, port, client_port, client_host = None): from pydev_console.thrift_communication import console_thrift + enable_thrift_logging() + client_service = console_thrift.IDE client, server_transport = make_rpc_client(client_service, client_host, client_port) - interpreter = InterpreterInterface(client_host, client_port, threading.currentThread(), client) + interpreter = InterpreterInterface(threading.currentThread(), None, client) # start_new_thread(start_console_server,(host, port, interpreter)) # we do not need to start the server in a new thread because it does not need to accept a client connection, it already has it @@ -355,9 +379,7 @@ def start_server(host, port, client_port, client_host = None): # `InterpreterInterface` implements all methods required for the handler server_handler = interpreter - server = make_rpc_server(server_transport, server_service, server_handler) - # todo as `server.serve()` is excessive we may want to get rid of `server` as `TThreadedServer` - server.serve() + start_rpc_server(server_transport, server_service, server_handler) process_exec_queue(interpreter)