PY-18029 Synchronous process of client requests from different threads

This commit is contained in:
Alexander Koshevoy
2018-08-22 23:16:40 +03:00
parent 29b064ada5
commit 51f617151a
3 changed files with 54 additions and 21 deletions
@@ -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
@@ -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)
+27 -5
View File
@@ -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)