Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
f17d82961d | ||
|
|
f603940f7c | ||
|
|
2c1fc582c7 | ||
|
|
2b99d77a57 | ||
|
|
5f9b9e1194 | ||
|
|
ba2d3e8e6e | ||
|
|
05c1f18d6a | ||
|
|
4acd6973bf | ||
|
|
548e4fe088 | ||
|
|
efa7063c52 | ||
|
|
48568ab6d7 | ||
|
|
0201dea577 | ||
|
|
926beaa3c4 | ||
|
|
019ad3e2e3 | ||
|
|
8d9dff84ce | ||
|
|
fba2e169f9 | ||
|
|
7777be7b0f | ||
|
|
91c1568b1a | ||
|
|
f9c2c1acc6 | ||
|
|
61135f1806 | ||
|
|
e77fba4d03 | ||
|
|
87fe813b3a | ||
|
|
2e0d241420 | ||
|
|
9f7e214e8c | ||
|
|
947077bbcb | ||
|
|
ff2e628f4e | ||
|
|
aefa9caacf | ||
|
|
896cf1a9d9 | ||
|
|
5abb12e04f | ||
|
|
c39a0f7c2a | ||
|
|
8a012df9c6 | ||
|
|
549409ffa5 | ||
|
|
47717222b6 | ||
|
|
d8da89ea5d | ||
|
|
4fcc564a97 | ||
|
|
27d9c729e5 | ||
|
|
de4fb877e8 | ||
|
|
1fade3d3b6 | ||
|
|
9f4d4eef8f | ||
|
|
8eeceabdc2 | ||
|
|
733af14822 | ||
|
|
c9cb699bfb | ||
|
|
b37099a145 | ||
|
|
8fe586d344 | ||
|
|
219e29af65 | ||
|
|
1980bcb7f6 | ||
|
|
57786b14e3 | ||
|
|
2211540d93 | ||
|
|
c85c84be6b | ||
|
|
85b2477099 | ||
|
|
a945e5c066 | ||
|
|
65a2e40ae4 | ||
|
|
d6c70f22cd | ||
|
|
012bcd49ee | ||
|
|
1198f7ceff |
+313
@@ -10,6 +10,319 @@ HIGH: There is a critical bug that may affect a subset of users. Upgrade!
|
||||
CRITICAL: There is a critical bug affecting MOST USERS. Upgrade ASAP.
|
||||
--------------------------------------------------------------------------------
|
||||
|
||||
================================================================================
|
||||
Redis 4.0.8 Released Fri Feb 2 11:17:40 CET 2018
|
||||
================================================================================
|
||||
|
||||
Upgrade urgency CRITICAL ONLY for Redis Cluster users. Otherwise no reason
|
||||
to upgrade at all.
|
||||
|
||||
Redis 4.0.8 fixes a single critical bug in the radix tree data structure
|
||||
used for Redis Cluster keys slot tracking. The problem was actually fixed
|
||||
10 months ago into unstable, but it was fixed in a commit related to Streams
|
||||
so it was never backported (for error) into the 4.0 branch.
|
||||
|
||||
The problem will crash Redis Cluster instances during deletions, but it is
|
||||
very hard to trigger: only when the node removed is in the edge of a memory
|
||||
mapped area there are the conditions to create an issue, because otherwise
|
||||
the code just accesses an out of range word in read-only way in an allocated
|
||||
structure: this is almost always harmless.
|
||||
|
||||
The single commit in this release:
|
||||
|
||||
f603940f Rax updated to latest antirez/rax commit. (Salvatore Sanfilippo)
|
||||
|
||||
Cheers,
|
||||
Salvatore
|
||||
|
||||
================================================================================
|
||||
Redis 4.0.7 Released Wed Jan 24 11:01:40 CET 2018
|
||||
================================================================================
|
||||
|
||||
Upgrade urgency MODERATE: Several bugs fixed, but none of critical level.
|
||||
|
||||
Dear Redis Users,
|
||||
|
||||
Redis 4.0.7 addresses a number of problems and adds a few things that are
|
||||
very useful to have and was worth to backport into a patchlevel release.
|
||||
Here is a list of the most important things, but you can find the full list
|
||||
of commits below as usually:
|
||||
|
||||
* Many 32 bit overflows were addressed in order to allow to use Redis with
|
||||
a very significant amount of data, memory size permitting. (zhaozhao.zz,
|
||||
Oran Agra)
|
||||
|
||||
* MEMORY USAGE fixed for the list type. (gnuhpc)
|
||||
|
||||
* Allow read-only scripts in Redis Cluster. (Salvatore Sanfilippo)
|
||||
|
||||
* Fix AOF pipes setup in edge case. (heqin)
|
||||
|
||||
* AUTH option for MIGRATE. (AlexStocks, Salvatore Sanfilippo, Fabio Nicotra)
|
||||
|
||||
* HyperLogLogs are no longer converted from sparse to dense in order
|
||||
to be merged. (Salvatore Sanfilippo)
|
||||
|
||||
* Fix AOF rewrite dead loop under edge cases. (heqin)
|
||||
|
||||
* Fix processing of large bulk strings (>= 2GB). (Oran Agra)
|
||||
|
||||
* Added RM_UnlinkKey in modules API. (Dvir Volk)
|
||||
|
||||
* Fix Redis Cluster crashes when certain commands with a variable number
|
||||
of arguments are called in an improper way. (Salvatore Sanfilippo)
|
||||
|
||||
* Fix memory leak in lazyfree engine. (zhaozhao.zz)
|
||||
|
||||
* Fix many potentially successful partial synchronizations that end
|
||||
doing a full SYNC, because of a bug destroying the replication
|
||||
backlog on the slave. So after a failover the slave was often not able
|
||||
to PSYNC with masters, and a full SYNC was triggered. The bug only
|
||||
happened after 1 hour of uptime so escaped the unit tests. (Oran Agra)
|
||||
|
||||
* Improve anti-affinity in master/slave allocation for Redis Cluster
|
||||
when the cluster is created. (Salvatore Sanfilippo)
|
||||
|
||||
* Improve output buffer handling for slaves, by not limiting the amount
|
||||
of writes a slave could receive. (Guy Benoish)
|
||||
|
||||
The full list of commits follow.
|
||||
|
||||
Enjoy,
|
||||
Salvatore
|
||||
|
||||
jianqingdu in commit 2b99d77a:
|
||||
fix not call va_end when syncWrite() failed
|
||||
1 file changed, 2 insertions(+), 2 deletions(-)
|
||||
|
||||
Yusaku Kaneta in commit 5f9b9e11:
|
||||
Fix the firstkey, lastkey, and keystep of moduleCommand
|
||||
1 file changed, 1 insertion(+), 1 deletion(-)
|
||||
|
||||
Mark Nunberg in commit ba2d3e8e:
|
||||
redismodule.h: Check ModuleNameBusy before calling it
|
||||
1 file changed, 1 insertion(+), 1 deletion(-)
|
||||
|
||||
antirez in commit 05c1f18d:
|
||||
Fix integration test NOREPLICAS error time dependent false positive.
|
||||
1 file changed, 6 insertions(+), 3 deletions(-)
|
||||
|
||||
antirez in commit 4acd6973:
|
||||
Fix migrateCommand() access of not initialized byte.
|
||||
1 file changed, 5 insertions(+), 2 deletions(-)
|
||||
|
||||
Guy Benoish in commit 548e4fe0:
|
||||
Replication buffer fills up on high rate traffic.
|
||||
1 file changed, 7 insertions(+), 2 deletions(-)
|
||||
|
||||
antirez in commit efa7063c:
|
||||
Cluster: improve anti-affinity algo in redis-trib.rb.
|
||||
1 file changed, 131 insertions(+), 1 deletion(-)
|
||||
|
||||
antirez in commit 48568ab6:
|
||||
Remove useless comment from serverCron().
|
||||
1 file changed, 2 insertions(+), 3 deletions(-)
|
||||
|
||||
heqin in commit 0201dea5:
|
||||
fixbug for #4545 dead loop aof rewrite
|
||||
1 file changed, 1 insertion(+), 1 deletion(-)
|
||||
|
||||
antirez in commit 926beaa3:
|
||||
Hopefully more clear comment to explain the change in #4607.
|
||||
1 file changed, 4 insertions(+), 3 deletions(-)
|
||||
|
||||
qinchao in commit 019ad3e2:
|
||||
fix assert problem in ZIP_DECODE_PREVLENSIZE , see issue: https://github.com/antirez/redis/issues/4587
|
||||
1 file changed, 1 insertion(+), 1 deletion(-)
|
||||
|
||||
Oran Agra in commit 8d9dff84:
|
||||
PSYNC2 fix - promoted slave should hold on to it's backlog
|
||||
1 file changed, 5 insertions(+)
|
||||
|
||||
zhaozhao.zz in commit fba2e169:
|
||||
aof: format code and comment
|
||||
1 file changed, 5 insertions(+), 5 deletions(-)
|
||||
|
||||
antirez in commit 7777be7b:
|
||||
Put more details in the comment introduced by #4601.
|
||||
1 file changed, 8 insertions(+), 3 deletions(-)
|
||||
|
||||
zhaozhao.zz in commit 91c1568b:
|
||||
lazyfree: fix memory leak for lazyfree-lazy-server-del
|
||||
1 file changed, 4 insertions(+), 3 deletions(-)
|
||||
|
||||
antirez in commit f9c2c1ac:
|
||||
Fix getKeysUsingCommandTable() in the case of nagative arity.
|
||||
1 file changed, 7 insertions(+), 5 deletions(-)
|
||||
|
||||
antirez in commit 61135f18:
|
||||
Document new protocol options in #4568 into redis.conf.
|
||||
1 file changed, 14 insertions(+)
|
||||
|
||||
antirez in commit e77fba4d:
|
||||
proto-max-querybuf-len -> client-query-buffer-limit.
|
||||
1 file changed, 4 insertions(+), 4 deletions(-)
|
||||
|
||||
antirez in commit 87fe813b:
|
||||
New config options about protocol prefixed with "proto".
|
||||
4 files changed, 13 insertions(+), 13 deletions(-)
|
||||
|
||||
gnuhpc in commit 2e0d2414:
|
||||
Fix a typo(maybe instruction?) in crash log
|
||||
1 file changed, 1 insertion(+), 1 deletion(-)
|
||||
|
||||
Dvir Volk in commit 9f7e214e:
|
||||
Added RM_UnlinkKey - a low level analog to UNLINK command
|
||||
3 files changed, 56 insertions(+)
|
||||
|
||||
zhaozhao.zz in commit 947077bb:
|
||||
redis-benchmark: bugfix - handle zero liveclients in right way
|
||||
1 file changed, 1 insertion(+), 1 deletion(-)
|
||||
|
||||
Oran Agra in commit ff2e628f:
|
||||
Add config options for max-bulk-len and max-querybuf-len mainly to support RESTORE of large keys
|
||||
4 files changed, 16 insertions(+), 1 deletion(-)
|
||||
|
||||
Oran Agra in commit aefa9caa:
|
||||
fix processing of large bulks (above 2GB)
|
||||
8 files changed, 39 insertions(+), 33 deletions(-)
|
||||
|
||||
heqin in commit 896cf1a9:
|
||||
fixbug for #4545 dead loop aof rewrite
|
||||
1 file changed, 3 insertions(+), 1 deletion(-)
|
||||
|
||||
antirez in commit 5abb12e0:
|
||||
Hyperloglog: refresh hdr variable correctly.
|
||||
1 file changed, 3 insertions(+), 1 deletion(-)
|
||||
|
||||
antirez in commit c39a0f7c:
|
||||
Hyperloglog: Support for PFMERGE sparse encoding as target.
|
||||
1 file changed, 14 insertions(+), 3 deletions(-)
|
||||
|
||||
antirez in commit 8a012df9:
|
||||
Hyperloglog: refactoring of sparse/dense add function.
|
||||
1 file changed, 38 insertions(+), 20 deletions(-)
|
||||
|
||||
antirez in commit 549409ff:
|
||||
Test: MIGRATE AUTH test added.
|
||||
1 file changed, 24 insertions(+)
|
||||
|
||||
antirez in commit 47717222:
|
||||
Rewrite MIGRATE AUTH option.
|
||||
1 file changed, 38 insertions(+), 12 deletions(-)
|
||||
|
||||
heqin in commit d8da89ea:
|
||||
fixbug for #4538 Error opening /setting AOF rewrite IPC pipes: No such file or directory
|
||||
1 file changed, 6 insertions(+), 4 deletions(-)
|
||||
|
||||
antirez in commit 4fcc564a:
|
||||
safe_write -> aofWrite. Function commented.
|
||||
1 file changed, 9 insertions(+), 2 deletions(-)
|
||||
|
||||
zhaozhao.zz in commit 27d9c729:
|
||||
aof: cast sdslen to ssize_t
|
||||
1 file changed, 1 insertion(+), 1 deletion(-)
|
||||
|
||||
zhaozhao.zz in commit de4fb877:
|
||||
aof: fix the short write
|
||||
1 file changed, 22 insertions(+), 1 deletion(-)
|
||||
|
||||
Tomasz Poradowski in commit 1fade3d3:
|
||||
always enable command history in redis-cli
|
||||
1 file changed, 2 insertions(+), 1 deletion(-)
|
||||
|
||||
antirez in commit 9f4d4eef:
|
||||
Cluster: allow read-only EVAL/EVALSHA in slaves.
|
||||
1 file changed, 2 insertions(+), 1 deletion(-)
|
||||
|
||||
nashe in commit 8eeceabd:
|
||||
Prevent off-by-one read in stringmatchlen() (fixes #4527)
|
||||
1 file changed, 1 insertion(+), 1 deletion(-)
|
||||
|
||||
gnuhpc in commit 733af148:
|
||||
Fix memory usage list bug
|
||||
1 file changed, 1 insertion(+), 1 deletion(-)
|
||||
|
||||
zhaozhao.zz in commit c9cb699b:
|
||||
dict: fix the int problem for defrag
|
||||
3 files changed, 5 insertions(+), 5 deletions(-)
|
||||
|
||||
zhaozhao.zz in commit b37099a1:
|
||||
dict: fix the int problem
|
||||
1 file changed, 9 insertions(+), 9 deletions(-)
|
||||
|
||||
zhaozhao.zz in commit 8fe586d3:
|
||||
set: fix the int problem for qsort
|
||||
1 file changed, 8 insertions(+), 2 deletions(-)
|
||||
|
||||
zhaozhao.zz in commit 219e29af:
|
||||
set: fix the int problem for SPOP & SRANDMEMBER
|
||||
1 file changed, 2 insertions(+), 2 deletions(-)
|
||||
|
||||
================================================================================
|
||||
Redis 4.0.6 Released Thu Dec 4 17:54:10 CET 2017
|
||||
================================================================================
|
||||
|
||||
Upgrade urgency CRITICAL: More errors in the fixes for PSYNC2 in Redis 4.0.5
|
||||
were identified.
|
||||
|
||||
This release fixes yet more errors present in the 4.0.5 fixes, that could
|
||||
affect slaves. Moreover another critical issue in quicklists, when they are
|
||||
used at a massive memory scale, was fixed in this release. Upgrading from
|
||||
any 4.0.x release, especially if you are running 4.0.4 or 4.0.5, is highly
|
||||
recommended.
|
||||
|
||||
Note that while this fix for 4.0.6 was written in an hurry as well, this
|
||||
time we took extra precautions in order to avoid writing a broken patch:
|
||||
|
||||
1. The code was reviewed by two developers independently.
|
||||
2. A regression test about the problem introduced in 4.0.4/5 was added.
|
||||
3. Resisting to duplicated Lua scripts loading into the Lua engine is now
|
||||
the default action of the loading function, thus it's simpler to stress
|
||||
its behavior.
|
||||
4. The code section was tested with Valgrind.
|
||||
|
||||
The following is the list of commits included in this release:
|
||||
|
||||
zhaozhao.zz in commit 57786b14:
|
||||
quicklist: change the len of quicklist to unsigned long
|
||||
2 files changed, 4 insertions(+), 4 deletions(-)
|
||||
|
||||
zhaozhao.zz in commit 2211540d:
|
||||
quicklist: fix the return value of quicklistCount
|
||||
2 files changed, 2 insertions(+), 2 deletions(-)
|
||||
|
||||
antirez in commit c85c84be:
|
||||
Refactoring: improve luaCreateFunction() API.
|
||||
3 files changed, 38 insertions(+), 58 deletions(-)
|
||||
|
||||
antirez in commit 85b24770:
|
||||
Remove useless variable check from luaCreateFunction().
|
||||
1 file changed, 1 insertion(+), 1 deletion(-)
|
||||
|
||||
antirez in commit a945e5c0:
|
||||
Fix issue #4505, Lua RDB AUX field loading of existing scripts.
|
||||
1 file changed, 9 insertions(+), 3 deletions(-)
|
||||
|
||||
antirez in commit 65a2e40a:
|
||||
Regression test for #4505 (Lua AUX field loading).
|
||||
1 file changed, 22 insertions(+), 1 deletion(-)
|
||||
|
||||
antirez in commit d6c70f22:
|
||||
DEBUG change-repl-id implemented.
|
||||
1 file changed, 7 insertions(+)
|
||||
|
||||
================================================================================
|
||||
Redis 4.0.5 Released Thu Dec 1 16:03:32 CET 2017
|
||||
================================================================================
|
||||
|
||||
Upgrade urgency CRITICAL: Redis 4.0.4 fix for PSYNC2 was broken, causing the
|
||||
slave to crash when receiving an RDB file from the
|
||||
master that contained a duplicated Lua script.
|
||||
|
||||
Please upgrade ASAP if you are with 4.0.4 and you use any form of Lua scripting
|
||||
because this problem will easily crash Redis.
|
||||
|
||||
================================================================================
|
||||
Redis 4.0.4 Released Thu Nov 30 18:42:12 CET 2017
|
||||
================================================================================
|
||||
|
||||
+14
@@ -1154,6 +1154,20 @@ client-output-buffer-limit normal 0 0 0
|
||||
client-output-buffer-limit slave 256mb 64mb 60
|
||||
client-output-buffer-limit pubsub 32mb 8mb 60
|
||||
|
||||
# Client query buffers accumulate new commands. They are limited to a fixed
|
||||
# amount by default in order to avoid that a protocol desynchronization (for
|
||||
# instance due to a bug in the client) will lead to unbound memory usage in
|
||||
# the query buffer. However you can configure it here if you have very special
|
||||
# needs, such us huge multi/exec requests or alike.
|
||||
#
|
||||
# client-query-buffer-limit 1gb
|
||||
|
||||
# In the Redis protocol, bulk requests, that are, elements representing single
|
||||
# strings, are normally limited ot 512 mb. However you can change this limit
|
||||
# here.
|
||||
#
|
||||
# proto-max-bulk-len 512mb
|
||||
|
||||
# Redis calls an internal function to perform many background tasks, like
|
||||
# closing connections of clients in timeout, purging expired keys that are
|
||||
# never requested, and so forth.
|
||||
|
||||
@@ -237,11 +237,11 @@ void stopAppendOnly(void) {
|
||||
* at runtime using the CONFIG command. */
|
||||
int startAppendOnly(void) {
|
||||
char cwd[MAXPATHLEN]; /* Current working dir path for error messages. */
|
||||
int newfd;
|
||||
|
||||
server.aof_last_fsync = server.unixtime;
|
||||
server.aof_fd = open(server.aof_filename,O_WRONLY|O_APPEND|O_CREAT,0644);
|
||||
newfd = open(server.aof_filename,O_WRONLY|O_APPEND|O_CREAT,0644);
|
||||
serverAssert(server.aof_state == AOF_OFF);
|
||||
if (server.aof_fd == -1) {
|
||||
if (newfd == -1) {
|
||||
char *cwdp = getcwd(cwd,MAXPATHLEN);
|
||||
|
||||
serverLog(LL_WARNING,
|
||||
@@ -256,16 +256,46 @@ int startAppendOnly(void) {
|
||||
server.aof_rewrite_scheduled = 1;
|
||||
serverLog(LL_WARNING,"AOF was enabled but there is already a child process saving an RDB file on disk. An AOF background was scheduled to start when possible.");
|
||||
} else if (rewriteAppendOnlyFileBackground() == C_ERR) {
|
||||
close(server.aof_fd);
|
||||
close(newfd);
|
||||
serverLog(LL_WARNING,"Redis needs to enable the AOF but can't trigger a background AOF rewrite operation. Check the above logs for more info about the error.");
|
||||
return C_ERR;
|
||||
}
|
||||
/* We correctly switched on AOF, now wait for the rewrite to be complete
|
||||
* in order to append data on disk. */
|
||||
server.aof_state = AOF_WAIT_REWRITE;
|
||||
server.aof_last_fsync = server.unixtime;
|
||||
server.aof_fd = newfd;
|
||||
return C_OK;
|
||||
}
|
||||
|
||||
/* This is a wrapper to the write syscall in order to retry on short writes
|
||||
* or if the syscall gets interrupted. It could look strange that we retry
|
||||
* on short writes given that we are writing to a block device: normally if
|
||||
* the first call is short, there is a end-of-space condition, so the next
|
||||
* is likely to fail. However apparently in modern systems this is no longer
|
||||
* true, and in general it looks just more resilient to retry the write. If
|
||||
* there is an actual error condition we'll get it at the next try. */
|
||||
ssize_t aofWrite(int fd, const char *buf, size_t len) {
|
||||
ssize_t nwritten = 0, totwritten = 0;
|
||||
|
||||
while(len) {
|
||||
nwritten = write(fd, buf, len);
|
||||
|
||||
if (nwritten < 0) {
|
||||
if (errno == EINTR) {
|
||||
continue;
|
||||
}
|
||||
return totwritten ? totwritten : -1;
|
||||
}
|
||||
|
||||
len -= nwritten;
|
||||
buf += nwritten;
|
||||
totwritten += nwritten;
|
||||
}
|
||||
|
||||
return totwritten;
|
||||
}
|
||||
|
||||
/* Write the append only file buffer on disk.
|
||||
*
|
||||
* Since we are required to write the AOF before replying to the client,
|
||||
@@ -323,7 +353,7 @@ void flushAppendOnlyFile(int force) {
|
||||
* or alike */
|
||||
|
||||
latencyStartMonitor(latency);
|
||||
nwritten = write(server.aof_fd,server.aof_buf,sdslen(server.aof_buf));
|
||||
nwritten = aofWrite(server.aof_fd,server.aof_buf,sdslen(server.aof_buf));
|
||||
latencyEndMonitor(latency);
|
||||
/* We want to capture different events for delayed writes:
|
||||
* when the delay happens with a pending fsync, or with a saving child
|
||||
@@ -342,7 +372,7 @@ void flushAppendOnlyFile(int force) {
|
||||
/* We performed the write so reset the postponed flush sentinel to zero. */
|
||||
server.aof_flush_postponed_start = 0;
|
||||
|
||||
if (nwritten != (signed)sdslen(server.aof_buf)) {
|
||||
if (nwritten != (ssize_t)sdslen(server.aof_buf)) {
|
||||
static time_t last_write_error_log = 0;
|
||||
int can_log = 0;
|
||||
|
||||
@@ -1285,7 +1315,7 @@ int aofCreatePipes(void) {
|
||||
|
||||
if (pipe(fds) == -1) goto error; /* parent -> children data. */
|
||||
if (pipe(fds+2) == -1) goto error; /* children -> parent ack. */
|
||||
if (pipe(fds+4) == -1) goto error; /* children -> parent ack. */
|
||||
if (pipe(fds+4) == -1) goto error; /* parent -> children ack. */
|
||||
/* Parent -> children data is non blocking. */
|
||||
if (anetNonBlock(NULL,fds[0]) != ANET_OK) goto error;
|
||||
if (anetNonBlock(NULL,fds[1]) != ANET_OK) goto error;
|
||||
@@ -1499,10 +1529,10 @@ void backgroundRewriteDoneHandler(int exitcode, int bysignal) {
|
||||
if (server.aof_fd == -1) {
|
||||
/* AOF disabled */
|
||||
|
||||
/* Don't care if this fails: oldfd will be -1 and we handle that.
|
||||
* One notable case of -1 return is if the old file does
|
||||
* not exist. */
|
||||
oldfd = open(server.aof_filename,O_RDONLY|O_NONBLOCK);
|
||||
/* Don't care if this fails: oldfd will be -1 and we handle that.
|
||||
* One notable case of -1 return is if the old file does
|
||||
* not exist. */
|
||||
oldfd = open(server.aof_filename,O_RDONLY|O_NONBLOCK);
|
||||
} else {
|
||||
/* AOF enabled */
|
||||
oldfd = -1; /* We'll set this to the current AOF filedes later. */
|
||||
|
||||
+43
-13
@@ -4867,14 +4867,16 @@ void migrateCloseTimedoutSockets(void) {
|
||||
dictReleaseIterator(di);
|
||||
}
|
||||
|
||||
/* MIGRATE host port key dbid timeout [COPY | REPLACE]
|
||||
/* MIGRATE host port key dbid timeout [COPY | REPLACE | AUTH password]
|
||||
*
|
||||
* On in the multiple keys form:
|
||||
*
|
||||
* MIGRATE host port "" dbid timeout [COPY | REPLACE] KEYS key1 key2 ... keyN */
|
||||
* MIGRATE host port "" dbid timeout [COPY | REPLACE | AUTH password] KEYS key1
|
||||
* key2 ... keyN */
|
||||
void migrateCommand(client *c) {
|
||||
migrateCachedSocket *cs;
|
||||
int copy, replace, j;
|
||||
int copy = 0, replace = 0, j;
|
||||
char *password = NULL;
|
||||
long timeout;
|
||||
long dbid;
|
||||
robj **ov = NULL; /* Objects to migrate. */
|
||||
@@ -4889,16 +4891,20 @@ void migrateCommand(client *c) {
|
||||
int first_key = 3; /* Argument index of the first key. */
|
||||
int num_keys = 1; /* By default only migrate the 'key' argument. */
|
||||
|
||||
/* Initialization */
|
||||
copy = 0;
|
||||
replace = 0;
|
||||
|
||||
/* Parse additional options */
|
||||
for (j = 6; j < c->argc; j++) {
|
||||
int moreargs = j < c->argc-1;
|
||||
if (!strcasecmp(c->argv[j]->ptr,"copy")) {
|
||||
copy = 1;
|
||||
} else if (!strcasecmp(c->argv[j]->ptr,"replace")) {
|
||||
replace = 1;
|
||||
} else if (!strcasecmp(c->argv[j]->ptr,"auth")) {
|
||||
if (!moreargs) {
|
||||
addReply(c,shared.syntaxerr);
|
||||
return;
|
||||
}
|
||||
j++;
|
||||
password = c->argv[j]->ptr;
|
||||
} else if (!strcasecmp(c->argv[j]->ptr,"keys")) {
|
||||
if (sdslen(c->argv[3]->ptr) != 0) {
|
||||
addReplyError(c,
|
||||
@@ -4957,6 +4963,14 @@ try_again:
|
||||
|
||||
rioInitWithBuffer(&cmd,sdsempty());
|
||||
|
||||
/* Authentication */
|
||||
if (password) {
|
||||
serverAssertWithInfo(c,NULL,rioWriteBulkCount(&cmd,'*',2));
|
||||
serverAssertWithInfo(c,NULL,rioWriteBulkString(&cmd,"AUTH",4));
|
||||
serverAssertWithInfo(c,NULL,rioWriteBulkString(&cmd,password,
|
||||
sdslen(password)));
|
||||
}
|
||||
|
||||
/* Send the SELECT command if the current DB is not already selected. */
|
||||
int select = cs->last_dbid != dbid; /* Should we emit SELECT? */
|
||||
if (select) {
|
||||
@@ -4974,7 +4988,9 @@ try_again:
|
||||
ttl = expireat-mstime();
|
||||
if (ttl < 1) ttl = 1;
|
||||
}
|
||||
serverAssertWithInfo(c,NULL,rioWriteBulkCount(&cmd,'*',replace ? 5 : 4));
|
||||
serverAssertWithInfo(c,NULL,
|
||||
rioWriteBulkCount(&cmd,'*',replace ? 5 : 4));
|
||||
|
||||
if (server.cluster_enabled)
|
||||
serverAssertWithInfo(c,NULL,
|
||||
rioWriteBulkString(&cmd,"RESTORE-ASKING",14));
|
||||
@@ -5017,9 +5033,14 @@ try_again:
|
||||
}
|
||||
}
|
||||
|
||||
char buf0[1024]; /* Auth reply. */
|
||||
char buf1[1024]; /* Select reply. */
|
||||
char buf2[1024]; /* Restore reply. */
|
||||
|
||||
/* Read the AUTH reply if needed. */
|
||||
if (password && syncReadLine(cs->fd, buf0, sizeof(buf0), timeout) <= 0)
|
||||
goto socket_err;
|
||||
|
||||
/* Read the SELECT reply if needed. */
|
||||
if (select && syncReadLine(cs->fd, buf1, sizeof(buf1), timeout) <= 0)
|
||||
goto socket_err;
|
||||
@@ -5036,13 +5057,21 @@ try_again:
|
||||
socket_error = 1;
|
||||
break;
|
||||
}
|
||||
if ((select && buf1[0] == '-') || buf2[0] == '-') {
|
||||
if ((password && buf0[0] == '-') ||
|
||||
(select && buf1[0] == '-') ||
|
||||
buf2[0] == '-')
|
||||
{
|
||||
/* On error assume that last_dbid is no longer valid. */
|
||||
if (!error_from_target) {
|
||||
cs->last_dbid = -1;
|
||||
addReplyErrorFormat(c,"Target instance replied with error: %s",
|
||||
(select && buf1[0] == '-') ? buf1+1 : buf2+1);
|
||||
char *errbuf;
|
||||
if (password && buf0[0] == '-') errbuf = buf0;
|
||||
else if (select && buf1[0] == '-') errbuf = buf1;
|
||||
else errbuf = buf2;
|
||||
|
||||
error_from_target = 1;
|
||||
addReplyErrorFormat(c,"Target instance replied with error: %s",
|
||||
errbuf+1);
|
||||
}
|
||||
} else {
|
||||
if (!copy) {
|
||||
@@ -5107,7 +5136,7 @@ try_again:
|
||||
addReply(c,shared.ok);
|
||||
} else {
|
||||
/* On error we already sent it in the for loop above, and set
|
||||
* the curretly selected socket to -1 to force SELECT the next time. */
|
||||
* the currently selected socket to -1 to force SELECT the next time. */
|
||||
}
|
||||
|
||||
sdsfree(cmd.io.buffer.ptr);
|
||||
@@ -5363,7 +5392,8 @@ clusterNode *getNodeByQuery(client *c, struct redisCommand *cmd, robj **argv, in
|
||||
* node is a slave and the request is about an hash slot our master
|
||||
* is serving, we can reply without redirection. */
|
||||
if (c->flags & CLIENT_READONLY &&
|
||||
cmd->flags & CMD_READONLY &&
|
||||
(cmd->flags & CMD_READONLY || cmd->proc == evalCommand ||
|
||||
cmd->proc == evalShaCommand) &&
|
||||
nodeIsSlave(myself) &&
|
||||
myself->slaveof == n)
|
||||
{
|
||||
|
||||
@@ -328,6 +328,10 @@ void loadServerConfigFromString(char *config) {
|
||||
err = "maxmemory-samples must be 1 or greater";
|
||||
goto loaderr;
|
||||
}
|
||||
} else if ((!strcasecmp(argv[0],"proto-max-bulk-len")) && argc == 2) {
|
||||
server.proto_max_bulk_len = memtoll(argv[1],NULL);
|
||||
} else if ((!strcasecmp(argv[0],"client-query-buffer-limit")) && argc == 2) {
|
||||
server.client_max_querybuf_len = memtoll(argv[1],NULL);
|
||||
} else if (!strcasecmp(argv[0],"lfu-log-factor") && argc == 2) {
|
||||
server.lfu_log_factor = atoi(argv[1]);
|
||||
if (server.lfu_log_factor < 0) {
|
||||
@@ -1132,6 +1136,10 @@ void configSetCommand(client *c) {
|
||||
}
|
||||
freeMemoryIfNeeded();
|
||||
}
|
||||
} config_set_memory_field(
|
||||
"proto-max-bulk-len",server.proto_max_bulk_len) {
|
||||
} config_set_memory_field(
|
||||
"client-query-buffer-limit",server.client_max_querybuf_len) {
|
||||
} config_set_memory_field("repl-backlog-size",ll) {
|
||||
resizeReplicationBacklog(ll);
|
||||
} config_set_memory_field("auto-aof-rewrite-min-size",ll) {
|
||||
@@ -1220,6 +1228,8 @@ void configGetCommand(client *c) {
|
||||
|
||||
/* Numerical values */
|
||||
config_get_numerical_field("maxmemory",server.maxmemory);
|
||||
config_get_numerical_field("proto-max-bulk-len",server.proto_max_bulk_len);
|
||||
config_get_numerical_field("client-query-buffer-limit",server.client_max_querybuf_len);
|
||||
config_get_numerical_field("maxmemory-samples",server.maxmemory_samples);
|
||||
config_get_numerical_field("lfu-log-factor",server.lfu_log_factor);
|
||||
config_get_numerical_field("lfu-decay-time",server.lfu_decay_time);
|
||||
@@ -1992,6 +2002,8 @@ int rewriteConfig(char *path) {
|
||||
rewriteConfigStringOption(state,"requirepass",server.requirepass,NULL);
|
||||
rewriteConfigNumericalOption(state,"maxclients",server.maxclients,CONFIG_DEFAULT_MAX_CLIENTS);
|
||||
rewriteConfigBytesOption(state,"maxmemory",server.maxmemory,CONFIG_DEFAULT_MAXMEMORY);
|
||||
rewriteConfigBytesOption(state,"proto-max-bulk-len",server.proto_max_bulk_len,CONFIG_DEFAULT_PROTO_MAX_BULK_LEN);
|
||||
rewriteConfigBytesOption(state,"client-query-buffer-limit",server.client_max_querybuf_len,PROTO_MAX_QUERYBUF_LEN);
|
||||
rewriteConfigEnumOption(state,"maxmemory-policy",server.maxmemory_policy,maxmemory_policy_enum,CONFIG_DEFAULT_MAXMEMORY_POLICY);
|
||||
rewriteConfigNumericalOption(state,"maxmemory-samples",server.maxmemory_samples,CONFIG_DEFAULT_MAXMEMORY_SAMPLES);
|
||||
rewriteConfigNumericalOption(state,"lfu-log-factor",server.lfu_log_factor,CONFIG_DEFAULT_LFU_LOG_FACTOR);
|
||||
|
||||
@@ -1151,11 +1151,13 @@ int *getKeysUsingCommandTable(struct redisCommand *cmd,robj **argv, int argc, in
|
||||
keys = zmalloc(sizeof(int)*((last - cmd->firstkey)+1));
|
||||
for (j = cmd->firstkey; j <= last; j += cmd->keystep) {
|
||||
if (j >= argc) {
|
||||
/* Modules command do not have dispatch time arity checks, so
|
||||
* we need to handle the case where the user passed an invalid
|
||||
* number of arguments here. In this case we return no keys
|
||||
* and expect the module command to report an arity error. */
|
||||
if (cmd->flags & CMD_MODULE) {
|
||||
/* Modules commands, and standard commands with a not fixed number
|
||||
* of arugments (negative arity parameter) do not have dispatch
|
||||
* time arity checks, so we need to handle the case where the user
|
||||
* passed an invalid number of arguments here. In this case we
|
||||
* return no keys and expect the command implementation to report
|
||||
* an arity or syntax error. */
|
||||
if (cmd->flags & CMD_MODULE || cmd->arity < 0) {
|
||||
zfree(keys);
|
||||
*numkeys = 0;
|
||||
return NULL;
|
||||
|
||||
+10
-3
@@ -308,6 +308,8 @@ void debugCommand(client *c) {
|
||||
"structsize -- Return the size of different Redis core C structures.");
|
||||
blen++; addReplyStatus(c,
|
||||
"htstats <dbid> -- Return hash table statistics of the specified Redis database.");
|
||||
blen++; addReplyStatus(c,
|
||||
"change-repl-id -- Change the replication IDs of the instance. Dangerous, should be used only for testing the replication subsystem.");
|
||||
setDeferredMultiBulkLength(c,blenp,blen);
|
||||
} else if (!strcasecmp(c->argv[1]->ptr,"segfault")) {
|
||||
*((char*)-1) = 'x';
|
||||
@@ -370,13 +372,13 @@ void debugCommand(client *c) {
|
||||
val = dictGetVal(de);
|
||||
strenc = strEncoding(val->encoding);
|
||||
|
||||
char extra[128] = {0};
|
||||
char extra[138] = {0};
|
||||
if (val->encoding == OBJ_ENCODING_QUICKLIST) {
|
||||
char *nextra = extra;
|
||||
int remaining = sizeof(extra);
|
||||
quicklist *ql = val->ptr;
|
||||
/* Add number of quicklist nodes */
|
||||
int used = snprintf(nextra, remaining, " ql_nodes:%u", ql->len);
|
||||
int used = snprintf(nextra, remaining, " ql_nodes:%lu", ql->len);
|
||||
nextra += used;
|
||||
remaining -= used;
|
||||
/* Add average quicklist fill factor */
|
||||
@@ -549,6 +551,11 @@ void debugCommand(client *c) {
|
||||
stats = sdscat(stats,buf);
|
||||
|
||||
addReplyBulkSds(c,stats);
|
||||
} else if (!strcasecmp(c->argv[1]->ptr,"change-repl-id") && c->argc == 2) {
|
||||
serverLog(LL_WARNING,"Changing replication IDs after receiving DEBUG change-repl-id");
|
||||
changeReplicationId();
|
||||
clearReplicationId2();
|
||||
addReply(c,shared.ok);
|
||||
} else {
|
||||
addReplyErrorFormat(c, "Unknown DEBUG subcommand or wrong number of arguments for '%s'",
|
||||
(char*)c->argv[1]->ptr);
|
||||
@@ -1023,7 +1030,7 @@ void sigsegvHandler(int sig, siginfo_t *info, void *secret) {
|
||||
"Redis %s crashed by signal: %d", REDIS_VERSION, sig);
|
||||
if (eip != NULL) {
|
||||
serverLog(LL_WARNING,
|
||||
"Crashed running the instuction at: %p", eip);
|
||||
"Crashed running the instruction at: %p", eip);
|
||||
}
|
||||
if (sig == SIGSEGV || sig == SIGBUS) {
|
||||
serverLog(LL_WARNING,
|
||||
|
||||
+1
-1
@@ -289,7 +289,7 @@ int defragKey(redisDb *db, dictEntry *de) {
|
||||
/* Dirty code:
|
||||
* I can't search in db->expires for that key after i already released
|
||||
* the pointer it holds it won't be able to do the string compare */
|
||||
unsigned int hash = dictGetHash(db->dict, de->key);
|
||||
uint64_t hash = dictGetHash(db->dict, de->key);
|
||||
replaceSateliteDictKeyPtrAndOrDefragDictEntry(db->expires, keysds, newsds, hash, &defragged);
|
||||
}
|
||||
|
||||
|
||||
+11
-11
@@ -66,7 +66,7 @@ static unsigned int dict_force_resize_ratio = 5;
|
||||
|
||||
static int _dictExpandIfNeeded(dict *ht);
|
||||
static unsigned long _dictNextPower(unsigned long size);
|
||||
static int _dictKeyIndex(dict *ht, const void *key, unsigned int hash, dictEntry **existing);
|
||||
static long _dictKeyIndex(dict *ht, const void *key, uint64_t hash, dictEntry **existing);
|
||||
static int _dictInit(dict *ht, dictType *type, void *privDataPtr);
|
||||
|
||||
/* -------------------------- hash functions -------------------------------- */
|
||||
@@ -202,7 +202,7 @@ int dictRehash(dict *d, int n) {
|
||||
de = d->ht[0].table[d->rehashidx];
|
||||
/* Move all the keys in this bucket from the old to the new hash HT */
|
||||
while(de) {
|
||||
unsigned int h;
|
||||
uint64_t h;
|
||||
|
||||
nextde = de->next;
|
||||
/* Get the index in the new hash table */
|
||||
@@ -291,7 +291,7 @@ int dictAdd(dict *d, void *key, void *val)
|
||||
*/
|
||||
dictEntry *dictAddRaw(dict *d, void *key, dictEntry **existing)
|
||||
{
|
||||
int index;
|
||||
long index;
|
||||
dictEntry *entry;
|
||||
dictht *ht;
|
||||
|
||||
@@ -362,7 +362,7 @@ dictEntry *dictAddOrFind(dict *d, void *key) {
|
||||
* dictDelete() and dictUnlink(), please check the top comment
|
||||
* of those functions. */
|
||||
static dictEntry *dictGenericDelete(dict *d, const void *key, int nofree) {
|
||||
unsigned int h, idx;
|
||||
uint64_t h, idx;
|
||||
dictEntry *he, *prevHe;
|
||||
int table;
|
||||
|
||||
@@ -476,7 +476,7 @@ void dictRelease(dict *d)
|
||||
dictEntry *dictFind(dict *d, const void *key)
|
||||
{
|
||||
dictEntry *he;
|
||||
unsigned int h, idx, table;
|
||||
uint64_t h, idx, table;
|
||||
|
||||
if (d->ht[0].used + d->ht[1].used == 0) return NULL; /* dict is empty */
|
||||
if (dictIsRehashing(d)) _dictRehashStep(d);
|
||||
@@ -610,7 +610,7 @@ void dictReleaseIterator(dictIterator *iter)
|
||||
dictEntry *dictGetRandomKey(dict *d)
|
||||
{
|
||||
dictEntry *he, *orighe;
|
||||
unsigned int h;
|
||||
unsigned long h;
|
||||
int listlen, listele;
|
||||
|
||||
if (dictSize(d) == 0) return NULL;
|
||||
@@ -955,9 +955,9 @@ static unsigned long _dictNextPower(unsigned long size)
|
||||
*
|
||||
* Note that if we are in the process of rehashing the hash table, the
|
||||
* index is always returned in the context of the second (new) hash table. */
|
||||
static int _dictKeyIndex(dict *d, const void *key, unsigned int hash, dictEntry **existing)
|
||||
static long _dictKeyIndex(dict *d, const void *key, uint64_t hash, dictEntry **existing)
|
||||
{
|
||||
unsigned int idx, table;
|
||||
unsigned long idx, table;
|
||||
dictEntry *he;
|
||||
if (existing) *existing = NULL;
|
||||
|
||||
@@ -995,7 +995,7 @@ void dictDisableResize(void) {
|
||||
dict_can_resize = 0;
|
||||
}
|
||||
|
||||
unsigned int dictGetHash(dict *d, const void *key) {
|
||||
uint64_t dictGetHash(dict *d, const void *key) {
|
||||
return dictHashKey(d, key);
|
||||
}
|
||||
|
||||
@@ -1004,9 +1004,9 @@ unsigned int dictGetHash(dict *d, const void *key) {
|
||||
* the hash value should be provided using dictGetHash.
|
||||
* no string / key comparison is performed.
|
||||
* return value is the reference to the dictEntry if found, or NULL if not found. */
|
||||
dictEntry **dictFindEntryRefByPtrAndHash(dict *d, const void *oldptr, unsigned int hash) {
|
||||
dictEntry **dictFindEntryRefByPtrAndHash(dict *d, const void *oldptr, uint64_t hash) {
|
||||
dictEntry *he, **heref;
|
||||
unsigned int idx, table;
|
||||
unsigned long idx, table;
|
||||
|
||||
if (d->ht[0].used + d->ht[1].used == 0) return NULL; /* dict is empty */
|
||||
for (table = 0; table <= 1; table++) {
|
||||
|
||||
+2
-2
@@ -178,8 +178,8 @@ int dictRehashMilliseconds(dict *d, int ms);
|
||||
void dictSetHashFunctionSeed(uint8_t *seed);
|
||||
uint8_t *dictGetHashFunctionSeed(void);
|
||||
unsigned long dictScan(dict *d, unsigned long v, dictScanFunction *fn, dictScanBucketFunction *bucketfn, void *privdata);
|
||||
unsigned int dictGetHash(dict *d, const void *key);
|
||||
dictEntry **dictFindEntryRefByPtrAndHash(dict *d, const void *oldptr, unsigned int hash);
|
||||
uint64_t dictGetHash(dict *d, const void *key);
|
||||
dictEntry **dictFindEntryRefByPtrAndHash(dict *d, const void *oldptr, uint64_t hash);
|
||||
|
||||
/* Hash table types */
|
||||
extern dictType dictTypeHeapStringCopyKey;
|
||||
|
||||
+55
-24
@@ -475,9 +475,8 @@ int hllPatLen(unsigned char *ele, size_t elesize, long *regp) {
|
||||
|
||||
/* ================== Dense representation implementation ================== */
|
||||
|
||||
/* "Add" the element in the dense hyperloglog data structure.
|
||||
* Actually nothing is added, but the max 0 pattern counter of the subset
|
||||
* the element belongs to is incremented if needed.
|
||||
/* Low level function to set the dense HLL register at 'index' to the
|
||||
* specified value if the current value is smaller than 'count'.
|
||||
*
|
||||
* 'registers' is expected to have room for HLL_REGISTERS plus an
|
||||
* additional byte on the right. This requirement is met by sds strings
|
||||
@@ -486,12 +485,9 @@ int hllPatLen(unsigned char *ele, size_t elesize, long *regp) {
|
||||
* The function always succeed, however if as a result of the operation
|
||||
* the approximated cardinality changed, 1 is returned. Otherwise 0
|
||||
* is returned. */
|
||||
int hllDenseAdd(uint8_t *registers, unsigned char *ele, size_t elesize) {
|
||||
uint8_t oldcount, count;
|
||||
long index;
|
||||
int hllDenseSet(uint8_t *registers, long index, uint8_t count) {
|
||||
uint8_t oldcount;
|
||||
|
||||
/* Update the register if this element produced a longer run of zeroes. */
|
||||
count = hllPatLen(ele,elesize,&index);
|
||||
HLL_DENSE_GET_REGISTER(oldcount,registers,index);
|
||||
if (count > oldcount) {
|
||||
HLL_DENSE_SET_REGISTER(registers,index,count);
|
||||
@@ -501,6 +497,19 @@ int hllDenseAdd(uint8_t *registers, unsigned char *ele, size_t elesize) {
|
||||
}
|
||||
}
|
||||
|
||||
/* "Add" the element in the dense hyperloglog data structure.
|
||||
* Actually nothing is added, but the max 0 pattern counter of the subset
|
||||
* the element belongs to is incremented if needed.
|
||||
*
|
||||
* This is just a wrapper to hllDenseSet(), performing the hashing of the
|
||||
* element in order to retrieve the index and zero-run count. */
|
||||
int hllDenseAdd(uint8_t *registers, unsigned char *ele, size_t elesize) {
|
||||
long index;
|
||||
uint8_t count = hllPatLen(ele,elesize,&index);
|
||||
/* Update the register if this element produced a longer run of zeroes. */
|
||||
return hllDenseSet(registers,index,count);
|
||||
}
|
||||
|
||||
/* Compute SUM(2^-reg) in the dense representation.
|
||||
* PE is an array with a pre-computer table of values 2^-reg indexed by reg.
|
||||
* As a side effect the integer pointed by 'ezp' is set to the number
|
||||
@@ -623,9 +632,8 @@ int hllSparseToDense(robj *o) {
|
||||
return C_OK;
|
||||
}
|
||||
|
||||
/* "Add" the element in the sparse hyperloglog data structure.
|
||||
* Actually nothing is added, but the max 0 pattern counter of the subset
|
||||
* the element belongs to is incremented if needed.
|
||||
/* Low level function to set the sparse HLL register at 'index' to the
|
||||
* specified value if the current value is smaller than 'count'.
|
||||
*
|
||||
* The object 'o' is the String object holding the HLL. The function requires
|
||||
* a reference to the object in order to be able to enlarge the string if
|
||||
@@ -639,15 +647,12 @@ int hllSparseToDense(robj *o) {
|
||||
* sparse to dense: this happens when a register requires to be set to a value
|
||||
* not representable with the sparse representation, or when the resulting
|
||||
* size would be greater than server.hll_sparse_max_bytes. */
|
||||
int hllSparseAdd(robj *o, unsigned char *ele, size_t elesize) {
|
||||
int hllSparseSet(robj *o, long index, uint8_t count) {
|
||||
struct hllhdr *hdr;
|
||||
uint8_t oldcount, count, *sparse, *end, *p, *prev, *next;
|
||||
long index, first, span;
|
||||
uint8_t oldcount, *sparse, *end, *p, *prev, *next;
|
||||
long first, span;
|
||||
long is_zero = 0, is_xzero = 0, is_val = 0, runlen = 0;
|
||||
|
||||
/* Update the register if this element produced a longer run of zeroes. */
|
||||
count = hllPatLen(ele,elesize,&index);
|
||||
|
||||
/* If the count is too big to be representable by the sparse representation
|
||||
* switch to dense representation. */
|
||||
if (count > HLL_SPARSE_VAL_MAX_VALUE) goto promote;
|
||||
@@ -880,11 +885,24 @@ promote: /* Promote to dense representation. */
|
||||
* Note that this in turn means that PFADD will make sure the command
|
||||
* is propagated to slaves / AOF, so if there is a sparse -> dense
|
||||
* convertion, it will be performed in all the slaves as well. */
|
||||
int dense_retval = hllDenseAdd(hdr->registers, ele, elesize);
|
||||
int dense_retval = hllDenseSet(hdr->registers,index,count);
|
||||
serverAssert(dense_retval == 1);
|
||||
return dense_retval;
|
||||
}
|
||||
|
||||
/* "Add" the element in the sparse hyperloglog data structure.
|
||||
* Actually nothing is added, but the max 0 pattern counter of the subset
|
||||
* the element belongs to is incremented if needed.
|
||||
*
|
||||
* This function is actually a wrapper for hllSparseSet(), it only performs
|
||||
* the hashshing of the elmenet to obtain the index and zeros run length. */
|
||||
int hllSparseAdd(robj *o, unsigned char *ele, size_t elesize) {
|
||||
long index;
|
||||
uint8_t count = hllPatLen(ele,elesize,&index);
|
||||
/* Update the register if this element produced a longer run of zeroes. */
|
||||
return hllSparseSet(o,index,count);
|
||||
}
|
||||
|
||||
/* Compute SUM(2^-reg) in the sparse representation.
|
||||
* PE is an array with a pre-computer table of values 2^-reg indexed by reg.
|
||||
* As a side effect the integer pointed by 'ezp' is set to the number
|
||||
@@ -1280,9 +1298,10 @@ void pfmergeCommand(client *c) {
|
||||
uint8_t max[HLL_REGISTERS];
|
||||
struct hllhdr *hdr;
|
||||
int j;
|
||||
int use_dense = 0; /* Use dense representation as target? */
|
||||
|
||||
/* Compute an HLL with M[i] = MAX(M[i]_j).
|
||||
* We we the maximum into the max array of registers. We'll write
|
||||
* We store the maximum into the max array of registers. We'll write
|
||||
* it to the target variable later. */
|
||||
memset(max,0,sizeof(max));
|
||||
for (j = 1; j < c->argc; j++) {
|
||||
@@ -1291,6 +1310,11 @@ void pfmergeCommand(client *c) {
|
||||
if (o == NULL) continue; /* Assume empty HLL for non existing var. */
|
||||
if (isHLLObjectOrReply(c,o) != C_OK) return;
|
||||
|
||||
/* If at least one involved HLL is dense, use the dense representation
|
||||
* as target ASAP to save time and avoid the conversion step. */
|
||||
hdr = o->ptr;
|
||||
if (hdr->encoding == HLL_DENSE) use_dense = 1;
|
||||
|
||||
/* Merge with this HLL with our 'max' HHL by setting max[i]
|
||||
* to MAX(max[i],hll[i]). */
|
||||
if (hllMerge(max,o) == C_ERR) {
|
||||
@@ -1314,22 +1338,29 @@ void pfmergeCommand(client *c) {
|
||||
o = dbUnshareStringValue(c->db,c->argv[1],o);
|
||||
}
|
||||
|
||||
/* Only support dense objects as destination. */
|
||||
if (hllSparseToDense(o) == C_ERR) {
|
||||
/* Convert the destination object to dense representation if at least
|
||||
* one of the inputs was dense. */
|
||||
if (use_dense && hllSparseToDense(o) == C_ERR) {
|
||||
addReplySds(c,sdsnew(invalid_hll_err));
|
||||
return;
|
||||
}
|
||||
|
||||
/* Write the resulting HLL to the destination HLL registers and
|
||||
* invalidate the cached value. */
|
||||
hdr = o->ptr;
|
||||
for (j = 0; j < HLL_REGISTERS; j++) {
|
||||
HLL_DENSE_SET_REGISTER(hdr->registers,j,max[j]);
|
||||
if (max[j] == 0) continue;
|
||||
hdr = o->ptr;
|
||||
switch(hdr->encoding) {
|
||||
case HLL_DENSE: hllDenseSet(hdr->registers,j,max[j]); break;
|
||||
case HLL_SPARSE: hllSparseSet(o,j,max[j]); break;
|
||||
}
|
||||
}
|
||||
hdr = o->ptr; /* o->ptr may be different now, as a side effect of
|
||||
last hllSparseSet() call. */
|
||||
HLL_INVALIDATE_CACHE(hdr);
|
||||
|
||||
signalModifiedKey(c->db,c->argv[1]);
|
||||
/* We generate an PFADD event for PFMERGE for semantical simplicity
|
||||
/* We generate a PFADD event for PFMERGE for semantical simplicity
|
||||
* since in theory this is a mass-add of elements. */
|
||||
notifyKeyspaceEvent(NOTIFY_STRING,"pfadd",c->argv[1],c->db->id);
|
||||
server.dirty++;
|
||||
|
||||
+9
-3
@@ -64,9 +64,15 @@ int dbAsyncDelete(redisDb *db, robj *key) {
|
||||
robj *val = dictGetVal(de);
|
||||
size_t free_effort = lazyfreeGetFreeEffort(val);
|
||||
|
||||
/* If releasing the object is too much work, let's put it into the
|
||||
* lazy free list. */
|
||||
if (free_effort > LAZYFREE_THRESHOLD) {
|
||||
/* If releasing the object is too much work, do it in the background
|
||||
* by adding the object to the lazy free list.
|
||||
* Note that if the object is shared, to reclaim it now it is not
|
||||
* possible. This rarely happens, however sometimes the implementation
|
||||
* of parts of the Redis core may call incrRefCount() to protect
|
||||
* objects, and then call dbDelete(). In this case we'll fall
|
||||
* through and reach the dictFreeUnlinkedEntry() call, that will be
|
||||
* equivalent to just calling decrRefCount(). */
|
||||
if (free_effort > LAZYFREE_THRESHOLD && val->refcount == 1) {
|
||||
atomicIncr(lazyfree_objects,1);
|
||||
bioCreateBackgroundJob(BIO_LAZY_FREE,val,NULL,NULL);
|
||||
dictSetVal(db->dict,de,NULL);
|
||||
|
||||
+17
-2
@@ -1456,6 +1456,20 @@ int RM_DeleteKey(RedisModuleKey *key) {
|
||||
return REDISMODULE_OK;
|
||||
}
|
||||
|
||||
/* If the key is open for writing, unlink it (that is delete it in a
|
||||
* non-blocking way, not reclaiming memory immediately) and setup the key to
|
||||
* accept new writes as an empty key (that will be created on demand).
|
||||
* On success REDISMODULE_OK is returned. If the key is not open for
|
||||
* writing REDISMODULE_ERR is returned. */
|
||||
int RM_UnlinkKey(RedisModuleKey *key) {
|
||||
if (!(key->mode & REDISMODULE_WRITE)) return REDISMODULE_ERR;
|
||||
if (key->value) {
|
||||
dbAsyncDelete(key->db,key->key);
|
||||
key->value = NULL;
|
||||
}
|
||||
return REDISMODULE_OK;
|
||||
}
|
||||
|
||||
/* Return the key expire value, as milliseconds of remaining TTL.
|
||||
* If no TTL is associated with the key or if the key is empty,
|
||||
* REDISMODULE_NO_EXPIRE is returned. */
|
||||
@@ -3024,7 +3038,7 @@ int64_t RM_LoadSigned(RedisModuleIO *io) {
|
||||
void RM_SaveString(RedisModuleIO *io, RedisModuleString *s) {
|
||||
if (io->error) return;
|
||||
/* Save opcode. */
|
||||
int retval = rdbSaveLen(io->rio, RDB_MODULE_OPCODE_STRING);
|
||||
ssize_t retval = rdbSaveLen(io->rio, RDB_MODULE_OPCODE_STRING);
|
||||
if (retval == -1) goto saveerr;
|
||||
io->bytes += retval;
|
||||
/* Save value. */
|
||||
@@ -3042,7 +3056,7 @@ saveerr:
|
||||
void RM_SaveStringBuffer(RedisModuleIO *io, const char *str, size_t len) {
|
||||
if (io->error) return;
|
||||
/* Save opcode. */
|
||||
int retval = rdbSaveLen(io->rio, RDB_MODULE_OPCODE_STRING);
|
||||
ssize_t retval = rdbSaveLen(io->rio, RDB_MODULE_OPCODE_STRING);
|
||||
if (retval == -1) goto saveerr;
|
||||
io->bytes += retval;
|
||||
/* Save value. */
|
||||
@@ -3960,6 +3974,7 @@ void moduleRegisterCoreAPI(void) {
|
||||
REGISTER_API(Replicate);
|
||||
REGISTER_API(ReplicateVerbatim);
|
||||
REGISTER_API(DeleteKey);
|
||||
REGISTER_API(UnlinkKey);
|
||||
REGISTER_API(StringSet);
|
||||
REGISTER_API(StringDMA);
|
||||
REGISTER_API(StringTruncate);
|
||||
|
||||
@@ -120,6 +120,38 @@ int TestStringPrintf(RedisModuleCtx *ctx, RedisModuleString **argv, int argc) {
|
||||
return REDISMODULE_OK;
|
||||
}
|
||||
|
||||
int failTest(RedisModuleCtx *ctx, const char *msg) {
|
||||
RedisModule_ReplyWithError(ctx, msg);
|
||||
return REDISMODULE_ERR;
|
||||
}
|
||||
int TestUnlink(RedisModuleCtx *ctx, RedisModuleString **argv, int argc) {
|
||||
RedisModule_AutoMemory(ctx);
|
||||
REDISMODULE_NOT_USED(argv);
|
||||
REDISMODULE_NOT_USED(argc);
|
||||
|
||||
RedisModuleKey *k = RedisModule_OpenKey(ctx, RedisModule_CreateStringPrintf(ctx, "unlinked"), REDISMODULE_WRITE | REDISMODULE_READ);
|
||||
if (!k) return failTest(ctx, "Could not create key");
|
||||
|
||||
if (REDISMODULE_ERR == RedisModule_StringSet(k, RedisModule_CreateStringPrintf(ctx, "Foobar"))) {
|
||||
return failTest(ctx, "Could not set string value");
|
||||
}
|
||||
|
||||
RedisModuleCallReply *rep = RedisModule_Call(ctx, "EXISTS", "c", "unlinked");
|
||||
if (!rep || RedisModule_CallReplyInteger(rep) != 1) {
|
||||
return failTest(ctx, "Key does not exist before unlink");
|
||||
}
|
||||
|
||||
if (REDISMODULE_ERR == RedisModule_UnlinkKey(k)) {
|
||||
return failTest(ctx, "Could not unlink key");
|
||||
}
|
||||
|
||||
rep = RedisModule_Call(ctx, "EXISTS", "c", "unlinked");
|
||||
if (!rep || RedisModule_CallReplyInteger(rep) != 0) {
|
||||
return failTest(ctx, "Could not verify key to be unlinked");
|
||||
}
|
||||
return RedisModule_ReplyWithSimpleString(ctx, "OK");
|
||||
|
||||
}
|
||||
|
||||
/* TEST.CTXFLAGS -- Test GetContextFlags. */
|
||||
int TestCtxFlags(RedisModuleCtx *ctx, RedisModuleString **argv, int argc) {
|
||||
@@ -269,6 +301,9 @@ int TestIt(RedisModuleCtx *ctx, RedisModuleString **argv, int argc) {
|
||||
T("test.string.append","");
|
||||
if (!TestAssertStringReply(ctx,reply,"foobar",6)) goto fail;
|
||||
|
||||
T("test.unlink","");
|
||||
if (!TestAssertStringReply(ctx,reply,"OK",2)) goto fail;
|
||||
|
||||
T("test.string.append.am","");
|
||||
if (!TestAssertStringReply(ctx,reply,"foobar",6)) goto fail;
|
||||
|
||||
@@ -310,6 +345,10 @@ int RedisModule_OnLoad(RedisModuleCtx *ctx, RedisModuleString **argv, int argc)
|
||||
if (RedisModule_CreateCommand(ctx,"test.ctxflags",
|
||||
TestCtxFlags,"readonly",1,1,1) == REDISMODULE_ERR)
|
||||
return REDISMODULE_ERR;
|
||||
|
||||
if (RedisModule_CreateCommand(ctx,"test.unlink",
|
||||
TestUnlink,"write deny-oom",1,1,1) == REDISMODULE_ERR)
|
||||
return REDISMODULE_ERR;
|
||||
|
||||
if (RedisModule_CreateCommand(ctx,"test.it",
|
||||
TestIt,"readonly",1,1,1) == REDISMODULE_ERR)
|
||||
|
||||
+15
-9
@@ -33,7 +33,7 @@
|
||||
#include <math.h>
|
||||
#include <ctype.h>
|
||||
|
||||
static void setProtocolError(const char *errstr, client *c, int pos);
|
||||
static void setProtocolError(const char *errstr, client *c, long pos);
|
||||
|
||||
/* Return the size consumed from the allocator, for the specified SDS string,
|
||||
* including internal fragmentation. This function is used in order to compute
|
||||
@@ -939,10 +939,15 @@ int writeToClient(int fd, client *c, int handler_installed) {
|
||||
* scenario think about 'KEYS *' against the loopback interface).
|
||||
*
|
||||
* However if we are over the maxmemory limit we ignore that and
|
||||
* just deliver as much data as it is possible to deliver. */
|
||||
* just deliver as much data as it is possible to deliver.
|
||||
*
|
||||
* Moreover, we also send as much as possible if the client is
|
||||
* a slave (otherwise, on high-speed traffic, the replication
|
||||
* buffer will grow indefinitely) */
|
||||
if (totwritten > NET_MAX_WRITES_PER_EVENT &&
|
||||
(server.maxmemory == 0 ||
|
||||
zmalloc_used_memory() < server.maxmemory)) break;
|
||||
zmalloc_used_memory() < server.maxmemory) &&
|
||||
!(c->flags & CLIENT_SLAVE)) break;
|
||||
}
|
||||
server.stat_net_output_bytes += totwritten;
|
||||
if (nwritten == -1) {
|
||||
@@ -1107,7 +1112,7 @@ int processInlineBuffer(client *c) {
|
||||
/* Helper function. Trims query buffer to make the function that processes
|
||||
* multi bulk requests idempotent. */
|
||||
#define PROTO_DUMP_LEN 128
|
||||
static void setProtocolError(const char *errstr, client *c, int pos) {
|
||||
static void setProtocolError(const char *errstr, client *c, long pos) {
|
||||
if (server.verbosity <= LL_VERBOSE) {
|
||||
sds client = catClientInfoString(sdsempty(),c);
|
||||
|
||||
@@ -1148,7 +1153,8 @@ static void setProtocolError(const char *errstr, client *c, int pos) {
|
||||
* to be '*'. Otherwise for inline commands processInlineBuffer() is called. */
|
||||
int processMultibulkBuffer(client *c) {
|
||||
char *newline = NULL;
|
||||
int pos = 0, ok;
|
||||
long pos = 0;
|
||||
int ok;
|
||||
long long ll;
|
||||
|
||||
if (c->multibulklen == 0) {
|
||||
@@ -1220,7 +1226,7 @@ int processMultibulkBuffer(client *c) {
|
||||
}
|
||||
|
||||
ok = string2ll(c->querybuf+pos+1,newline-(c->querybuf+pos+1),&ll);
|
||||
if (!ok || ll < 0 || ll > 512*1024*1024) {
|
||||
if (!ok || ll < 0 || ll > server.proto_max_bulk_len) {
|
||||
addReplyError(c,"Protocol error: invalid bulk length");
|
||||
setProtocolError("invalid bulk length",c,pos);
|
||||
return C_ERR;
|
||||
@@ -1246,7 +1252,7 @@ int processMultibulkBuffer(client *c) {
|
||||
}
|
||||
|
||||
/* Read bulk argument */
|
||||
if (sdslen(c->querybuf)-pos < (unsigned)(c->bulklen+2)) {
|
||||
if (sdslen(c->querybuf)-pos < (size_t)(c->bulklen+2)) {
|
||||
/* Not enough data (+2 == trailing \r\n) */
|
||||
break;
|
||||
} else {
|
||||
@@ -1255,7 +1261,7 @@ int processMultibulkBuffer(client *c) {
|
||||
* just use the current sds string. */
|
||||
if (pos == 0 &&
|
||||
c->bulklen >= PROTO_MBULK_BIG_ARG &&
|
||||
(signed) sdslen(c->querybuf) == c->bulklen+2)
|
||||
sdslen(c->querybuf) == (size_t)(c->bulklen+2))
|
||||
{
|
||||
c->argv[c->argc++] = createObject(OBJ_STRING,c->querybuf);
|
||||
sdsIncrLen(c->querybuf,-2); /* remove CRLF */
|
||||
@@ -1366,7 +1372,7 @@ void readQueryFromClient(aeEventLoop *el, int fd, void *privdata, int mask) {
|
||||
if (c->reqtype == PROTO_REQ_MULTIBULK && c->multibulklen && c->bulklen != -1
|
||||
&& c->bulklen >= PROTO_MBULK_BIG_ARG)
|
||||
{
|
||||
int remaining = (unsigned)(c->bulklen+2)-sdslen(c->querybuf);
|
||||
ssize_t remaining = (size_t)(c->bulklen+2)-sdslen(c->querybuf);
|
||||
|
||||
if (remaining < readlen) readlen = remaining;
|
||||
}
|
||||
|
||||
+1
-1
@@ -727,7 +727,7 @@ size_t objectComputeSize(robj *o, size_t sample_size) {
|
||||
elesize += sizeof(quicklistNode)+ziplistBlobLen(node->zl);
|
||||
samples++;
|
||||
} while ((node = node->next) && samples < sample_size);
|
||||
asize += (double)elesize/samples*listTypeLength(o);
|
||||
asize += (double)elesize/samples*ql->len;
|
||||
} else if (o->encoding == OBJ_ENCODING_ZIPLIST) {
|
||||
asize = sizeof(*o)+ziplistBlobLen(o->ptr);
|
||||
} else {
|
||||
|
||||
+1
-1
@@ -149,7 +149,7 @@ REDIS_STATIC quicklistNode *quicklistCreateNode(void) {
|
||||
}
|
||||
|
||||
/* Return cached quicklist count */
|
||||
unsigned int quicklistCount(const quicklist *ql) { return ql->count; }
|
||||
unsigned long quicklistCount(const quicklist *ql) { return ql->count; }
|
||||
|
||||
/* Free entire quicklist. */
|
||||
void quicklistRelease(quicklist *quicklist) {
|
||||
|
||||
+3
-3
@@ -64,7 +64,7 @@ typedef struct quicklistLZF {
|
||||
char compressed[];
|
||||
} quicklistLZF;
|
||||
|
||||
/* quicklist is a 32 byte struct (on 64-bit systems) describing a quicklist.
|
||||
/* quicklist is a 40 byte struct (on 64-bit systems) describing a quicklist.
|
||||
* 'count' is the number of total entries.
|
||||
* 'len' is the number of quicklist nodes.
|
||||
* 'compress' is: -1 if compression disabled, otherwise it's the number
|
||||
@@ -74,7 +74,7 @@ typedef struct quicklist {
|
||||
quicklistNode *head;
|
||||
quicklistNode *tail;
|
||||
unsigned long count; /* total count of all entries in all ziplists */
|
||||
unsigned int len; /* number of quicklistNodes */
|
||||
unsigned long len; /* number of quicklistNodes */
|
||||
int fill : 16; /* fill factor for individual nodes */
|
||||
unsigned int compress : 16; /* depth of end nodes not to compress;0=off */
|
||||
} quicklist;
|
||||
@@ -154,7 +154,7 @@ int quicklistPopCustom(quicklist *quicklist, int where, unsigned char **data,
|
||||
void *(*saver)(unsigned char *data, unsigned int sz));
|
||||
int quicklistPop(quicklist *quicklist, int where, unsigned char **data,
|
||||
unsigned int *sz, long long *slong);
|
||||
unsigned int quicklistCount(const quicklist *ql);
|
||||
unsigned long quicklistCount(const quicklist *ql);
|
||||
int quicklistCompare(unsigned char *p1, unsigned char *p2, int p2_len);
|
||||
size_t quicklistGetLzf(const quicklistNode *node, void **data);
|
||||
|
||||
|
||||
@@ -131,7 +131,7 @@ static inline void raxStackFree(raxStack *ts) {
|
||||
}
|
||||
|
||||
/* ----------------------------------------------------------------------------
|
||||
* Radis tree implementation
|
||||
* Radix tree implementation
|
||||
* --------------------------------------------------------------------------*/
|
||||
|
||||
/* Allocate a new non compressed node with the specified number of children.
|
||||
@@ -186,10 +186,10 @@ raxNode *raxReallocForData(raxNode *n, void *data) {
|
||||
void raxSetData(raxNode *n, void *data) {
|
||||
n->iskey = 1;
|
||||
if (data != NULL) {
|
||||
n->isnull = 0;
|
||||
void **ndata = (void**)
|
||||
((char*)n+raxNodeCurrentLength(n)-sizeof(void*));
|
||||
memcpy(ndata,&data,sizeof(data));
|
||||
n->isnull = 0;
|
||||
} else {
|
||||
n->isnull = 1;
|
||||
}
|
||||
@@ -396,6 +396,7 @@ static inline size_t raxLowWalk(rax *rax, unsigned char *s, size_t len, raxNode
|
||||
position to 0 to signal this node represents
|
||||
the searched key. */
|
||||
}
|
||||
debugnode("Lookup stop node is",h);
|
||||
if (stopnode) *stopnode = h;
|
||||
if (plink) *plink = parentlink;
|
||||
if (splitpos && h->iscompr) *splitpos = j;
|
||||
@@ -424,18 +425,21 @@ int raxInsert(rax *rax, unsigned char *s, size_t len, void *data, void **old) {
|
||||
* our key. We have just to reallocate the node and make space for the
|
||||
* data pointer. */
|
||||
if (i == len && (!h->iscompr || j == 0 /* not in the middle if j is 0 */)) {
|
||||
debugf("### Insert: node representing key exists\n");
|
||||
if (!h->iskey || h->isnull) {
|
||||
h = raxReallocForData(h,data);
|
||||
if (h) memcpy(parentlink,&h,sizeof(h));
|
||||
}
|
||||
if (h == NULL) {
|
||||
errno = ENOMEM;
|
||||
return 0;
|
||||
}
|
||||
if (h->iskey) {
|
||||
if (old) *old = raxGetData(h);
|
||||
raxSetData(h,data);
|
||||
errno = 0;
|
||||
return 0; /* Element already exists. */
|
||||
}
|
||||
h = raxReallocForData(h,data);
|
||||
if (h == NULL) {
|
||||
errno = ENOMEM;
|
||||
return 0;
|
||||
}
|
||||
memcpy(parentlink,&h,sizeof(h));
|
||||
raxSetData(h,data);
|
||||
rax->numele++;
|
||||
return 1; /* Element inserted. */
|
||||
@@ -734,9 +738,7 @@ int raxInsert(rax *rax, unsigned char *s, size_t len, void *data, void **old) {
|
||||
}
|
||||
|
||||
/* We walked the radix tree as far as we could, but still there are left
|
||||
* chars in our string. We need to insert the missing nodes.
|
||||
* Note: while loop never entered if the node was split by ALGO2,
|
||||
* since i == len. */
|
||||
* chars in our string. We need to insert the missing nodes. */
|
||||
while(i < len) {
|
||||
raxNode *child;
|
||||
|
||||
@@ -871,7 +873,8 @@ raxNode *raxRemoveChild(raxNode *parent, raxNode *child) {
|
||||
memmove(((char*)cp)-1,cp,(parent->size-taillen-1)*sizeof(raxNode**));
|
||||
|
||||
/* Move the remaining "tail" pointer at the right position as well. */
|
||||
memmove(((char*)c)-1,c+1,taillen*sizeof(raxNode**)+parent->iskey*sizeof(void*));
|
||||
size_t valuelen = (parent->iskey && !parent->isnull) ? sizeof(void*) : 0;
|
||||
memmove(((char*)c)-1,c+1,taillen*sizeof(raxNode**)+valuelen);
|
||||
|
||||
/* 4. Update size. */
|
||||
parent->size--;
|
||||
@@ -1090,27 +1093,36 @@ int raxRemove(rax *rax, unsigned char *s, size_t len, void **old) {
|
||||
|
||||
/* This is the core of raxFree(): performs a depth-first scan of the
|
||||
* tree and releases all the nodes found. */
|
||||
void raxRecursiveFree(rax *rax, raxNode *n) {
|
||||
void raxRecursiveFree(rax *rax, raxNode *n, void (*free_callback)(void*)) {
|
||||
debugnode("free traversing",n);
|
||||
int numchildren = n->iscompr ? 1 : n->size;
|
||||
raxNode **cp = raxNodeLastChildPtr(n);
|
||||
while(numchildren--) {
|
||||
raxNode *child;
|
||||
memcpy(&child,cp,sizeof(child));
|
||||
raxRecursiveFree(rax,child);
|
||||
raxRecursiveFree(rax,child,free_callback);
|
||||
cp--;
|
||||
}
|
||||
debugnode("free depth-first",n);
|
||||
if (free_callback && n->iskey && !n->isnull)
|
||||
free_callback(raxGetData(n));
|
||||
rax_free(n);
|
||||
rax->numnodes--;
|
||||
}
|
||||
|
||||
/* Free a whole radix tree. */
|
||||
void raxFree(rax *rax) {
|
||||
raxRecursiveFree(rax,rax->head);
|
||||
/* Free a whole radix tree, calling the specified callback in order to
|
||||
* free the auxiliary data. */
|
||||
void raxFreeWithCallback(rax *rax, void (*free_callback)(void*)) {
|
||||
raxRecursiveFree(rax,rax->head,free_callback);
|
||||
assert(rax->numnodes == 0);
|
||||
rax_free(rax);
|
||||
}
|
||||
|
||||
/* Free a whole radix tree. */
|
||||
void raxFree(rax *rax) {
|
||||
raxFreeWithCallback(rax,NULL);
|
||||
}
|
||||
|
||||
/* ------------------------------- Iterator --------------------------------- */
|
||||
|
||||
/* Initialize a Rax iterator. This call should be performed a single time
|
||||
@@ -1172,7 +1184,7 @@ void raxIteratorDelChars(raxIterator *it, size_t count) {
|
||||
* The function returns 1 on success or 0 on out of memory. */
|
||||
int raxIteratorNextStep(raxIterator *it, int noup) {
|
||||
if (it->flags & RAX_ITER_EOF) {
|
||||
return 0;
|
||||
return 1;
|
||||
} else if (it->flags & RAX_ITER_JUST_SEEKED) {
|
||||
it->flags &= ~RAX_ITER_JUST_SEEKED;
|
||||
return 1;
|
||||
@@ -1184,10 +1196,6 @@ int raxIteratorNextStep(raxIterator *it, int noup) {
|
||||
size_t orig_stack_items = it->stack.items;
|
||||
raxNode *orig_node = it->node;
|
||||
|
||||
/* Clear the EOF flag: it will be set again if the EOF condition
|
||||
* is still valid. */
|
||||
it->flags &= ~RAX_ITER_EOF;
|
||||
|
||||
while(1) {
|
||||
int children = it->node->iscompr ? 1 : it->node->size;
|
||||
if (!noup && children) {
|
||||
@@ -1288,7 +1296,7 @@ int raxSeekGreatest(raxIterator *it) {
|
||||
* effect to the one of raxIteratorPrevSte(). */
|
||||
int raxIteratorPrevStep(raxIterator *it, int noup) {
|
||||
if (it->flags & RAX_ITER_EOF) {
|
||||
return 0;
|
||||
return 1;
|
||||
} else if (it->flags & RAX_ITER_JUST_SEEKED) {
|
||||
it->flags &= ~RAX_ITER_JUST_SEEKED;
|
||||
return 1;
|
||||
@@ -1409,6 +1417,7 @@ int raxSeek(raxIterator *it, const char *op, unsigned char *ele, size_t len) {
|
||||
it->node = it->rt->head;
|
||||
if (!raxSeekGreatest(it)) return 0;
|
||||
assert(it->node->iskey);
|
||||
it->data = raxGetData(it->node);
|
||||
return 1;
|
||||
}
|
||||
|
||||
@@ -1427,6 +1436,7 @@ int raxSeek(raxIterator *it, const char *op, unsigned char *ele, size_t len) {
|
||||
/* We found our node, since the key matches and we have an
|
||||
* "equal" condition. */
|
||||
if (!raxIteratorAddChars(it,ele,len)) return 0; /* OOM. */
|
||||
it->data = raxGetData(it->node);
|
||||
} else if (lt || gt) {
|
||||
/* Exact key not found or eq flag not set. We have to set as current
|
||||
* key the one represented by the node we stopped at, and perform
|
||||
@@ -1499,6 +1509,7 @@ int raxSeek(raxIterator *it, const char *op, unsigned char *ele, size_t len) {
|
||||
* the previous sub-tree. */
|
||||
if (nodechar < keychar) {
|
||||
if (!raxSeekGreatest(it)) return 0;
|
||||
it->data = raxGetData(it->node);
|
||||
} else {
|
||||
if (!raxIteratorAddChars(it,it->node->data,it->node->size))
|
||||
return 0;
|
||||
@@ -1615,8 +1626,8 @@ int raxCompare(raxIterator *iter, const char *op, unsigned char *key, size_t key
|
||||
int eq = 0, lt = 0, gt = 0;
|
||||
|
||||
if (op[0] == '=' || op[1] == '=') eq = 1;
|
||||
if (op[1] == '>') gt = 1;
|
||||
else if (op[1] == '<') lt = 1;
|
||||
if (op[0] == '>') gt = 1;
|
||||
else if (op[0] == '<') lt = 1;
|
||||
else if (op[1] != '=') return 0; /* Syntax error. */
|
||||
|
||||
size_t minlen = key_len < iter->key_len ? key_len : iter->key_len;
|
||||
@@ -1644,6 +1655,19 @@ void raxStop(raxIterator *it) {
|
||||
raxStackFree(&it->stack);
|
||||
}
|
||||
|
||||
/* Return if the iterator is in an EOF state. This happens when raxSeek()
|
||||
* failed to seek an appropriate element, so that raxNext() or raxPrev()
|
||||
* will return zero, or when an EOF condition was reached while iterating
|
||||
* with raxNext() and raxPrev(). */
|
||||
int raxEOF(raxIterator *it) {
|
||||
return it->flags & RAX_ITER_EOF;
|
||||
}
|
||||
|
||||
/* Return the number of elements inside the radix tree. */
|
||||
uint64_t raxSize(rax *rax) {
|
||||
return rax->numele;
|
||||
}
|
||||
|
||||
/* ----------------------------- Introspection ------------------------------ */
|
||||
|
||||
/* This function is mostly used for debugging and learning purposes.
|
||||
|
||||
@@ -148,6 +148,7 @@ int raxInsert(rax *rax, unsigned char *s, size_t len, void *data, void **old);
|
||||
int raxRemove(rax *rax, unsigned char *s, size_t len, void **old);
|
||||
void *raxFind(rax *rax, unsigned char *s, size_t len);
|
||||
void raxFree(rax *rax);
|
||||
void raxFreeWithCallback(rax *rax, void (*free_callback)(void*));
|
||||
void raxStart(raxIterator *it, rax *rt);
|
||||
int raxSeek(raxIterator *it, const char *op, unsigned char *ele, size_t len);
|
||||
int raxNext(raxIterator *it);
|
||||
@@ -155,6 +156,8 @@ int raxPrev(raxIterator *it);
|
||||
int raxRandomWalk(raxIterator *it, size_t steps);
|
||||
int raxCompare(raxIterator *iter, const char *op, unsigned char *key, size_t key_len);
|
||||
void raxStop(raxIterator *it);
|
||||
int raxEOF(raxIterator *it);
|
||||
void raxShow(rax *rax);
|
||||
uint64_t raxSize(rax *rax);
|
||||
|
||||
#endif
|
||||
|
||||
@@ -424,7 +424,7 @@ ssize_t rdbSaveLongLongAsStringObject(rio *rdb, long long value) {
|
||||
}
|
||||
|
||||
/* Like rdbSaveRawString() gets a Redis object instead. */
|
||||
int rdbSaveStringObject(rio *rdb, robj *obj) {
|
||||
ssize_t rdbSaveStringObject(rio *rdb, robj *obj) {
|
||||
/* Avoid to decode the object, then encode it again, if the
|
||||
* object is already integer encoded. */
|
||||
if (obj->encoding == OBJ_ENCODING_INT) {
|
||||
@@ -826,21 +826,25 @@ int rdbSaveKeyValuePair(rio *rdb, robj *key, robj *val,
|
||||
}
|
||||
|
||||
/* Save an AUX field. */
|
||||
int rdbSaveAuxField(rio *rdb, void *key, size_t keylen, void *val, size_t vallen) {
|
||||
if (rdbSaveType(rdb,RDB_OPCODE_AUX) == -1) return -1;
|
||||
if (rdbSaveRawString(rdb,key,keylen) == -1) return -1;
|
||||
if (rdbSaveRawString(rdb,val,vallen) == -1) return -1;
|
||||
return 1;
|
||||
ssize_t rdbSaveAuxField(rio *rdb, void *key, size_t keylen, void *val, size_t vallen) {
|
||||
ssize_t ret, len = 0;
|
||||
if ((ret = rdbSaveType(rdb,RDB_OPCODE_AUX)) == -1) return -1;
|
||||
len += ret;
|
||||
if ((ret = rdbSaveRawString(rdb,key,keylen) == -1)) return -1;
|
||||
len += ret;
|
||||
if ((ret = rdbSaveRawString(rdb,val,vallen) == -1)) return -1;
|
||||
len += ret;
|
||||
return len;
|
||||
}
|
||||
|
||||
/* Wrapper for rdbSaveAuxField() used when key/val length can be obtained
|
||||
* with strlen(). */
|
||||
int rdbSaveAuxFieldStrStr(rio *rdb, char *key, char *val) {
|
||||
ssize_t rdbSaveAuxFieldStrStr(rio *rdb, char *key, char *val) {
|
||||
return rdbSaveAuxField(rdb,key,strlen(key),val,strlen(val));
|
||||
}
|
||||
|
||||
/* Wrapper for strlen(key) + integer type (up to long long range). */
|
||||
int rdbSaveAuxFieldStrInt(rio *rdb, char *key, long long val) {
|
||||
ssize_t rdbSaveAuxFieldStrInt(rio *rdb, char *key, long long val) {
|
||||
char buf[LONG_STR_SIZE];
|
||||
int vlen = ll2string(buf,sizeof(buf),val);
|
||||
return rdbSaveAuxField(rdb,key,strlen(key),buf,vlen);
|
||||
@@ -1605,7 +1609,7 @@ int rdbLoadRio(rio *rdb, rdbSaveInfo *rsi) {
|
||||
if (rsi) rsi->repl_offset = strtoll(auxval->ptr,NULL,10);
|
||||
} else if (!strcasecmp(auxkey->ptr,"lua")) {
|
||||
/* Load the script back in memory. */
|
||||
if (luaCreateFunction(NULL,server.lua,NULL,auxval) == C_ERR) {
|
||||
if (luaCreateFunction(NULL,server.lua,auxval) == NULL) {
|
||||
rdbExitReportCorruptRDB(
|
||||
"Can't load Lua script from RDB file! "
|
||||
"BODY: %s", auxval->ptr);
|
||||
|
||||
@@ -139,7 +139,7 @@ robj *rdbLoadObject(int type, rio *rdb);
|
||||
void backgroundSaveDoneHandler(int exitcode, int bysignal);
|
||||
int rdbSaveKeyValuePair(rio *rdb, robj *key, robj *val, long long expiretime, long long now);
|
||||
robj *rdbLoadStringObject(rio *rdb);
|
||||
int rdbSaveStringObject(rio *rdb, robj *obj);
|
||||
ssize_t rdbSaveStringObject(rio *rdb, robj *obj);
|
||||
ssize_t rdbSaveRawString(rio *rdb, unsigned char *s, size_t len);
|
||||
void *rdbGenericLoadStringObject(rio *rdb, int flags, size_t *lenptr);
|
||||
int rdbSaveBinaryDoubleValue(rio *rdb, double val);
|
||||
|
||||
@@ -614,7 +614,7 @@ int showThroughput(struct aeEventLoop *eventLoop, long long id, void *clientData
|
||||
UNUSED(id);
|
||||
UNUSED(clientData);
|
||||
|
||||
if (config.liveclients == 0) {
|
||||
if (config.liveclients == 0 && config.requests_finished != config.requests) {
|
||||
fprintf(stderr,"All clients disconnected... aborting.\n");
|
||||
exit(1);
|
||||
}
|
||||
|
||||
+2
-1
@@ -1386,8 +1386,9 @@ static void repl(void) {
|
||||
/* Only use history and load the rc file when stdin is a tty. */
|
||||
if (isatty(fileno(stdin))) {
|
||||
historyfile = getDotfilePath(REDIS_CLI_HISTFILE_ENV,REDIS_CLI_HISTFILE_DEFAULT);
|
||||
//keep in-memory history always regardless if history file can be determined
|
||||
history = 1;
|
||||
if (historyfile != NULL) {
|
||||
history = 1;
|
||||
linenoiseHistoryLoad(historyfile);
|
||||
}
|
||||
cliLoadPreferences();
|
||||
|
||||
+131
-1
@@ -701,7 +701,12 @@ class RedisTrib
|
||||
|
||||
masters.each{|m| puts m}
|
||||
|
||||
# Alloc slots on masters
|
||||
# Rotating the list sometimes helps to get better initial
|
||||
# anti-affinity before the optimizer runs.
|
||||
interleaved.push interleaved.shift
|
||||
|
||||
# Alloc slots on masters. After interleaving to get just the first N
|
||||
# should be optimal. With slaves is more complex, see later...
|
||||
slots_per_node = ClusterHashSlots.to_f / masters_count
|
||||
first = 0
|
||||
cursor = 0.0
|
||||
@@ -769,6 +774,131 @@ class RedisTrib
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
optimize_anti_affinity
|
||||
end
|
||||
|
||||
def optimize_anti_affinity
|
||||
score,aux = get_anti_affinity_score
|
||||
return if score == 0
|
||||
|
||||
xputs ">>> Trying to optimize slaves allocation for anti-affinity"
|
||||
|
||||
maxiter = 500*@nodes.length # Effort is proportional to cluster size...
|
||||
while maxiter > 0
|
||||
score,offenders = get_anti_affinity_score
|
||||
break if score == 0 # Optimal anti affinity reached
|
||||
|
||||
# We'll try to randomly swap a slave's assigned master causing
|
||||
# an affinity problem with another random slave, to see if we
|
||||
# can improve the affinity.
|
||||
first = offenders.shuffle.first
|
||||
nodes = @nodes.select{|n| n != first && n.info[:replicate]}
|
||||
break if nodes.length == 0
|
||||
second = nodes.shuffle.first
|
||||
|
||||
first_master = first.info[:replicate]
|
||||
second_master = second.info[:replicate]
|
||||
first.set_as_replica(second_master)
|
||||
second.set_as_replica(first_master)
|
||||
|
||||
new_score,aux = get_anti_affinity_score
|
||||
# If the change actually makes thing worse, revert. Otherwise
|
||||
# leave as it is becuase the best solution may need a few
|
||||
# combined swaps.
|
||||
if new_score > score
|
||||
first.set_as_replica(first_master)
|
||||
second.set_as_replica(second_master)
|
||||
end
|
||||
|
||||
maxiter -= 1
|
||||
end
|
||||
|
||||
score,aux = get_anti_affinity_score
|
||||
if score == 0
|
||||
xputs "[OK] Perfect anti-affinity obtained!"
|
||||
elsif score >= 10000
|
||||
puts "[WARNING] Some slaves are in the same host as their master"
|
||||
else
|
||||
puts "[WARNING] Some slaves of the same master are in the same host"
|
||||
end
|
||||
end
|
||||
|
||||
# Return the anti-affinity score, which is a measure of the amount of
|
||||
# violations of anti-affinity in the current cluster layout, that is, how
|
||||
# badly the masters and slaves are distributed in the different IP
|
||||
# addresses so that slaves of the same master are not in the master
|
||||
# host and are also in different hosts.
|
||||
#
|
||||
# The score is calculated as follows:
|
||||
#
|
||||
# SAME_AS_MASTER = 10000 * each slave in the same IP of its master.
|
||||
# SAME_AS_SLAVE = 1 * each slave having the same IP as another slave
|
||||
# of the same master.
|
||||
# FINAL_SCORE = SAME_AS_MASTER + SAME_AS_SLAVE
|
||||
#
|
||||
# So a greater score means a worse anti-affinity level, while zero
|
||||
# means perfect anti-affinity.
|
||||
#
|
||||
# The anti affinity optimizator will try to get a score as low as
|
||||
# possible. Since we do not want to sacrifice the fact that slaves should
|
||||
# not be in the same host as the master, we assign 10000 times the score
|
||||
# to this violation, so that we'll optimize for the second factor only
|
||||
# if it does not impact the first one.
|
||||
#
|
||||
# The function returns two things: the above score, and the list of
|
||||
# offending slaves, so that the optimizer can try changing the
|
||||
# configuration of the slaves violating the anti-affinity goals.
|
||||
def get_anti_affinity_score
|
||||
score = 0
|
||||
offending = [] # List of offending slaves to return to the caller
|
||||
|
||||
# First, split nodes by host
|
||||
host_to_node = {}
|
||||
@nodes.each{|n|
|
||||
host = n.info[:host]
|
||||
host_to_node[host] = [] if host_to_node[host] == nil
|
||||
host_to_node[host] << n
|
||||
}
|
||||
|
||||
# Then, for each set of nodes in the same host, split by
|
||||
# related nodes (masters and slaves which are involved in
|
||||
# replication of each other)
|
||||
host_to_node.each{|host,nodes|
|
||||
related = {}
|
||||
nodes.each{|n|
|
||||
if !n.info[:replicate]
|
||||
name = n.info[:name]
|
||||
related[name] = [] if related[name] == nil
|
||||
related[name] << :m
|
||||
else
|
||||
name = n.info[:replicate]
|
||||
related[name] = [] if related[name] == nil
|
||||
related[name] << :s
|
||||
end
|
||||
}
|
||||
|
||||
# Now it's trivial to check, for each related group having the
|
||||
# same host, what is their local score.
|
||||
related.each{|id,types|
|
||||
next if types.length < 2
|
||||
types.sort! # Make sure :m if the first if any
|
||||
if types[0] == :m
|
||||
score += 10000 * (types.length-1)
|
||||
else
|
||||
score += 1 * types.length
|
||||
end
|
||||
|
||||
# Populate the list of offending nodes
|
||||
@nodes.each{|n|
|
||||
if n.info[:replicate] == id &&
|
||||
n.info[:host] == host
|
||||
offending << n
|
||||
end
|
||||
}
|
||||
}
|
||||
}
|
||||
return score,offending
|
||||
end
|
||||
|
||||
def flush_nodes_config
|
||||
|
||||
+3
-1
@@ -185,6 +185,7 @@ int REDISMODULE_API_FUNC(RedisModule_ReplicateVerbatim)(RedisModuleCtx *ctx);
|
||||
const char *REDISMODULE_API_FUNC(RedisModule_CallReplyStringPtr)(RedisModuleCallReply *reply, size_t *len);
|
||||
RedisModuleString *REDISMODULE_API_FUNC(RedisModule_CreateStringFromCallReply)(RedisModuleCallReply *reply);
|
||||
int REDISMODULE_API_FUNC(RedisModule_DeleteKey)(RedisModuleKey *key);
|
||||
int REDISMODULE_API_FUNC(RedisModule_UnlinkKey)(RedisModuleKey *key);
|
||||
int REDISMODULE_API_FUNC(RedisModule_StringSet)(RedisModuleKey *key, RedisModuleString *str);
|
||||
char *REDISMODULE_API_FUNC(RedisModule_StringDMA)(RedisModuleKey *key, size_t *len, int mode);
|
||||
int REDISMODULE_API_FUNC(RedisModule_StringTruncate)(RedisModuleKey *key, size_t newlen);
|
||||
@@ -306,6 +307,7 @@ static int RedisModule_Init(RedisModuleCtx *ctx, const char *name, int ver, int
|
||||
REDISMODULE_GET_API(Replicate);
|
||||
REDISMODULE_GET_API(ReplicateVerbatim);
|
||||
REDISMODULE_GET_API(DeleteKey);
|
||||
REDISMODULE_GET_API(UnlinkKey);
|
||||
REDISMODULE_GET_API(StringSet);
|
||||
REDISMODULE_GET_API(StringDMA);
|
||||
REDISMODULE_GET_API(StringTruncate);
|
||||
@@ -372,7 +374,7 @@ static int RedisModule_Init(RedisModuleCtx *ctx, const char *name, int ver, int
|
||||
REDISMODULE_GET_API(AbortBlock);
|
||||
#endif
|
||||
|
||||
if (RedisModule_IsModuleNameBusy(name)) return REDISMODULE_ERR;
|
||||
if (RedisModule_IsModuleNameBusy && RedisModule_IsModuleNameBusy(name)) return REDISMODULE_ERR;
|
||||
RedisModule_SetModuleAttribs(ctx,name,ver,apiver);
|
||||
return REDISMODULE_OK;
|
||||
}
|
||||
|
||||
+8
-2
@@ -1330,7 +1330,8 @@ char *sendSynchronousCommand(int flags, int fd, ...) {
|
||||
cmd = sdscat(cmd,arg);
|
||||
}
|
||||
cmd = sdscatlen(cmd,"\r\n",2);
|
||||
|
||||
va_end(ap);
|
||||
|
||||
/* Transfer command to the server. */
|
||||
if (syncWrite(fd,cmd,sdslen(cmd),server.repl_syncio_timeout*1000)
|
||||
== -1)
|
||||
@@ -1340,7 +1341,6 @@ char *sendSynchronousCommand(int flags, int fd, ...) {
|
||||
strerror(errno));
|
||||
}
|
||||
sdsfree(cmd);
|
||||
va_end(ap);
|
||||
}
|
||||
|
||||
/* Read the reply from the server. */
|
||||
@@ -1970,6 +1970,12 @@ void replicationUnsetMaster(void) {
|
||||
* with PSYNC version 2, there is no need for full resync after a
|
||||
* master switch. */
|
||||
server.slaveseldb = -1;
|
||||
|
||||
/* Once we turn from slave to master, we consider the starting time without
|
||||
* slaves (that is used to count the replication backlog time to live) as
|
||||
* starting from now. Otherwise the backlog will be freed after a
|
||||
* failover if slaves do not connect immediately. */
|
||||
server.repl_no_slaves_since = server.unixtime;
|
||||
}
|
||||
|
||||
/* This function is called when the slave lose the connection with the
|
||||
|
||||
@@ -310,7 +310,7 @@ void rioSetAutoSync(rio *r, off_t bytes) {
|
||||
* generating the Redis protocol for the Append Only File. */
|
||||
|
||||
/* Write multi bulk count in the format: "*<count>\r\n". */
|
||||
size_t rioWriteBulkCount(rio *r, char prefix, int count) {
|
||||
size_t rioWriteBulkCount(rio *r, char prefix, long count) {
|
||||
char cbuf[128];
|
||||
int clen;
|
||||
|
||||
|
||||
@@ -130,7 +130,7 @@ void rioInitWithFdset(rio *r, int *fds, int numfds);
|
||||
|
||||
void rioFreeFdset(rio *r);
|
||||
|
||||
size_t rioWriteBulkCount(rio *r, char prefix, int count);
|
||||
size_t rioWriteBulkCount(rio *r, char prefix, long count);
|
||||
size_t rioWriteBulkString(rio *r, const char *buf, size_t len);
|
||||
size_t rioWriteBulkLongLong(rio *r, long long l);
|
||||
size_t rioWriteBulkDouble(rio *r, double d);
|
||||
|
||||
+37
-45
@@ -1141,33 +1141,38 @@ int redis_math_randomseed (lua_State *L) {
|
||||
* EVAL and SCRIPT commands implementation
|
||||
* ------------------------------------------------------------------------- */
|
||||
|
||||
/* Define a lua function with the specified function name and body.
|
||||
* The function name musts be a 42 characters long string, since all the
|
||||
* functions we defined in the Lua context are in the form:
|
||||
/* Define a Lua function with the specified body.
|
||||
* The function name will be generated in the following form:
|
||||
*
|
||||
* f_<hex sha1 sum>
|
||||
*
|
||||
* If 'funcname' is NULL, the function name is created by the function
|
||||
* on the fly doing the SHA1 of the body, this means that passing the funcname
|
||||
* is just an optimization in case it's already at hand.
|
||||
*
|
||||
* The function increments the reference count of the 'body' object as a
|
||||
* side effect of a successful call.
|
||||
*
|
||||
* On success C_OK is returned, and nothing is left on the Lua stack.
|
||||
* On error C_ERR is returned and an appropriate error is set in the
|
||||
* client context. */
|
||||
int luaCreateFunction(client *c, lua_State *lua, char *funcname, robj *body) {
|
||||
sds funcdef = sdsempty();
|
||||
char fname[43];
|
||||
* On success a pointer to an SDS string representing the function SHA1 of the
|
||||
* just added function is returned (and will be valid until the next call
|
||||
* to scriptingReset() function), otherwise NULL is returned.
|
||||
*
|
||||
* The function handles the fact of being called with a script that already
|
||||
* exists, and in such a case, it behaves like in the success case.
|
||||
*
|
||||
* If 'c' is not NULL, on error the client is informed with an appropriate
|
||||
* error describing the nature of the problem and the Lua interpreter error. */
|
||||
sds luaCreateFunction(client *c, lua_State *lua, robj *body) {
|
||||
char funcname[43];
|
||||
dictEntry *de;
|
||||
|
||||
if (funcname == NULL) {
|
||||
fname[0] = 'f';
|
||||
fname[1] = '_';
|
||||
sha1hex(fname+2,body->ptr,sdslen(body->ptr));
|
||||
funcname = fname;
|
||||
funcname[0] = 'f';
|
||||
funcname[1] = '_';
|
||||
sha1hex(funcname+2,body->ptr,sdslen(body->ptr));
|
||||
|
||||
sds sha = sdsnewlen(funcname+2,40);
|
||||
if ((de = dictFind(server.lua_scripts,sha)) != NULL) {
|
||||
sdsfree(sha);
|
||||
return dictGetKey(de);
|
||||
}
|
||||
|
||||
sds funcdef = sdsempty();
|
||||
funcdef = sdscat(funcdef,"function ");
|
||||
funcdef = sdscatlen(funcdef,funcname,42);
|
||||
funcdef = sdscatlen(funcdef,"() ",3);
|
||||
@@ -1181,29 +1186,29 @@ int luaCreateFunction(client *c, lua_State *lua, char *funcname, robj *body) {
|
||||
lua_tostring(lua,-1));
|
||||
}
|
||||
lua_pop(lua,1);
|
||||
sdsfree(sha);
|
||||
sdsfree(funcdef);
|
||||
return C_ERR;
|
||||
return NULL;
|
||||
}
|
||||
sdsfree(funcdef);
|
||||
|
||||
if (lua_pcall(lua,0,0,0)) {
|
||||
if (c != NULL) {
|
||||
addReplyErrorFormat(c,"Error running script (new function): %s\n",
|
||||
lua_tostring(lua,-1));
|
||||
}
|
||||
lua_pop(lua,1);
|
||||
return C_ERR;
|
||||
sdsfree(sha);
|
||||
return NULL;
|
||||
}
|
||||
|
||||
/* We also save a SHA1 -> Original script map in a dictionary
|
||||
* so that we can replicate / write in the AOF all the
|
||||
* EVALSHA commands as EVAL using the original script. */
|
||||
{
|
||||
int retval = dictAdd(server.lua_scripts,
|
||||
sdsnewlen(funcname+2,40),body);
|
||||
serverAssertWithInfo(c ? c : server.lua_client,NULL,retval == DICT_OK);
|
||||
incrRefCount(body);
|
||||
}
|
||||
return C_OK;
|
||||
int retval = dictAdd(server.lua_scripts,sha,body);
|
||||
serverAssertWithInfo(c ? c : server.lua_client,NULL,retval == DICT_OK);
|
||||
incrRefCount(body);
|
||||
return sha;
|
||||
}
|
||||
|
||||
/* This is the Lua script "count" hook that we use to detect scripts timeout. */
|
||||
@@ -1302,10 +1307,10 @@ void evalGenericCommand(client *c, int evalsha) {
|
||||
addReply(c, shared.noscripterr);
|
||||
return;
|
||||
}
|
||||
if (luaCreateFunction(c,lua,funcname,c->argv[1]) == C_ERR) {
|
||||
if (luaCreateFunction(c,lua,c->argv[1]) == NULL) {
|
||||
lua_pop(lua,1); /* remove the error handler from the stack. */
|
||||
/* The error is sent to the client by luaCreateFunction()
|
||||
* itself when it returns C_ERR. */
|
||||
* itself when it returns NULL. */
|
||||
return;
|
||||
}
|
||||
/* Now the following is guaranteed to return non nil */
|
||||
@@ -1466,22 +1471,9 @@ void scriptCommand(client *c) {
|
||||
addReply(c,shared.czero);
|
||||
}
|
||||
} else if (c->argc == 3 && !strcasecmp(c->argv[1]->ptr,"load")) {
|
||||
char funcname[43];
|
||||
sds sha;
|
||||
|
||||
funcname[0] = 'f';
|
||||
funcname[1] = '_';
|
||||
sha1hex(funcname+2,c->argv[2]->ptr,sdslen(c->argv[2]->ptr));
|
||||
sha = sdsnewlen(funcname+2,40);
|
||||
if (dictFind(server.lua_scripts,sha) == NULL) {
|
||||
if (luaCreateFunction(c,server.lua,funcname,c->argv[2])
|
||||
== C_ERR) {
|
||||
sdsfree(sha);
|
||||
return;
|
||||
}
|
||||
}
|
||||
addReplyBulkCBuffer(c,funcname+2,40);
|
||||
sdsfree(sha);
|
||||
sds sha = luaCreateFunction(c,server.lua,c->argv[2]);
|
||||
if (sha == NULL) return; /* The error was sent by luaCreateFunction(). */
|
||||
addReplyBulkCBuffer(c,sha,40);
|
||||
forceCommandPropagation(c,PROPAGATE_REPL|PROPAGATE_AOF);
|
||||
} else if (c->argc == 2 && !strcasecmp(c->argv[1]->ptr,"kill")) {
|
||||
if (server.lua_caller == NULL) {
|
||||
|
||||
@@ -175,7 +175,7 @@ void sdsfree(sds s) {
|
||||
* the output will be "6" as the string was modified but the logical length
|
||||
* remains 6 bytes. */
|
||||
void sdsupdatelen(sds s) {
|
||||
int reallen = strlen(s);
|
||||
size_t reallen = strlen(s);
|
||||
sdssetlen(s, reallen);
|
||||
}
|
||||
|
||||
@@ -319,7 +319,7 @@ void *sdsAllocPtr(sds s) {
|
||||
* ... check for nread <= 0 and handle it ...
|
||||
* sdsIncrLen(s, nread);
|
||||
*/
|
||||
void sdsIncrLen(sds s, int incr) {
|
||||
void sdsIncrLen(sds s, ssize_t incr) {
|
||||
unsigned char flags = s[-1];
|
||||
size_t len;
|
||||
switch(flags&SDS_TYPE_MASK) {
|
||||
@@ -589,7 +589,7 @@ sds sdscatprintf(sds s, const char *fmt, ...) {
|
||||
sds sdscatfmt(sds s, char const *fmt, ...) {
|
||||
size_t initlen = sdslen(s);
|
||||
const char *f = fmt;
|
||||
int i;
|
||||
long i;
|
||||
va_list ap;
|
||||
|
||||
va_start(ap,fmt);
|
||||
@@ -721,7 +721,7 @@ sds sdstrim(sds s, const char *cset) {
|
||||
* s = sdsnew("Hello World");
|
||||
* sdsrange(s,1,-1); => "ello World"
|
||||
*/
|
||||
void sdsrange(sds s, int start, int end) {
|
||||
void sdsrange(sds s, ssize_t start, ssize_t end) {
|
||||
size_t newlen, len = sdslen(s);
|
||||
|
||||
if (len == 0) return;
|
||||
@@ -735,9 +735,9 @@ void sdsrange(sds s, int start, int end) {
|
||||
}
|
||||
newlen = (start > end) ? 0 : (end-start)+1;
|
||||
if (newlen != 0) {
|
||||
if (start >= (signed)len) {
|
||||
if (start >= (ssize_t)len) {
|
||||
newlen = 0;
|
||||
} else if (end >= (signed)len) {
|
||||
} else if (end >= (ssize_t)len) {
|
||||
end = len-1;
|
||||
newlen = (start > end) ? 0 : (end-start)+1;
|
||||
}
|
||||
@@ -751,14 +751,14 @@ void sdsrange(sds s, int start, int end) {
|
||||
|
||||
/* Apply tolower() to every character of the sds string 's'. */
|
||||
void sdstolower(sds s) {
|
||||
int len = sdslen(s), j;
|
||||
size_t len = sdslen(s), j;
|
||||
|
||||
for (j = 0; j < len; j++) s[j] = tolower(s[j]);
|
||||
}
|
||||
|
||||
/* Apply toupper() to every character of the sds string 's'. */
|
||||
void sdstoupper(sds s) {
|
||||
int len = sdslen(s), j;
|
||||
size_t len = sdslen(s), j;
|
||||
|
||||
for (j = 0; j < len; j++) s[j] = toupper(s[j]);
|
||||
}
|
||||
@@ -782,7 +782,7 @@ int sdscmp(const sds s1, const sds s2) {
|
||||
l2 = sdslen(s2);
|
||||
minlen = (l1 < l2) ? l1 : l2;
|
||||
cmp = memcmp(s1,s2,minlen);
|
||||
if (cmp == 0) return l1-l2;
|
||||
if (cmp == 0) return l1>l2? 1: (l1<l2? -1: 0);
|
||||
return cmp;
|
||||
}
|
||||
|
||||
@@ -802,8 +802,9 @@ int sdscmp(const sds s1, const sds s2) {
|
||||
* requires length arguments. sdssplit() is just the
|
||||
* same function but for zero-terminated strings.
|
||||
*/
|
||||
sds *sdssplitlen(const char *s, int len, const char *sep, int seplen, int *count) {
|
||||
int elements = 0, slots = 5, start = 0, j;
|
||||
sds *sdssplitlen(const char *s, ssize_t len, const char *sep, int seplen, int *count) {
|
||||
int elements = 0, slots = 5;
|
||||
long start = 0, j;
|
||||
sds *tokens;
|
||||
|
||||
if (seplen < 1 || len < 0) return NULL;
|
||||
|
||||
@@ -236,11 +236,11 @@ sds sdscatprintf(sds s, const char *fmt, ...);
|
||||
|
||||
sds sdscatfmt(sds s, char const *fmt, ...);
|
||||
sds sdstrim(sds s, const char *cset);
|
||||
void sdsrange(sds s, int start, int end);
|
||||
void sdsrange(sds s, ssize_t start, ssize_t end);
|
||||
void sdsupdatelen(sds s);
|
||||
void sdsclear(sds s);
|
||||
int sdscmp(const sds s1, const sds s2);
|
||||
sds *sdssplitlen(const char *s, int len, const char *sep, int seplen, int *count);
|
||||
sds *sdssplitlen(const char *s, ssize_t len, const char *sep, int seplen, int *count);
|
||||
void sdsfreesplitres(sds *tokens, int count);
|
||||
void sdstolower(sds s);
|
||||
void sdstoupper(sds s);
|
||||
@@ -253,7 +253,7 @@ sds sdsjoinsds(sds *argv, int argc, const char *sep, size_t seplen);
|
||||
|
||||
/* Low level functions exposed to the user API */
|
||||
sds sdsMakeRoomFor(sds s, size_t addlen);
|
||||
void sdsIncrLen(sds s, int incr);
|
||||
void sdsIncrLen(sds s, ssize_t incr);
|
||||
sds sdsRemoveFreeSpace(sds s);
|
||||
size_t sdsAllocSize(sds s);
|
||||
void *sdsAllocPtr(sds s);
|
||||
|
||||
+6
-4
@@ -125,7 +125,7 @@ volatile unsigned long lru_clock; /* Server global current LRU time. */
|
||||
* are not fast commands.
|
||||
*/
|
||||
struct redisCommand redisCommandTable[] = {
|
||||
{"module",moduleCommand,-2,"as",0,NULL,1,1,1,0,0},
|
||||
{"module",moduleCommand,-2,"as",0,NULL,0,0,0,0,0},
|
||||
{"get",getCommand,2,"rF",0,NULL,1,1,1,0,0},
|
||||
{"set",setCommand,-3,"wm",0,NULL,1,1,1,0,0},
|
||||
{"setnx",setnxCommand,3,"wmF",0,NULL,1,1,1,0,0},
|
||||
@@ -1079,7 +1079,7 @@ int serverCron(struct aeEventLoop *eventLoop, long long id, void *clientData) {
|
||||
}
|
||||
} else {
|
||||
/* If there is not a background saving/rewrite in progress check if
|
||||
* we have to save/rewrite now */
|
||||
* we have to save/rewrite now. */
|
||||
for (j = 0; j < server.saveparamslen; j++) {
|
||||
struct saveparam *sp = server.saveparams+j;
|
||||
|
||||
@@ -1102,8 +1102,9 @@ int serverCron(struct aeEventLoop *eventLoop, long long id, void *clientData) {
|
||||
}
|
||||
}
|
||||
|
||||
/* Trigger an AOF rewrite if needed */
|
||||
if (server.rdb_child_pid == -1 &&
|
||||
/* Trigger an AOF rewrite if needed. */
|
||||
if (server.aof_state == AOF_ON &&
|
||||
server.rdb_child_pid == -1 &&
|
||||
server.aof_child_pid == -1 &&
|
||||
server.aof_rewrite_perc &&
|
||||
server.aof_current_size > server.aof_rewrite_min_size)
|
||||
@@ -1372,6 +1373,7 @@ void initServerConfig(void) {
|
||||
server.active_defrag_threshold_upper = CONFIG_DEFAULT_DEFRAG_THRESHOLD_UPPER;
|
||||
server.active_defrag_cycle_min = CONFIG_DEFAULT_DEFRAG_CYCLE_MIN;
|
||||
server.active_defrag_cycle_max = CONFIG_DEFAULT_DEFRAG_CYCLE_MAX;
|
||||
server.proto_max_bulk_len = CONFIG_DEFAULT_PROTO_MAX_BULK_LEN;
|
||||
server.client_max_querybuf_len = PROTO_MAX_QUERYBUF_LEN;
|
||||
server.saveparams = NULL;
|
||||
server.loading = 0;
|
||||
|
||||
+3
-1
@@ -160,6 +160,7 @@ typedef long long mstime_t; /* millisecond time type. */
|
||||
#define CONFIG_DEFAULT_DEFRAG_IGNORE_BYTES (100<<20) /* don't defrag if frag overhead is below 100mb */
|
||||
#define CONFIG_DEFAULT_DEFRAG_CYCLE_MIN 25 /* 25% CPU min (at lower threshold) */
|
||||
#define CONFIG_DEFAULT_DEFRAG_CYCLE_MAX 75 /* 75% CPU max (at upper threshold) */
|
||||
#define CONFIG_DEFAULT_PROTO_MAX_BULK_LEN (512ll*1024*1024) /* Bulk request max size */
|
||||
|
||||
#define ACTIVE_EXPIRE_CYCLE_LOOKUPS_PER_LOOP 20 /* Loopkups per loop. */
|
||||
#define ACTIVE_EXPIRE_CYCLE_FAST_DURATION 1000 /* Microseconds */
|
||||
@@ -1120,6 +1121,7 @@ struct redisServer {
|
||||
int maxmemory_samples; /* Pricision of random sampling */
|
||||
int lfu_log_factor; /* LFU logarithmic counter factor. */
|
||||
int lfu_decay_time; /* LFU counter decay factor. */
|
||||
long long proto_max_bulk_len; /* Protocol bulk length maximum size. */
|
||||
/* Blocked clients */
|
||||
unsigned int bpop_blocked_clients; /* Number of clients blocked by lists */
|
||||
list *unblocked_clients; /* list of clients to unblock before next loop */
|
||||
@@ -1781,7 +1783,7 @@ void scriptingInit(int setup);
|
||||
int ldbRemoveChild(pid_t pid);
|
||||
void ldbKillForkedSessions(void);
|
||||
int ldbPendingChildren(void);
|
||||
int luaCreateFunction(client *c, lua_State *lua, char *funcname, robj *body);
|
||||
sds luaCreateFunction(client *c, lua_State *lua, robj *body);
|
||||
|
||||
/* Blocked clients */
|
||||
void processUnblockedClients(void);
|
||||
|
||||
+10
-4
@@ -407,7 +407,7 @@ void spopWithCountCommand(client *c) {
|
||||
/* Get the count argument */
|
||||
if (getLongFromObjectOrReply(c,c->argv[2],&l,NULL) != C_OK) return;
|
||||
if (l >= 0) {
|
||||
count = (unsigned) l;
|
||||
count = (unsigned long) l;
|
||||
} else {
|
||||
addReply(c,shared.outofrangeerr);
|
||||
return;
|
||||
@@ -626,7 +626,7 @@ void srandmemberWithCountCommand(client *c) {
|
||||
|
||||
if (getLongFromObjectOrReply(c,c->argv[2],&l,NULL) != C_OK) return;
|
||||
if (l >= 0) {
|
||||
count = (unsigned) l;
|
||||
count = (unsigned long) l;
|
||||
} else {
|
||||
/* A negative count means: return the same elements multiple times
|
||||
* (i.e. don't remove the extracted element after every extraction). */
|
||||
@@ -774,15 +774,21 @@ void srandmemberCommand(client *c) {
|
||||
}
|
||||
|
||||
int qsortCompareSetsByCardinality(const void *s1, const void *s2) {
|
||||
return setTypeSize(*(robj**)s1)-setTypeSize(*(robj**)s2);
|
||||
if (setTypeSize(*(robj**)s1) > setTypeSize(*(robj**)s2)) return 1;
|
||||
if (setTypeSize(*(robj**)s1) < setTypeSize(*(robj**)s2)) return -1;
|
||||
return 0;
|
||||
}
|
||||
|
||||
/* This is used by SDIFF and in this case we can receive NULL that should
|
||||
* be handled as empty sets. */
|
||||
int qsortCompareSetsByRevCardinality(const void *s1, const void *s2) {
|
||||
robj *o1 = *(robj**)s1, *o2 = *(robj**)s2;
|
||||
unsigned long first = o1 ? setTypeSize(o1) : 0;
|
||||
unsigned long second = o2 ? setTypeSize(o2) : 0;
|
||||
|
||||
return (o2 ? setTypeSize(o2) : 0) - (o1 ? setTypeSize(o1) : 0);
|
||||
if (first < second) return 1;
|
||||
if (first > second) return -1;
|
||||
return 0;
|
||||
}
|
||||
|
||||
void sinterGenericCommand(client *c, robj **setkeys,
|
||||
|
||||
+1
-1
@@ -84,7 +84,7 @@ int stringmatchlen(const char *pattern, int patternLen,
|
||||
}
|
||||
match = 0;
|
||||
while(1) {
|
||||
if (pattern[0] == '\\') {
|
||||
if (pattern[0] == '\\' && patternLen >= 2) {
|
||||
pattern++;
|
||||
patternLen--;
|
||||
if (pattern[0] == string[0])
|
||||
|
||||
+1
-1
@@ -1 +1 @@
|
||||
#define REDIS_VERSION "4.0.4"
|
||||
#define REDIS_VERSION "4.0.8"
|
||||
|
||||
+1
-1
@@ -440,7 +440,7 @@ unsigned int zipStorePrevEntryLength(unsigned char *p, unsigned int len) {
|
||||
if ((prevlensize) == 1) { \
|
||||
(prevlen) = (ptr)[0]; \
|
||||
} else if ((prevlensize) == 5) { \
|
||||
assert(sizeof((prevlensize)) == 4); \
|
||||
assert(sizeof((prevlen)) == 4); \
|
||||
memcpy(&(prevlen), ((char*)(ptr)) + 1, 4); \
|
||||
memrev32ifbe(&prevlen); \
|
||||
} \
|
||||
|
||||
@@ -2,9 +2,12 @@ start_server {tags {"repl"}} {
|
||||
start_server {} {
|
||||
test {First server should have role slave after SLAVEOF} {
|
||||
r -1 slaveof [srv 0 host] [srv 0 port]
|
||||
after 1000
|
||||
s -1 role
|
||||
} {slave}
|
||||
wait_for_condition 50 100 {
|
||||
[s -1 master_link_status] eq {up}
|
||||
} else {
|
||||
fail "Replication not started."
|
||||
}
|
||||
}
|
||||
|
||||
test {If min-slaves-to-write is honored, write is accepted} {
|
||||
r config set min-slaves-to-write 1
|
||||
|
||||
@@ -100,7 +100,6 @@ start_server {tags {"repl"}} {
|
||||
close $fd
|
||||
puts "Master - Slave inconsistency"
|
||||
puts "Run diff -u against /tmp/repldump*.txt for more info"
|
||||
|
||||
}
|
||||
|
||||
set old_digest [r debug digest]
|
||||
@@ -109,5 +108,27 @@ start_server {tags {"repl"}} {
|
||||
set new_digest [r debug digest]
|
||||
assert {$old_digest eq $new_digest}
|
||||
}
|
||||
|
||||
test {SLAVE can reload "lua" AUX RDB fields of duplicated scripts} {
|
||||
# Force a Slave full resynchronization
|
||||
r debug change-repl-id
|
||||
r -1 client kill type master
|
||||
|
||||
# Check that after a full resync the slave can still load
|
||||
# correctly the RDB file: such file will contain "lua" AUX
|
||||
# sections with scripts already in the memory of the master.
|
||||
|
||||
wait_for_condition 50 100 {
|
||||
[s -1 master_link_status] eq {up}
|
||||
} else {
|
||||
fail "Replication not started."
|
||||
}
|
||||
|
||||
wait_for_condition 50 100 {
|
||||
[r debug digest] eq [r -1 debug digest]
|
||||
} else {
|
||||
fail "DEBUG DIGEST mismatch after full SYNC with many scripts"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -308,4 +308,28 @@ start_server {tags {"dump"}} {
|
||||
}
|
||||
}
|
||||
|
||||
test {MIGRATE AUTH: correct and wrong password cases} {
|
||||
set first [srv 0 client]
|
||||
r del list
|
||||
r lpush list a b c d
|
||||
start_server {tags {"repl"}} {
|
||||
set second [srv 0 client]
|
||||
set second_host [srv 0 host]
|
||||
set second_port [srv 0 port]
|
||||
$second config set requirepass foobar
|
||||
$second auth foobar
|
||||
|
||||
assert {[$first exists list] == 1}
|
||||
assert {[$second exists list] == 0}
|
||||
set ret [r -1 migrate $second_host $second_port list 9 5000 AUTH foobar]
|
||||
assert {$ret eq {OK}}
|
||||
assert {[$second exists list] == 1}
|
||||
assert {[$second lrange list 0 -1] eq {d c b a}}
|
||||
|
||||
r -1 lpush list a b c d
|
||||
$second config set requirepass foobar2
|
||||
catch {r -1 migrate $second_host $second_port list 9 5000 AUTH foobar} err
|
||||
assert_match {*invalid password*} $err
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user