Redis: optional pipelining and better tested transactions
This commit is contained in:
parent
8b82004359
commit
5da463e1f7
1 changed files with 234 additions and 120 deletions
|
|
@ -19,10 +19,17 @@ import sockets, os, strutils, parseutils
|
||||||
const
|
const
|
||||||
redisNil* = "\0\0"
|
redisNil* = "\0\0"
|
||||||
|
|
||||||
|
type
|
||||||
|
TPipeline = object
|
||||||
|
enabled: bool
|
||||||
|
buffer: ref string
|
||||||
|
expected: int ## number of replies expected if pipelined
|
||||||
|
|
||||||
type
|
type
|
||||||
TRedis* {.pure, final.} = object
|
TRedis* {.pure, final.} = object
|
||||||
socket: TSocket
|
socket: TSocket
|
||||||
connected: bool
|
connected: bool
|
||||||
|
pipeline: ref TPipeline
|
||||||
|
|
||||||
TRedisStatus* = string
|
TRedisStatus* = string
|
||||||
TRedisInteger* = biggestInt
|
TRedisInteger* = biggestInt
|
||||||
|
|
@ -32,12 +39,20 @@ type
|
||||||
EInvalidReply* = object of ESynch ## Invalid reply from redis
|
EInvalidReply* = object of ESynch ## Invalid reply from redis
|
||||||
ERedis* = object of ESynch ## Error in redis
|
ERedis* = object of ESynch ## Error in redis
|
||||||
|
|
||||||
|
proc newPipeline(): ref TPipeLine =
|
||||||
|
new(result)
|
||||||
|
result.buffer = new string
|
||||||
|
result.buffer[] = ""
|
||||||
|
result.enabled = false
|
||||||
|
result.expected = 0
|
||||||
|
|
||||||
proc open*(host = "localhost", port = 6379.TPort): TRedis =
|
proc open*(host = "localhost", port = 6379.TPort): TRedis =
|
||||||
## Opens a connection to the redis server.
|
## Opens a connection to the redis server.
|
||||||
result.socket = socket(buffered = false)
|
result.socket = socket(buffered = false)
|
||||||
if result.socket == InvalidSocket:
|
if result.socket == InvalidSocket:
|
||||||
OSError(OSLastError())
|
OSError(OSLastError())
|
||||||
result.socket.connect(host, port)
|
result.socket.connect(host, port)
|
||||||
|
result.pipeline = newPipeline()
|
||||||
|
|
||||||
proc raiseInvalidReply(expected, got: char) =
|
proc raiseInvalidReply(expected, got: char) =
|
||||||
raise newException(EInvalidReply,
|
raise newException(EInvalidReply,
|
||||||
|
|
@ -48,8 +63,12 @@ proc raiseNoOK(status: string) =
|
||||||
if status != "QUEUED" and status != "OK":
|
if status != "QUEUED" and status != "OK":
|
||||||
raise newException(EInvalidReply, "Expected \"OK\" got \"$1\"" % status)
|
raise newException(EInvalidReply, "Expected \"OK\" got \"$1\"" % status)
|
||||||
|
|
||||||
proc parseStatus(r: TRedis): TRedisStatus =
|
proc parseStatus(r: TRedis, lineIn: string = ""): TRedisStatus =
|
||||||
var line = ""
|
if r.pipeline.enabled:
|
||||||
|
return "OK"
|
||||||
|
|
||||||
|
var line = lineIn
|
||||||
|
if line == "":
|
||||||
r.socket.readLine(line)
|
r.socket.readLine(line)
|
||||||
if line == "":
|
if line == "":
|
||||||
raise newException(ERedis, "Server closed connection prematurely")
|
raise newException(ERedis, "Server closed connection prematurely")
|
||||||
|
|
@ -61,8 +80,11 @@ proc parseStatus(r: TRedis): TRedisStatus =
|
||||||
|
|
||||||
return line.substr(1) # Strip '+'
|
return line.substr(1) # Strip '+'
|
||||||
|
|
||||||
proc parseInteger(r: TRedis): TRedisInteger =
|
proc parseInteger(r: TRedis, lineIn: string = ""): TRedisInteger =
|
||||||
var line = ""
|
if r.pipeline.enabled: return -1
|
||||||
|
|
||||||
|
var line = lineIn
|
||||||
|
if line == "":
|
||||||
r.socket.readLine(line)
|
r.socket.readLine(line)
|
||||||
|
|
||||||
if line == "+QUEUED": # inside of multi
|
if line == "+QUEUED": # inside of multi
|
||||||
|
|
@ -85,7 +107,9 @@ proc recv(sock: TSocket, size: int): TaintedString =
|
||||||
if sock.recv(cstring(result), size) != size:
|
if sock.recv(cstring(result), size) != size:
|
||||||
raise newException(EInvalidReply, "recv failed")
|
raise newException(EInvalidReply, "recv failed")
|
||||||
|
|
||||||
proc parseSingle(r: TRedis, line:string, allowMBNil = False): TRedisString =
|
proc parseSingleString(r: TRedis, line:string, allowMBNil = False): TRedisString =
|
||||||
|
if r.pipeline.enabled: return ""
|
||||||
|
|
||||||
# Error.
|
# Error.
|
||||||
if line[0] == '-':
|
if line[0] == '-':
|
||||||
raise newException(ERedis, strip(line))
|
raise newException(ERedis, strip(line))
|
||||||
|
|
@ -95,9 +119,6 @@ proc parseSingle(r: TRedis, line:string, allowMBNil = False): TRedisString =
|
||||||
if line == "*-1":
|
if line == "*-1":
|
||||||
return RedisNil
|
return RedisNil
|
||||||
|
|
||||||
if line == "+QUEUED" or line == "+OK" : # inside of a transaction (multi)
|
|
||||||
return nil
|
|
||||||
|
|
||||||
if line[0] != '$':
|
if line[0] != '$':
|
||||||
raiseInvalidReply('$', line[0])
|
raiseInvalidReply('$', line[0])
|
||||||
|
|
||||||
|
|
@ -108,41 +129,86 @@ proc parseSingle(r: TRedis, line:string, allowMBNil = False): TRedisString =
|
||||||
var s = r.socket.recv(numBytes+2)
|
var s = r.socket.recv(numBytes+2)
|
||||||
result = strip(s.string)
|
result = strip(s.string)
|
||||||
|
|
||||||
proc parseMultiLines(r: TRedis, countLine:string): TRedisList =
|
proc parseNext(r: TRedis): TRedisList
|
||||||
|
|
||||||
|
proc parseArrayLines(r: TRedis, countLine:string): TRedisList =
|
||||||
if countLine.string[0] != '*':
|
if countLine.string[0] != '*':
|
||||||
raiseInvalidReply('*', countLine.string[0])
|
raiseInvalidReply('*', countLine.string[0])
|
||||||
|
|
||||||
var numElems = parseInt(countLine.string.substr(1))
|
var numElems = parseInt(countLine.string.substr(1))
|
||||||
if numElems == -1: return nil
|
if numElems == -1: return nil
|
||||||
result = @[]
|
result = @[]
|
||||||
|
|
||||||
for i in 1..numElems:
|
for i in 1..numElems:
|
||||||
var line = ""
|
var parsed = r.parseNext()
|
||||||
r.socket.readLine(line.TaintedString)
|
if not isNil(parsed):
|
||||||
if line[0] == '*': # after exec() may contain more multi-bulk replies
|
|
||||||
var parsed = r.parseMultiLines(line)
|
|
||||||
for item in parsed:
|
for item in parsed:
|
||||||
result.add(item)
|
result.add(item)
|
||||||
else:
|
|
||||||
result.add(r.parseSingle(line))
|
|
||||||
|
|
||||||
proc parseBulk(r: TRedis, allowMBNil = False): TRedisString =
|
proc parseBulkString(r: TRedis, allowMBNil = False, lineIn:string = ""): TRedisString =
|
||||||
var line = ""
|
if r.pipeline.enabled: return ""
|
||||||
|
|
||||||
|
var line = lineIn
|
||||||
|
if line == "":
|
||||||
r.socket.readLine(line.TaintedString)
|
r.socket.readLine(line.TaintedString)
|
||||||
|
|
||||||
if line == "+QUEUED" or line == "+OK": # inside of a transaction (multi)
|
return r.parseSingleString(line, allowMBNil)
|
||||||
return nil
|
|
||||||
|
|
||||||
return r.parseSingle(line, allowMBNil)
|
proc parseArray(r: TRedis): TRedisList =
|
||||||
|
if r.pipeline.enabled: return @[]
|
||||||
proc parseMultiBulk(r: TRedis): TRedisList =
|
|
||||||
var line = TaintedString""
|
var line = TaintedString""
|
||||||
r.socket.readLine(line)
|
r.socket.readLine(line)
|
||||||
|
|
||||||
if line == "+QUEUED": # inside of a transaction (multi)
|
return r.parseArrayLines(line)
|
||||||
return nil
|
|
||||||
|
|
||||||
return r.parseMultiLines(line)
|
proc parseNext(r: TRedis): TRedisList =
|
||||||
|
if r.pipeline.enabled: return @[]
|
||||||
|
var line = TaintedString""
|
||||||
|
r.socket.readLine(line)
|
||||||
|
|
||||||
|
var res = case line[0]
|
||||||
|
of '+': @[r.parseStatus(line)]
|
||||||
|
of '-': @[r.parseStatus(line)]
|
||||||
|
of ':': @[$(r.parseInteger(line))]
|
||||||
|
of '$': @[r.parseBulkString(true,line)]
|
||||||
|
of '*': r.parseArrayLines(line)
|
||||||
|
else:
|
||||||
|
raise newException(EInvalidReply, "parseNext failed on line: " & line)
|
||||||
|
nil
|
||||||
|
r.pipeline.expected -= 1
|
||||||
|
return res
|
||||||
|
|
||||||
|
proc flushPipeline*(r: TRedis, wasMulti = false): TRedisList =
|
||||||
|
## Send buffered commands, clear buffer, return results
|
||||||
|
if r.pipeline.buffer[].len > 0:
|
||||||
|
r.socket.send(r.pipeline.buffer[])
|
||||||
|
r.pipeline.buffer[] = ""
|
||||||
|
|
||||||
|
var prevState = r.pipeline.enabled
|
||||||
|
r.pipeline.enabled = false
|
||||||
|
result = @[]
|
||||||
|
|
||||||
|
var tot = r.pipeline.expected
|
||||||
|
|
||||||
|
for i in 0..tot-1:
|
||||||
|
var ret = r.parseNext()
|
||||||
|
if ret.len == 1 and (ret[0] == "OK" or ret[0] == "QUEUED"):
|
||||||
|
# Skip acknowledgement replies in multi
|
||||||
|
if not wasMulti: result.add(ret)
|
||||||
|
else:
|
||||||
|
result.add(ret)
|
||||||
|
|
||||||
|
r.pipeline.expected = 0
|
||||||
|
r.pipeline.enabled = prevState
|
||||||
|
|
||||||
|
proc setPipeline*(r: TRedis, state: bool) =
|
||||||
|
## Enable or disable command pipelining (reduces network roundtrips).
|
||||||
|
## Note that when enabled, you must call flushPipeline to actually send commands, except
|
||||||
|
## for multi/exec() which enable and flush the pipeline automatically.
|
||||||
|
## Commands return immediately with dummy values; actual results returned from
|
||||||
|
## flushPipeline() or exec()
|
||||||
|
r.pipeline.expected = 0
|
||||||
|
r.pipeline.enabled = state
|
||||||
|
|
||||||
proc sendCommand(r: TRedis, cmd: string, args: varargs[string]) =
|
proc sendCommand(r: TRedis, cmd: string, args: varargs[string]) =
|
||||||
var request = "*" & $(1 + args.len()) & "\c\L"
|
var request = "*" & $(1 + args.len()) & "\c\L"
|
||||||
|
|
@ -151,6 +217,11 @@ proc sendCommand(r: TRedis, cmd: string, args: varargs[string]) =
|
||||||
for i in items(args):
|
for i in items(args):
|
||||||
request.add("$" & $i.len() & "\c\L")
|
request.add("$" & $i.len() & "\c\L")
|
||||||
request.add(i & "\c\L")
|
request.add(i & "\c\L")
|
||||||
|
|
||||||
|
if r.pipeline.enabled:
|
||||||
|
r.pipeline.buffer[].add(request)
|
||||||
|
r.pipeline.expected += 1
|
||||||
|
else:
|
||||||
r.socket.send(request)
|
r.socket.send(request)
|
||||||
|
|
||||||
proc sendCommand(r: TRedis, cmd: string, arg1: string,
|
proc sendCommand(r: TRedis, cmd: string, arg1: string,
|
||||||
|
|
@ -163,6 +234,11 @@ proc sendCommand(r: TRedis, cmd: string, arg1: string,
|
||||||
for i in items(args):
|
for i in items(args):
|
||||||
request.add("$" & $i.len() & "\c\L")
|
request.add("$" & $i.len() & "\c\L")
|
||||||
request.add(i & "\c\L")
|
request.add(i & "\c\L")
|
||||||
|
|
||||||
|
if r.pipeline.enabled:
|
||||||
|
r.pipeline.expected += 1
|
||||||
|
r.pipeline.buffer[].add(request)
|
||||||
|
else:
|
||||||
r.socket.send(request)
|
r.socket.send(request)
|
||||||
|
|
||||||
# Keys
|
# Keys
|
||||||
|
|
@ -192,7 +268,7 @@ proc expireAt*(r: TRedis, key: string, timestamp: int): bool =
|
||||||
proc keys*(r: TRedis, pattern: string): TRedisList =
|
proc keys*(r: TRedis, pattern: string): TRedisList =
|
||||||
## Find all keys matching the given pattern
|
## Find all keys matching the given pattern
|
||||||
r.sendCommand("KEYS", pattern)
|
r.sendCommand("KEYS", pattern)
|
||||||
return r.parseMultiBulk()
|
return r.parseArray()
|
||||||
|
|
||||||
proc move*(r: TRedis, key: string, db: int): bool =
|
proc move*(r: TRedis, key: string, db: int): bool =
|
||||||
## Move a key to another database. Returns `true` on a successful move.
|
## Move a key to another database. Returns `true` on a successful move.
|
||||||
|
|
@ -208,7 +284,7 @@ proc persist*(r: TRedis, key: string): bool =
|
||||||
proc randomKey*(r: TRedis): TRedisString =
|
proc randomKey*(r: TRedis): TRedisString =
|
||||||
## Return a random key from the keyspace
|
## Return a random key from the keyspace
|
||||||
r.sendCommand("RANDOMKEY")
|
r.sendCommand("RANDOMKEY")
|
||||||
return r.parseBulk()
|
return r.parseBulkString()
|
||||||
|
|
||||||
proc rename*(r: TRedis, key, newkey: string): TRedisStatus =
|
proc rename*(r: TRedis, key, newkey: string): TRedisStatus =
|
||||||
## Rename a key.
|
## Rename a key.
|
||||||
|
|
@ -254,7 +330,7 @@ proc decrBy*(r: TRedis, key: string, decrement: int): TRedisInteger =
|
||||||
proc get*(r: TRedis, key: string): TRedisString =
|
proc get*(r: TRedis, key: string): TRedisString =
|
||||||
## Get the value of a key. Returns `redisNil` when `key` doesn't exist.
|
## Get the value of a key. Returns `redisNil` when `key` doesn't exist.
|
||||||
r.sendCommand("GET", key)
|
r.sendCommand("GET", key)
|
||||||
return r.parseBulk()
|
return r.parseBulkString()
|
||||||
|
|
||||||
proc getBit*(r: TRedis, key: string, offset: int): TRedisInteger =
|
proc getBit*(r: TRedis, key: string, offset: int): TRedisInteger =
|
||||||
## Returns the bit value at offset in the string value stored at key
|
## Returns the bit value at offset in the string value stored at key
|
||||||
|
|
@ -264,13 +340,13 @@ proc getBit*(r: TRedis, key: string, offset: int): TRedisInteger =
|
||||||
proc getRange*(r: TRedis, key: string, start, stop: int): TRedisString =
|
proc getRange*(r: TRedis, key: string, start, stop: int): TRedisString =
|
||||||
## Get a substring of the string stored at a key
|
## Get a substring of the string stored at a key
|
||||||
r.sendCommand("GETRANGE", key, $start, $stop)
|
r.sendCommand("GETRANGE", key, $start, $stop)
|
||||||
return r.parseBulk()
|
return r.parseBulkString()
|
||||||
|
|
||||||
proc getSet*(r: TRedis, key: string, value: string): TRedisString =
|
proc getSet*(r: TRedis, key: string, value: string): TRedisString =
|
||||||
## Set the string value of a key and return its old value. Returns `redisNil`
|
## Set the string value of a key and return its old value. Returns `redisNil`
|
||||||
## when key doesn't exist.
|
## when key doesn't exist.
|
||||||
r.sendCommand("GETSET", key, value)
|
r.sendCommand("GETSET", key, value)
|
||||||
return r.parseBulk()
|
return r.parseBulkString()
|
||||||
|
|
||||||
proc incr*(r: TRedis, key: string): TRedisInteger =
|
proc incr*(r: TRedis, key: string): TRedisInteger =
|
||||||
## Increment the integer value of a key by one.
|
## Increment the integer value of a key by one.
|
||||||
|
|
@ -332,12 +408,12 @@ proc hExists*(r: TRedis, key, field: string): bool =
|
||||||
proc hGet*(r: TRedis, key, field: string): TRedisString =
|
proc hGet*(r: TRedis, key, field: string): TRedisString =
|
||||||
## Get the value of a hash field
|
## Get the value of a hash field
|
||||||
r.sendCommand("HGET", key, field)
|
r.sendCommand("HGET", key, field)
|
||||||
return r.parseBulk()
|
return r.parseBulkString()
|
||||||
|
|
||||||
proc hGetAll*(r: TRedis, key: string): TRedisList =
|
proc hGetAll*(r: TRedis, key: string): TRedisList =
|
||||||
## Get all the fields and values in a hash
|
## Get all the fields and values in a hash
|
||||||
r.sendCommand("HGETALL", key)
|
r.sendCommand("HGETALL", key)
|
||||||
return r.parseMultiBulk()
|
return r.parseArray()
|
||||||
|
|
||||||
proc hIncrBy*(r: TRedis, key, field: string, incr: int): TRedisInteger =
|
proc hIncrBy*(r: TRedis, key, field: string, incr: int): TRedisInteger =
|
||||||
## Increment the integer value of a hash field by the given number
|
## Increment the integer value of a hash field by the given number
|
||||||
|
|
@ -347,7 +423,7 @@ proc hIncrBy*(r: TRedis, key, field: string, incr: int): TRedisInteger =
|
||||||
proc hKeys*(r: TRedis, key: string): TRedisList =
|
proc hKeys*(r: TRedis, key: string): TRedisList =
|
||||||
## Get all the fields in a hash
|
## Get all the fields in a hash
|
||||||
r.sendCommand("HKEYS", key)
|
r.sendCommand("HKEYS", key)
|
||||||
return r.parseMultiBulk()
|
return r.parseArray()
|
||||||
|
|
||||||
proc hLen*(r: TRedis, key: string): TRedisInteger =
|
proc hLen*(r: TRedis, key: string): TRedisInteger =
|
||||||
## Get the number of fields in a hash
|
## Get the number of fields in a hash
|
||||||
|
|
@ -357,7 +433,7 @@ proc hLen*(r: TRedis, key: string): TRedisInteger =
|
||||||
proc hMGet*(r: TRedis, key: string, fields: varargs[string]): TRedisList =
|
proc hMGet*(r: TRedis, key: string, fields: varargs[string]): TRedisList =
|
||||||
## Get the values of all the given hash fields
|
## Get the values of all the given hash fields
|
||||||
r.sendCommand("HMGET", key, fields)
|
r.sendCommand("HMGET", key, fields)
|
||||||
return r.parseMultiBulk()
|
return r.parseArray()
|
||||||
|
|
||||||
proc hMSet*(r: TRedis, key: string,
|
proc hMSet*(r: TRedis, key: string,
|
||||||
fieldValues: openarray[tuple[field, value: string]]) =
|
fieldValues: openarray[tuple[field, value: string]]) =
|
||||||
|
|
@ -382,7 +458,7 @@ proc hSetNX*(r: TRedis, key, field, value: string): TRedisInteger =
|
||||||
proc hVals*(r: TRedis, key: string): TRedisList =
|
proc hVals*(r: TRedis, key: string): TRedisList =
|
||||||
## Get all the values in a hash
|
## Get all the values in a hash
|
||||||
r.sendCommand("HVALS", key)
|
r.sendCommand("HVALS", key)
|
||||||
return r.parseMultiBulk()
|
return r.parseArray()
|
||||||
|
|
||||||
# Lists
|
# Lists
|
||||||
|
|
||||||
|
|
@ -393,7 +469,7 @@ proc bLPop*(r: TRedis, keys: varargs[string], timeout: int): TRedisList =
|
||||||
for i in items(keys): args.add(i)
|
for i in items(keys): args.add(i)
|
||||||
args.add($timeout)
|
args.add($timeout)
|
||||||
r.sendCommand("BLPOP", args)
|
r.sendCommand("BLPOP", args)
|
||||||
return r.parseMultiBulk()
|
return r.parseArray()
|
||||||
|
|
||||||
proc bRPop*(r: TRedis, keys: varargs[string], timeout: int): TRedisList =
|
proc bRPop*(r: TRedis, keys: varargs[string], timeout: int): TRedisList =
|
||||||
## Remove and get the *last* element in a list, or block until one
|
## Remove and get the *last* element in a list, or block until one
|
||||||
|
|
@ -402,7 +478,7 @@ proc bRPop*(r: TRedis, keys: varargs[string], timeout: int): TRedisList =
|
||||||
for i in items(keys): args.add(i)
|
for i in items(keys): args.add(i)
|
||||||
args.add($timeout)
|
args.add($timeout)
|
||||||
r.sendCommand("BRPOP", args)
|
r.sendCommand("BRPOP", args)
|
||||||
return r.parseMultiBulk()
|
return r.parseArray()
|
||||||
|
|
||||||
proc bRPopLPush*(r: TRedis, source, destination: string,
|
proc bRPopLPush*(r: TRedis, source, destination: string,
|
||||||
timeout: int): TRedisString =
|
timeout: int): TRedisString =
|
||||||
|
|
@ -411,12 +487,12 @@ proc bRPopLPush*(r: TRedis, source, destination: string,
|
||||||
##
|
##
|
||||||
## http://redis.io/commands/brpoplpush
|
## http://redis.io/commands/brpoplpush
|
||||||
r.sendCommand("BRPOPLPUSH", source, destination, $timeout)
|
r.sendCommand("BRPOPLPUSH", source, destination, $timeout)
|
||||||
return r.parseBulk(true) # Multi-Bulk nil allowed.
|
return r.parseBulkString(true) # Multi-Bulk nil allowed.
|
||||||
|
|
||||||
proc lIndex*(r: TRedis, key: string, index: int): TRedisString =
|
proc lIndex*(r: TRedis, key: string, index: int): TRedisString =
|
||||||
## Get an element from a list by its index
|
## Get an element from a list by its index
|
||||||
r.sendCommand("LINDEX", key, $index)
|
r.sendCommand("LINDEX", key, $index)
|
||||||
return r.parseBulk()
|
return r.parseBulkString()
|
||||||
|
|
||||||
proc lInsert*(r: TRedis, key: string, before: bool, pivot, value: string):
|
proc lInsert*(r: TRedis, key: string, before: bool, pivot, value: string):
|
||||||
TRedisInteger =
|
TRedisInteger =
|
||||||
|
|
@ -433,7 +509,7 @@ proc lLen*(r: TRedis, key: string): TRedisInteger =
|
||||||
proc lPop*(r: TRedis, key: string): TRedisString =
|
proc lPop*(r: TRedis, key: string): TRedisString =
|
||||||
## Remove and get the first element in a list
|
## Remove and get the first element in a list
|
||||||
r.sendCommand("LPOP", key)
|
r.sendCommand("LPOP", key)
|
||||||
return r.parseBulk()
|
return r.parseBulkString()
|
||||||
|
|
||||||
proc lPush*(r: TRedis, key, value: string, create: bool = True): TRedisInteger =
|
proc lPush*(r: TRedis, key, value: string, create: bool = True): TRedisInteger =
|
||||||
## Prepend a value to a list. Returns the length of the list after the push.
|
## Prepend a value to a list. Returns the length of the list after the push.
|
||||||
|
|
@ -450,7 +526,7 @@ proc lRange*(r: TRedis, key: string, start, stop: int): TRedisList =
|
||||||
## Get a range of elements from a list. Returns `nil` when `key`
|
## Get a range of elements from a list. Returns `nil` when `key`
|
||||||
## doesn't exist.
|
## doesn't exist.
|
||||||
r.sendCommand("LRANGE", key, $start, $stop)
|
r.sendCommand("LRANGE", key, $start, $stop)
|
||||||
return r.parseMultiBulk()
|
return r.parseArray()
|
||||||
|
|
||||||
proc lRem*(r: TRedis, key: string, value: string, count: int = 0): TRedisInteger =
|
proc lRem*(r: TRedis, key: string, value: string, count: int = 0): TRedisInteger =
|
||||||
## Remove elements from a list. Returns the number of elements that have been
|
## Remove elements from a list. Returns the number of elements that have been
|
||||||
|
|
@ -471,12 +547,12 @@ proc lTrim*(r: TRedis, key: string, start, stop: int) =
|
||||||
proc rPop*(r: TRedis, key: string): TRedisString =
|
proc rPop*(r: TRedis, key: string): TRedisString =
|
||||||
## Remove and get the last element in a list
|
## Remove and get the last element in a list
|
||||||
r.sendCommand("RPOP", key)
|
r.sendCommand("RPOP", key)
|
||||||
return r.parseBulk()
|
return r.parseBulkString()
|
||||||
|
|
||||||
proc rPopLPush*(r: TRedis, source, destination: string): TRedisString =
|
proc rPopLPush*(r: TRedis, source, destination: string): TRedisString =
|
||||||
## Remove the last element in a list, append it to another list and return it
|
## Remove the last element in a list, append it to another list and return it
|
||||||
r.sendCommand("RPOPLPUSH", source, destination)
|
r.sendCommand("RPOPLPUSH", source, destination)
|
||||||
return r.parseBulk()
|
return r.parseBulkString()
|
||||||
|
|
||||||
proc rPush*(r: TRedis, key, value: string, create: bool = True): TRedisInteger =
|
proc rPush*(r: TRedis, key, value: string, create: bool = True): TRedisInteger =
|
||||||
## Append a value to a list. Returns the length of the list after the push.
|
## Append a value to a list. Returns the length of the list after the push.
|
||||||
|
|
@ -504,7 +580,7 @@ proc scard*(r: TRedis, key: string): TRedisInteger =
|
||||||
proc sdiff*(r: TRedis, keys: varargs[string]): TRedisList =
|
proc sdiff*(r: TRedis, keys: varargs[string]): TRedisList =
|
||||||
## Subtract multiple sets
|
## Subtract multiple sets
|
||||||
r.sendCommand("SDIFF", keys)
|
r.sendCommand("SDIFF", keys)
|
||||||
return r.parseMultiBulk()
|
return r.parseArray()
|
||||||
|
|
||||||
proc sdiffstore*(r: TRedis, destination: string,
|
proc sdiffstore*(r: TRedis, destination: string,
|
||||||
keys: varargs[string]): TRedisInteger =
|
keys: varargs[string]): TRedisInteger =
|
||||||
|
|
@ -515,7 +591,7 @@ proc sdiffstore*(r: TRedis, destination: string,
|
||||||
proc sinter*(r: TRedis, keys: varargs[string]): TRedisList =
|
proc sinter*(r: TRedis, keys: varargs[string]): TRedisList =
|
||||||
## Intersect multiple sets
|
## Intersect multiple sets
|
||||||
r.sendCommand("SINTER", keys)
|
r.sendCommand("SINTER", keys)
|
||||||
return r.parseMultiBulk()
|
return r.parseArray()
|
||||||
|
|
||||||
proc sinterstore*(r: TRedis, destination: string,
|
proc sinterstore*(r: TRedis, destination: string,
|
||||||
keys: varargs[string]): TRedisInteger =
|
keys: varargs[string]): TRedisInteger =
|
||||||
|
|
@ -531,7 +607,7 @@ proc sismember*(r: TRedis, key: string, member: string): TRedisInteger =
|
||||||
proc smembers*(r: TRedis, key: string): TRedisList =
|
proc smembers*(r: TRedis, key: string): TRedisList =
|
||||||
## Get all the members in a set
|
## Get all the members in a set
|
||||||
r.sendCommand("SMEMBERS", key)
|
r.sendCommand("SMEMBERS", key)
|
||||||
return r.parseMultiBulk()
|
return r.parseArray()
|
||||||
|
|
||||||
proc smove*(r: TRedis, source: string, destination: string,
|
proc smove*(r: TRedis, source: string, destination: string,
|
||||||
member: string): TRedisInteger =
|
member: string): TRedisInteger =
|
||||||
|
|
@ -542,12 +618,12 @@ proc smove*(r: TRedis, source: string, destination: string,
|
||||||
proc spop*(r: TRedis, key: string): TRedisString =
|
proc spop*(r: TRedis, key: string): TRedisString =
|
||||||
## Remove and return a random member from a set
|
## Remove and return a random member from a set
|
||||||
r.sendCommand("SPOP", key)
|
r.sendCommand("SPOP", key)
|
||||||
return r.parseBulk()
|
return r.parseBulkString()
|
||||||
|
|
||||||
proc srandmember*(r: TRedis, key: string): TRedisString =
|
proc srandmember*(r: TRedis, key: string): TRedisString =
|
||||||
## Get a random member from a set
|
## Get a random member from a set
|
||||||
r.sendCommand("SRANDMEMBER", key)
|
r.sendCommand("SRANDMEMBER", key)
|
||||||
return r.parseBulk()
|
return r.parseBulkString()
|
||||||
|
|
||||||
proc srem*(r: TRedis, key: string, member: string): TRedisInteger =
|
proc srem*(r: TRedis, key: string, member: string): TRedisInteger =
|
||||||
## Remove a member from a set
|
## Remove a member from a set
|
||||||
|
|
@ -557,7 +633,7 @@ proc srem*(r: TRedis, key: string, member: string): TRedisInteger =
|
||||||
proc sunion*(r: TRedis, keys: varargs[string]): TRedisList =
|
proc sunion*(r: TRedis, keys: varargs[string]): TRedisList =
|
||||||
## Add multiple sets
|
## Add multiple sets
|
||||||
r.sendCommand("SUNION", keys)
|
r.sendCommand("SUNION", keys)
|
||||||
return r.parseMultiBulk()
|
return r.parseArray()
|
||||||
|
|
||||||
proc sunionstore*(r: TRedis, destination: string,
|
proc sunionstore*(r: TRedis, destination: string,
|
||||||
key: varargs[string]): TRedisInteger =
|
key: varargs[string]): TRedisInteger =
|
||||||
|
|
@ -586,7 +662,7 @@ proc zincrby*(r: TRedis, key: string, increment: string,
|
||||||
member: string): TRedisString =
|
member: string): TRedisString =
|
||||||
## Increment the score of a member in a sorted set
|
## Increment the score of a member in a sorted set
|
||||||
r.sendCommand("ZINCRBY", key, increment, member)
|
r.sendCommand("ZINCRBY", key, increment, member)
|
||||||
return r.parseBulk()
|
return r.parseBulkString()
|
||||||
|
|
||||||
proc zinterstore*(r: TRedis, destination: string, numkeys: string,
|
proc zinterstore*(r: TRedis, destination: string, numkeys: string,
|
||||||
keys: openarray[string], weights: openarray[string] = [],
|
keys: openarray[string], weights: openarray[string] = [],
|
||||||
|
|
@ -614,7 +690,7 @@ proc zrange*(r: TRedis, key: string, start: string, stop: string,
|
||||||
r.sendCommand("ZRANGE", key, start, stop)
|
r.sendCommand("ZRANGE", key, start, stop)
|
||||||
else:
|
else:
|
||||||
r.sendCommand("ZRANGE", "WITHSCORES", key, start, stop)
|
r.sendCommand("ZRANGE", "WITHSCORES", key, start, stop)
|
||||||
return r.parseMultiBulk()
|
return r.parseArray()
|
||||||
|
|
||||||
proc zrangebyscore*(r: TRedis, key: string, min: string, max: string,
|
proc zrangebyscore*(r: TRedis, key: string, min: string, max: string,
|
||||||
withScore: bool = false, limit: bool = False,
|
withScore: bool = false, limit: bool = False,
|
||||||
|
|
@ -629,12 +705,12 @@ proc zrangebyscore*(r: TRedis, key: string, min: string, max: string,
|
||||||
args.add($limitCount)
|
args.add($limitCount)
|
||||||
|
|
||||||
r.sendCommand("ZRANGEBYSCORE", args)
|
r.sendCommand("ZRANGEBYSCORE", args)
|
||||||
return r.parseMultiBulk()
|
return r.parseArray()
|
||||||
|
|
||||||
proc zrank*(r: TRedis, key: string, member: string): TRedisString =
|
proc zrank*(r: TRedis, key: string, member: string): TRedisString =
|
||||||
## Determine the index of a member in a sorted set
|
## Determine the index of a member in a sorted set
|
||||||
r.sendCommand("ZRANK", key, member)
|
r.sendCommand("ZRANK", key, member)
|
||||||
return r.parseBulk()
|
return r.parseBulkString()
|
||||||
|
|
||||||
proc zrem*(r: TRedis, key: string, member: string): TRedisInteger =
|
proc zrem*(r: TRedis, key: string, member: string): TRedisInteger =
|
||||||
## Remove a member from a sorted set
|
## Remove a member from a sorted set
|
||||||
|
|
@ -660,7 +736,7 @@ proc zrevrange*(r: TRedis, key: string, start: string, stop: string,
|
||||||
if withScore:
|
if withScore:
|
||||||
r.sendCommand("ZREVRANGE", "WITHSCORE", key, start, stop)
|
r.sendCommand("ZREVRANGE", "WITHSCORE", key, start, stop)
|
||||||
else: r.sendCommand("ZREVRANGE", key, start, stop)
|
else: r.sendCommand("ZREVRANGE", key, start, stop)
|
||||||
return r.parseMultiBulk()
|
return r.parseArray()
|
||||||
|
|
||||||
proc zrevrangebyscore*(r: TRedis, key: string, min: string, max: string,
|
proc zrevrangebyscore*(r: TRedis, key: string, min: string, max: string,
|
||||||
withScore: bool = false, limit: bool = False,
|
withScore: bool = false, limit: bool = False,
|
||||||
|
|
@ -676,18 +752,18 @@ proc zrevrangebyscore*(r: TRedis, key: string, min: string, max: string,
|
||||||
args.add($limitCount)
|
args.add($limitCount)
|
||||||
|
|
||||||
r.sendCommand("ZREVRANGEBYSCORE", args)
|
r.sendCommand("ZREVRANGEBYSCORE", args)
|
||||||
return r.parseMultiBulk()
|
return r.parseArray()
|
||||||
|
|
||||||
proc zrevrank*(r: TRedis, key: string, member: string): TRedisString =
|
proc zrevrank*(r: TRedis, key: string, member: string): TRedisString =
|
||||||
## Determine the index of a member in a sorted set, with
|
## Determine the index of a member in a sorted set, with
|
||||||
## scores ordered from high to low
|
## scores ordered from high to low
|
||||||
r.sendCommand("ZREVRANK", key, member)
|
r.sendCommand("ZREVRANK", key, member)
|
||||||
return r.parseBulk()
|
return r.parseBulkString()
|
||||||
|
|
||||||
proc zscore*(r: TRedis, key: string, member: string): TRedisString =
|
proc zscore*(r: TRedis, key: string, member: string): TRedisString =
|
||||||
## Get the score associated with the given member in a sorted set
|
## Get the score associated with the given member in a sorted set
|
||||||
r.sendCommand("ZSCORE", key, member)
|
r.sendCommand("ZSCORE", key, member)
|
||||||
return r.parseBulk()
|
return r.parseBulkString()
|
||||||
|
|
||||||
proc zunionstore*(r: TRedis, destination: string, numkeys: string,
|
proc zunionstore*(r: TRedis, destination: string, numkeys: string,
|
||||||
keys: openarray[string], weights: openarray[string] = [],
|
keys: openarray[string], weights: openarray[string] = [],
|
||||||
|
|
@ -749,11 +825,14 @@ proc discardMulti*(r: TRedis) =
|
||||||
proc exec*(r: TRedis): TRedisList =
|
proc exec*(r: TRedis): TRedisList =
|
||||||
## Execute all commands issued after MULTI
|
## Execute all commands issued after MULTI
|
||||||
r.sendCommand("EXEC")
|
r.sendCommand("EXEC")
|
||||||
|
r.pipeline.enabled = false
|
||||||
return r.parseMultiBulk()
|
# Will reply with +OK for MULTI/EXEC and +QUEUED for every command
|
||||||
|
# between, then with the results
|
||||||
|
return r.flushPipeline(true)
|
||||||
|
|
||||||
proc multi*(r: TRedis) =
|
proc multi*(r: TRedis) =
|
||||||
## Mark the start of a transaction block
|
## Mark the start of a transaction block
|
||||||
|
r.setPipeline(true)
|
||||||
r.sendCommand("MULTI")
|
r.sendCommand("MULTI")
|
||||||
raiseNoOK(r.parseStatus())
|
raiseNoOK(r.parseStatus())
|
||||||
|
|
||||||
|
|
@ -777,7 +856,7 @@ proc auth*(r: TRedis, password: string) =
|
||||||
proc echoServ*(r: TRedis, message: string): TRedisString =
|
proc echoServ*(r: TRedis, message: string): TRedisString =
|
||||||
## Echo the given string
|
## Echo the given string
|
||||||
r.sendCommand("ECHO", message)
|
r.sendCommand("ECHO", message)
|
||||||
return r.parseBulk()
|
return r.parseBulkString()
|
||||||
|
|
||||||
proc ping*(r: TRedis): TRedisStatus =
|
proc ping*(r: TRedis): TRedisStatus =
|
||||||
## Ping the server
|
## Ping the server
|
||||||
|
|
@ -809,7 +888,7 @@ proc bgsave*(r: TRedis) =
|
||||||
proc configGet*(r: TRedis, parameter: string): TRedisList =
|
proc configGet*(r: TRedis, parameter: string): TRedisList =
|
||||||
## Get the value of a configuration parameter
|
## Get the value of a configuration parameter
|
||||||
r.sendCommand("CONFIG", "GET", parameter)
|
r.sendCommand("CONFIG", "GET", parameter)
|
||||||
return r.parseMultiBulk()
|
return r.parseArray()
|
||||||
|
|
||||||
proc configSet*(r: TRedis, parameter: string, value: string) =
|
proc configSet*(r: TRedis, parameter: string, value: string) =
|
||||||
## Set a configuration parameter to the given value
|
## Set a configuration parameter to the given value
|
||||||
|
|
@ -848,7 +927,7 @@ proc flushdb*(r: TRedis): TRedisStatus =
|
||||||
proc info*(r: TRedis): TRedisString =
|
proc info*(r: TRedis): TRedisString =
|
||||||
## Get information and statistics about the server
|
## Get information and statistics about the server
|
||||||
r.sendCommand("INFO")
|
r.sendCommand("INFO")
|
||||||
return r.parseBulk()
|
return r.parseBulkString()
|
||||||
|
|
||||||
proc lastsave*(r: TRedis): TRedisInteger =
|
proc lastsave*(r: TRedis): TRedisInteger =
|
||||||
## Get the UNIX time stamp of the last successful save to disk
|
## Get the UNIX time stamp of the last successful save to disk
|
||||||
|
|
@ -891,32 +970,67 @@ iterator hPairs*(r: TRedis, key: string): tuple[key, value: string] =
|
||||||
yield (k, i)
|
yield (k, i)
|
||||||
k = ""
|
k = ""
|
||||||
|
|
||||||
|
proc someTests(r: TRedis) =
|
||||||
when false:
|
#r.auth("pass")
|
||||||
# sorry, deactivated for the test suite
|
|
||||||
var r = open()
|
|
||||||
r.auth("pass")
|
|
||||||
|
|
||||||
r.setk("nim:test", "Testing something.")
|
r.setk("nim:test", "Testing something.")
|
||||||
r.setk("nim:utf8", "こんにちは")
|
r.setk("nim:utf8", "こんにちは")
|
||||||
r.setk("nim:esc", "\\ths ągt\\")
|
r.setk("nim:esc", "\\ths ągt\\")
|
||||||
|
r.setk("nim:int", "1")
|
||||||
echo r.get("nim:esc")
|
echo(r.get("nim:esc"))
|
||||||
echo r.incr("nim:int")
|
echo(r.incr("nim:int"))
|
||||||
echo r.incr("nim:int")
|
|
||||||
echo r.get("nim:int")
|
echo r.get("nim:int")
|
||||||
echo r.get("nim:utf8")
|
echo r.get("nim:utf8")
|
||||||
|
echo r.hSet("test1", "name", "A Test")
|
||||||
|
var res = r.hGetAll("test1")
|
||||||
echo repr(r.get("blahasha"))
|
echo repr(r.get("blahasha"))
|
||||||
echo r.randomKey()
|
echo r.randomKey()
|
||||||
|
discard r.lpush("mylist","itema")
|
||||||
|
discard r.lpush("mylist","itemb")
|
||||||
|
r.ltrim("mylist",0,1)
|
||||||
var p = r.lrange("mylist", 0, -1)
|
var p = r.lrange("mylist", 0, -1)
|
||||||
|
|
||||||
for i in items(p):
|
for i in items(p):
|
||||||
|
if not isNil(i):
|
||||||
echo(" ", i)
|
echo(" ", i)
|
||||||
|
|
||||||
echo(r.debugObject("test"))
|
echo(r.debugObject("mylist"))
|
||||||
|
|
||||||
r.configSet("timeout", "299")
|
r.configSet("timeout", "299")
|
||||||
for i in items(r.configGet("timeout")): echo ">> ", i
|
for i in items(r.configGet("timeout")): echo ">> ", i
|
||||||
|
|
||||||
echo r.echoServ("BLAH")
|
echo r.echoServ("BLAH")
|
||||||
|
|
||||||
|
|
||||||
|
when false:
|
||||||
|
var r = open()
|
||||||
|
|
||||||
|
# Test with no pipelining
|
||||||
|
echo("----------------------------------------------")
|
||||||
|
echo("Testing without pipelining.")
|
||||||
|
r.someTests()
|
||||||
|
|
||||||
|
# Test with pipelining enabled
|
||||||
|
echo("//////////////////////////////////////////////")
|
||||||
|
echo()
|
||||||
|
echo("Testing with pipelining.")
|
||||||
|
echo()
|
||||||
|
r.setPipeline(true)
|
||||||
|
r.someTests()
|
||||||
|
var list = r.flushPipeline()
|
||||||
|
r.setPipeline(false)
|
||||||
|
echo("-- list length is " & $list.len & " --")
|
||||||
|
for item in list:
|
||||||
|
if not isNil(item):
|
||||||
|
echo item
|
||||||
|
|
||||||
|
# Test with multi/exec() (automatic pipelining)
|
||||||
|
echo("************************************************")
|
||||||
|
echo("Testing with transaction (automatic pipelining)")
|
||||||
|
r.multi()
|
||||||
|
r.someTests()
|
||||||
|
list = r.exec()
|
||||||
|
echo("-- list length is " & $list.len & " --")
|
||||||
|
for item in list:
|
||||||
|
if not isNil(item):
|
||||||
|
echo item
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue