diff --git a/python/helpers/pydev/_pydev_comm/io.py b/python/helpers/pydev/_pydev_comm/io.py index 482a3832d3fc..549df348ea9b 100644 --- a/python/helpers/pydev/_pydev_comm/io.py +++ b/python/helpers/pydev/_pydev_comm/io.py @@ -14,6 +14,8 @@ class PipeIO(object): self.buffer = bytearray() self.read_pos = 0 + self._closed = False + def _bytes_available(self): return self.read_pos < len(self.buffer) @@ -32,6 +34,9 @@ class PipeIO(object): self.lock.acquire() try: while not self._bytes_available(): + if self._closed: + return bytes() + self.bytes_produced.wait() read_until_pos = min(self.read_pos + sz, len(self.buffer)) @@ -74,3 +79,39 @@ class PipeIO(object): buf_pos = new_buf_pos finally: self.lock.release() + + def close(self): + """Gracefully closes the `PipeIO` + + Allows to read remaining bytes from this `PipeIO`. + + Note that `close()` is expected to be invoked from the same thread as + `write()`. + """ + self.lock.acquire() + try: + self._closed = True + # wake up the reader to let him find out that no more bytes will be + # available + self.bytes_produced.notifyAll() + finally: + self.lock.release() + + +def readall(read_fn, sz): + """Reads `sz` bytes using `read_fn` + + Raises `EOFError` if `read_fn` returned the empty byte array while reading + all `sz` bytes. + """ + buff = b'' + have = 0 + while have < sz: + chunk = read_fn(sz - have) + have += len(chunk) + buff += chunk + + if len(chunk) == 0: + raise EOFError + + return buff diff --git a/python/helpers/pydev/_pydev_comm/transport.py b/python/helpers/pydev/_pydev_comm/transport.py index 0a0ec2e56176..992918bb5a41 100644 --- a/python/helpers/pydev/_pydev_comm/transport.py +++ b/python/helpers/pydev/_pydev_comm/transport.py @@ -2,9 +2,9 @@ import socket import struct import threading +from _pydev_comm.io import PipeIO, readall from _shaded_thriftpy.thrift import TClient -from _shaded_thriftpy.transport import TTransportBase, readall -from _pydev_comm.io import PipeIO +from _shaded_thriftpy.transport import TTransportBase REQUEST = 0 RESPONSE = 1 @@ -59,8 +59,18 @@ class MultiplexedSocketReader(object): t.start() def _read_forever(self): - while True: - self._read_and_dispatch_next_frame() + try: + while True: + self._read_and_dispatch_next_frame() + except EOFError: + # normal Python Console termination + pass + finally: + self._close_pipes() + + def _close_pipes(self): + self._request_pipe.close() + self._response_pipe.close() class SocketWriter(object):