PY-18029 Fix TTransportException on Python Console restart

This commit is contained in:
Alexander Koshevoy
2018-08-22 23:16:40 +03:00
parent f2b2262f27
commit 6ec3ff167f
2 changed files with 55 additions and 4 deletions
+41
View File
@@ -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
+14 -4
View File
@@ -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):