summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorAndy McCurdy <andy@andymccurdy.com>2011-05-17 01:44:34 -0700
committerAndy McCurdy <andy@andymccurdy.com>2011-05-17 01:44:34 -0700
commit2a3e05c66f2aecae4682da1fca0f627e7aca1ded (patch)
treeee493bce0f0744b112ef90e4c15ca8c12a3b1ddb
parentc650073a51c37e80815b6d5db901e7ae3c893411 (diff)
downloadredis-py-2a3e05c66f2aecae4682da1fca0f627e7aca1ded.tar.gz
all tests pass now except pub/sub. connection_pool's get_connection now always received the command name for the next command. still need to pass keys.
-rw-r--r--redis/client.py24
-rw-r--r--redis/connection.py13
2 files changed, 18 insertions, 19 deletions
diff --git a/redis/client.py b/redis/client.py
index a5b2f9b..b38c5ad 100644
--- a/redis/client.py
+++ b/redis/client.py
@@ -191,19 +191,6 @@ class Redis(threading.local):
encoding_errors=errors
)
- #### Legacy accessors of connection information ####
- def _get_host(self):
- return self.connection.host
- host = property(_get_host)
-
- def _get_port(self):
- return self.connection.port
- port = property(_get_port)
-
- def _get_db(self):
- return self.connection.db
- db = property(_get_db)
-
def pipeline(self, transaction=True):
"""
Return a new pipeline object that can queue multiple commands for
@@ -230,9 +217,9 @@ class Redis(threading.local):
#### COMMAND EXECUTION AND PROTOCOL PARSING ####
def execute_command(self, *args, **options):
- connection = self.connection_pool.get_connection()
+ command_name = args[0]
+ connection = self.connection_pool.get_connection(command_name)
try:
- command_name = args[0]
subscription_command = command_name in self.SUBSCRIPTION_COMMANDS
if self.subscribed and not subscription_command:
raise PubSubError("Cannot issue commands other than SUBSCRIBE "
@@ -1121,6 +1108,9 @@ class Redis(threading.local):
"Return the list of values within hash ``name``"
return self.execute_command('HVALS', name)
+ def pubsub(self):
+ return PubSub(self.connection_pool)
+
# channels
def psubscribe(self, patterns):
@@ -1236,7 +1226,7 @@ class Pipeline(Redis):
return self
def _execute_transaction(self, commands):
- connection = self.connection_pool.get_connection()
+ connection = self.connection_pool.get_connection('MULTI')
try:
all_cmds = ''.join(starmap(connection.pack_command,
[args for args, options in commands]))
@@ -1272,7 +1262,7 @@ class Pipeline(Redis):
def _execute_pipeline(self, commands):
# build up all commands into a single request to increase network perf
- connection = self.connection_pool.get_connection()
+ connection = self.connection_pool.get_connection('MULTI')
try:
all_cmds = ''.join(starmap(connection.pack_command,
[args for args, options in commands]))
diff --git a/redis/connection.py b/redis/connection.py
index 47c4f2f..dc4d40c 100644
--- a/redis/connection.py
+++ b/redis/connection.py
@@ -244,16 +244,25 @@ class ConnectionPool(object):
self.connection_class = connection_class
self.kwargs = kwargs
self._connection = None
+ self._in_use = False
- def get_connection(self, *args, **kwargs):
+ def copy(self):
+ "Return a new instance of this class with the same parameters"
+ return self.__class__(self.connection_class, **self.kwargs)
+
+ def get_connection(self, command_name, *keys):
"Get a connection from the pool"
+ if self._in_use:
+ raise ConnectionError("Connection already in-use")
if not self._connection:
self._connection = self.connection_class(**self.kwargs)
+ self._in_use = True
return self._connection
def release(self, connection):
"Releases the connection back to the pool"
- pass
+ assert self._connection == connection
+ self._in_use = False
def disconnect(self):
"Disconnects all connections in the pool"