diff options
| author | Andy McCurdy <andy@andymccurdy.com> | 2019-06-01 14:55:25 -0700 |
|---|---|---|
| committer | Andy McCurdy <andy@andymccurdy.com> | 2019-06-01 14:55:25 -0700 |
| commit | 0b81453efe61708dcd072eb4831e9107fe590077 (patch) | |
| tree | 08a9b5a84246b350842a155c2c87f6f48a5b6796 | |
| parent | e498e1c3b5bca5df27452e29d89b6883c6a44d30 (diff) | |
| download | redis-py-0b81453efe61708dcd072eb4831e9107fe590077.tar.gz | |
easier to just pass timeout values around rather than confusing block bools
| -rwxr-xr-x | redis/connection.py | 67 |
1 files changed, 37 insertions, 30 deletions
diff --git a/redis/connection.py b/redis/connection.py index da0f47f..e4773ec 100755 --- a/redis/connection.py +++ b/redis/connection.py @@ -61,6 +61,8 @@ SYM_EMPTY = b'' SERVER_CLOSED_CONNECTION_ERROR = "Connection closed by server." +SENTINEL = object() + class Encoder(object): "Encode strings to bytes and decode bytes to strings" @@ -137,23 +139,25 @@ class SocketBuffer(object): def length(self): return self.bytes_written - self.bytes_read - def _read_from_socket(self, length=None, block=True): + def _read_from_socket(self, length=None, timeout=SENTINEL, + raise_on_socket_error=True): sock = self._sock socket_read_size = self.socket_read_size buf = self._buffer buf.seek(self.bytes_written) marker = 0 + custom_timeout = timeout is not SENTINEL try: - if not block: - sock.setblocking(False) + if custom_timeout: + sock.settimeout(timeout) while True: data = recv(self._sock, socket_read_size) # an empty string indicates the server shutdown the socket if isinstance(data, bytes) and len(data) == 0: - if not block: - return False - raise socket.error(SERVER_CLOSED_CONNECTION_ERROR) + if raise_on_socket_error: + raise socket.error(SERVER_CLOSED_CONNECTION_ERROR) + return False buf.write(data) data_length = len(data) self.bytes_written += data_length @@ -167,15 +171,17 @@ class SocketBuffer(object): # blocking error, simply return False indicating that # there's no data to be read. otherwise raise the # original exception. - if not block and ex.errno == EWOULDBLOCK: - return False - raise + if raise_on_socket_error or ex.errno != EWOULDBLOCK: + raise + return False finally: - if not block: + if custom_timeout: sock.settimeout(self.socket_timeout) - def can_read(self): - return bool(self.length) or self._read_from_socket(block=False) + def can_read(self, timeout): + return bool(self.length) or \ + self._read_from_socket(timeout=timeout, + raise_on_socket_error=False) def read(self, length): length = length + 2 # make sure to read the \r\n terminator @@ -264,8 +270,8 @@ class PythonParser(BaseParser): self._buffer = None self.encoder = None - def can_read(self): - return self._buffer and self._buffer.can_read() + def can_read(self, timeout): + return self._buffer and self._buffer.can_read(timeout) def read_response(self): response = self._buffer.readline() @@ -354,35 +360,36 @@ class HiredisParser(BaseParser): self._reader = None self._next_response = False - def can_read(self): + def can_read(self, timeout): if not self._reader: raise ConnectionError(SERVER_CLOSED_CONNECTION_ERROR) if self._next_response is False: self._next_response = self._reader.gets() if self._next_response is False: - return self.read_from_socket(block=False) + return self.read_from_socket(timeout=timeout) return True - def read_from_socket(self, block=True): + def read_from_socket(self, timeout=SENTINEL, raise_on_socket_error=True): sock = self._sock + custom_timeout = timeout is not SENTINEL try: - if not block: - sock.setblocking(False) + if custom_timeout: + sock.settimeout(timeout) if HIREDIS_USE_BYTE_BUFFER: bufflen = recv_into(self._sock, self._buffer) if bufflen == 0: - if not block: - return False - raise socket.error(SERVER_CLOSED_CONNECTION_ERROR) + if raise_on_socket_error: + raise socket.error(SERVER_CLOSED_CONNECTION_ERROR) + return False self._reader.feed(self._buffer, 0, bufflen) else: buffer = recv(self._sock, self.socket_read_size) # an empty string indicates the server shutdown the socket if not isinstance(buffer, bytes) or len(buffer) == 0: - if not block: - return False - raise socket.error(SERVER_CLOSED_CONNECTION_ERROR) + if raise_on_socket_error: + raise socket.error(SERVER_CLOSED_CONNECTION_ERROR) + return False self._reader.feed(buffer) # data was read from the socket and added to the buffer. # return True to indicate that data was read. @@ -392,11 +399,11 @@ class HiredisParser(BaseParser): # blocking error, simply return False indicating that # there's no data to be read. otherwise raise the # original exception. - if not block and ex.errno == EWOULDBLOCK: - return False - raise + if raise_on_socket_error or ex.errno != EWOULDBLOCK: + raise + return False finally: - if not block: + if custom_timeout: sock.settimeout(self._socket_timeout) def read_response(self): @@ -625,7 +632,7 @@ class Connection(object): if not sock: self.connect() sock = self._sock - return self._parser.can_read() + return self._parser.can_read(timeout) def read_response(self): "Read the response from a previously sent command" |
