summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorAndy McCurdy <andy@andymccurdy.com>2019-06-01 14:55:25 -0700
committerAndy McCurdy <andy@andymccurdy.com>2019-06-01 14:55:25 -0700
commit0b81453efe61708dcd072eb4831e9107fe590077 (patch)
tree08a9b5a84246b350842a155c2c87f6f48a5b6796
parente498e1c3b5bca5df27452e29d89b6883c6a44d30 (diff)
downloadredis-py-0b81453efe61708dcd072eb4831e9107fe590077.tar.gz
easier to just pass timeout values around rather than confusing block bools
-rwxr-xr-xredis/connection.py67
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"