Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1980bcb7f6 | ||
|
|
57786b14e3 | ||
|
|
2211540d93 | ||
|
|
c85c84be6b | ||
|
|
85b2477099 | ||
|
|
a945e5c066 | ||
|
|
65a2e40ae4 | ||
|
|
d6c70f22cd | ||
|
|
012bcd49ee | ||
|
|
1198f7ceff | ||
|
|
cb2f001f52 | ||
|
|
8449227f55 | ||
|
|
eeac1d35d9 | ||
|
|
fb0441a8a2 | ||
|
|
0429db3c65 | ||
|
|
d06fbbdd54 | ||
|
|
ab3d3aca48 | ||
|
|
b7c7edf912 | ||
|
|
cbb6f45784 | ||
|
|
d766322e67 | ||
|
|
6544796ab4 | ||
|
|
e2355c192c | ||
|
|
22969a13a0 | ||
|
|
6b71f714bd | ||
|
|
2090052ef2 | ||
|
|
a75f2025f2 | ||
|
|
76aab08f1f | ||
|
|
b6fe5074e1 | ||
|
|
eda5cb0a04 | ||
|
|
4a60fbd8b5 | ||
|
|
060eb3b2d0 | ||
|
|
3c942b1269 | ||
|
|
6b6a83c7ab | ||
|
|
048097ada8 | ||
|
|
906134fe52 | ||
|
|
03657e88fe | ||
|
|
52fda0132d | ||
|
|
15bc8e97a9 | ||
|
|
f30454c1f6 | ||
|
|
1e7227f429 | ||
|
|
9524fce0ac | ||
|
|
2a27da1cf5 | ||
|
|
e0c2a0ecfd | ||
|
|
2eca8aed14 | ||
|
|
3594238350 | ||
|
|
be1b9ee0d5 | ||
|
|
9f69e179cf | ||
|
|
0205dd01e2 | ||
|
|
3cce566eb1 | ||
|
|
d01f163ce0 | ||
|
|
9a3e15c6a2 | ||
|
|
fa87879bab | ||
|
|
bc7076b090 | ||
|
|
7675b00a64 | ||
|
|
f31d9b12fd | ||
|
|
897d857115 | ||
|
|
1ee6af4d87 | ||
|
|
1740300f35 | ||
|
|
b25c245156 | ||
|
|
1847b987c6 | ||
|
|
c94cd1d831 | ||
|
|
193e4acc07 | ||
|
|
d131921c41 | ||
|
|
2e71edccfc | ||
|
|
44053df0a4 | ||
|
|
1c60b7a671 | ||
|
|
368124e8fa | ||
|
|
79567b6e66 | ||
|
|
f119464929 | ||
|
|
097a555677 | ||
|
|
f1a2cbfd6e | ||
|
|
0c0b77d149 | ||
|
|
fa6bd1b230 | ||
|
|
ad0ddcf390 | ||
|
|
8651e5d50d | ||
|
|
f2b2897f80 | ||
|
|
363be78397 |
+382
@@ -10,6 +10,382 @@ 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.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
|
||||
================================================================================
|
||||
|
||||
Upgrade urgency CRITICAL: Several PSYNC2 bugs can corrupt the slave data set
|
||||
after a restart and a successful PSYNC2 handshake.
|
||||
|
||||
This is a quick followup to Redis 4.0.3 since I forgot to add a few fixes...
|
||||
that are actually described in the 4.0.3 changelog (but not in the list of
|
||||
commits). Basically it's the following commits, implementing the ability
|
||||
to persist scripts into RDB files for a successful PSYNC, otherwise a corruption
|
||||
could happen when a slave is restarted and receives EVALSHA from the master
|
||||
about scripts it does not know:
|
||||
|
||||
8449227f PSYNC2: Fix off by one buffer size in luaCreateFunction().
|
||||
eeac1d35 PSYNC2: just store script bodies into RDB.
|
||||
fb0441a8 PSYNC2: luaCreateFunction() should handle NULL client parameter.
|
||||
0429db3c PSYNC2: Save Lua scripts state into RDB file.
|
||||
d06fbbdd Regression test: Slave restart with EVALSHA in backlog issue #4483.
|
||||
ab3d3aca Prevent corruption of server.executable after DEBUG RESTART.
|
||||
b7c7edf9 Be more verbose when DEBUG RESTART fails.
|
||||
|
||||
Please upgrade ASAP to 4.0.4 becuase 4.0.3 had an incomplete set of fixes.
|
||||
|
||||
Cheers and sorry for the 4.0.3 fiasco ;-)
|
||||
Salvatore
|
||||
|
||||
================================================================================
|
||||
Redis 4.0.3 Released Thu Nov 30 13:14:50 CET 2017
|
||||
================================================================================
|
||||
|
||||
Upgrade urgency CRITICAL: Several PSYNC2 bugs can corrupt the slave data set
|
||||
after a restart and a successful PSYNC2 handshake.
|
||||
|
||||
Hi all, Redis 4.0.3 contains several bug fixes to different parts of Redis 4.0,
|
||||
but the highlight is definitely in the "PSYNC after restart" that the new
|
||||
RDB format, containing replication metadata information, was able to provide
|
||||
to Redis 4.0. There were several bugs that are addressed in this release.
|
||||
Moreover several LFU fixes improve the ability of Redis to correctly estimate
|
||||
the popularity of keys. This release also fixes important bugs in Redis modules,
|
||||
including bugs related to replication of modules commands, reloading the same
|
||||
module multiple times, and other related things. Finally there is even a
|
||||
security fix related to loading a corrupted Cluster state from a corrupted
|
||||
file. We advice to upgrade ASAP. Check the list of commits for credits, several
|
||||
people helped a lot in this release. I'm grateful to each of them.
|
||||
|
||||
Cheers,
|
||||
Salvatore
|
||||
|
||||
antirez in commit d766322e:
|
||||
LFU: Fix LFUDecrAndReturn() to just decrement.
|
||||
1 file changed, 3 insertions(+), 13 deletions(-)
|
||||
|
||||
zhaozhao.zz in commit 6544796a:
|
||||
LFU: add hotkeys option to redis-cli
|
||||
1 file changed, 135 insertions(+)
|
||||
|
||||
zhaozhao.zz in commit e2355c19:
|
||||
LFU: do some changes about LFU to find hotkeys
|
||||
4 files changed, 39 insertions(+), 19 deletions(-)
|
||||
|
||||
zhaozhao.zz in commit 22969a13:
|
||||
LFU: change lfu* parameters to int
|
||||
2 files changed, 3 insertions(+), 3 deletions(-)
|
||||
|
||||
zhaozhao.zz in commit 6b71f714:
|
||||
LFU: fix the missing of config get and rewrite
|
||||
1 file changed, 6 insertions(+), 2 deletions(-)
|
||||
|
||||
Felix Krause in commit 2090052e:
|
||||
Update link to https and use inline link
|
||||
1 file changed, 1 insertion(+), 1 deletion(-)
|
||||
|
||||
Bo Cai in commit a75f2025:
|
||||
redis-cli.c typo: Requets -> Requests.
|
||||
1 file changed, 1 insertion(+), 1 deletion(-)
|
||||
|
||||
Bo Cai in commit 76aab08f:
|
||||
redis-cli.c typo: helpe -> helper.
|
||||
1 file changed, 1 insertion(+), 1 deletion(-)
|
||||
|
||||
Sébastien Fievet in commit b6fe5074:
|
||||
Fix some typos
|
||||
1 file changed, 3 insertions(+), 3 deletions(-)
|
||||
|
||||
antirez in commit eda5cb0a:
|
||||
t_hash.c: clarify calling two times the same function.
|
||||
1 file changed, 2 insertions(+), 2 deletions(-)
|
||||
|
||||
antirez in commit 4a60fbd8:
|
||||
adlist: fix listJoin() in the case the second list is empty.
|
||||
1 file changed, 1 insertion(+), 1 deletion(-)
|
||||
|
||||
Chris Lamb in commit 060eb3b2:
|
||||
Correct spelling of "faield".
|
||||
1 file changed, 1 insertion(+), 1 deletion(-)
|
||||
|
||||
antirez in commit 3c942b12:
|
||||
Improve OBJECT HELP descriptions.
|
||||
1 file changed, 2 insertions(+), 2 deletions(-)
|
||||
|
||||
antirez in commit 6b6a83c7:
|
||||
Fix entry command table entry for OBJECT for HELP option.
|
||||
1 file changed, 1 insertion(+), 1 deletion(-)
|
||||
|
||||
Itamar Haber in commit 048097ad:
|
||||
Adds `OBJECT help`
|
||||
1 file changed, 18 insertions(+), 3 deletions(-)
|
||||
|
||||
David Carlier in commit 906134fe:
|
||||
Fix undefined behavior constant defined.
|
||||
2 files changed, 10 insertions(+), 2 deletions(-)
|
||||
|
||||
rouzier in commit 03657e88:
|
||||
Fix file descriptor leak and error handling
|
||||
1 file changed, 6 insertions(+), 3 deletions(-)
|
||||
|
||||
Itamar Haber in commit 52fda013:
|
||||
Prevents `OBJECT freq` with `noeviction`
|
||||
1 file changed, 2 insertions(+), 2 deletions(-)
|
||||
|
||||
Itamar Haber in commit 15bc8e97:
|
||||
Adds -u <uri> option to redis-cli.
|
||||
1 file changed, 89 insertions(+)
|
||||
|
||||
antirez in commit f30454c1:
|
||||
Test: regression test for latency expire events logging bug.
|
||||
1 file changed, 14 insertions(+)
|
||||
|
||||
zhaozhao.zz in commit 1e7227f4:
|
||||
expire & latency: fix the missing latency records generated by expire
|
||||
1 file changed, 11 insertions(+), 8 deletions(-)
|
||||
|
||||
antirez in commit 9524fce0:
|
||||
Modules: fix memory leak in RM_IsModuleNameBusy().
|
||||
1 file changed, 3 insertions(+), 7 deletions(-)
|
||||
|
||||
antirez in commit 2a27da1c:
|
||||
PSYNC2: reorganize comments related to recent fixes.
|
||||
2 files changed, 24 insertions(+), 26 deletions(-)
|
||||
|
||||
zhaozhao.zz in commit e0c2a0ec:
|
||||
PSYNC2: persist cached_master's dbid inside the RDB
|
||||
1 file changed, 16 insertions(+), 2 deletions(-)
|
||||
|
||||
zhaozhao.zz in commit 2eca8aed:
|
||||
PSYNC2: make repl_stream_db never be -1
|
||||
1 file changed, 6 insertions(+), 9 deletions(-)
|
||||
|
||||
zhaozhao.zz in commit 35942383:
|
||||
PSYNC2: clarify the scenario when repl_stream_db can be -1
|
||||
2 files changed, 21 insertions(+), 9 deletions(-)
|
||||
|
||||
zhaozhao.zz in commit be1b9ee0:
|
||||
PSYNC2 & RDB: fix the missing rdbSaveInfo for BGSAVE
|
||||
1 file changed, 4 insertions(+), 1 deletion(-)
|
||||
|
||||
zhaozhao.zz in commit 9f69e179:
|
||||
PSYNC2: safe free backlog when reach the time limit
|
||||
1 file changed, 12 insertions(+)
|
||||
|
||||
zhaozhao.zz in commit 0205dd01:
|
||||
Modules: handle the busy module name
|
||||
2 files changed, 19 insertions(+), 2 deletions(-)
|
||||
|
||||
zhaozhao.zz in commit 3cce566e:
|
||||
Modules: handle the conflict of registering commands
|
||||
1 file changed, 28 insertions(+), 21 deletions(-)
|
||||
|
||||
Oran Agra in commit d01f163c:
|
||||
fix string to double conversion, stopped parsing on \0 even if the string has more data.
|
||||
2 files changed, 9 insertions(+), 2 deletions(-)
|
||||
|
||||
antirez in commit 9a3e15c6:
|
||||
Modules: fix for scripting replication of modules commands.
|
||||
2 files changed, 9 insertions(+), 7 deletions(-)
|
||||
|
||||
Yossi Gottlieb in commit fa87879b:
|
||||
Nested MULTI/EXEC may replicate in different cases.
|
||||
2 files changed, 10 insertions(+)
|
||||
|
||||
zhaozhao.zz in commit bc7076b0:
|
||||
rehash: handle one db until finished
|
||||
1 file changed, 5 insertions(+), 2 deletions(-)
|
||||
|
||||
kmiku7 in commit 7675b00a:
|
||||
fix boundary case for _dictNextPower
|
||||
1 file changed, 1 insertion(+), 1 deletion(-)
|
||||
|
||||
Itamar Haber in commit f31d9b12:
|
||||
Fixes an off-by-one in argument handling of `MEMORY USAGE`
|
||||
1 file changed, 1 insertion(+), 1 deletion(-)
|
||||
|
||||
antirez in commit 897d8571:
|
||||
SDS: improve sdsRemoveFreeSpace() to avoid useless data copy.
|
||||
1 file changed, 12 insertions(+), 5 deletions(-)
|
||||
|
||||
antirez in commit 1ee6af4d:
|
||||
Fix saving of zero-length lists.
|
||||
1 file changed, 3 insertions(+), 2 deletions(-)
|
||||
|
||||
antirez in commit 1740300f:
|
||||
Fix buffer overflows occurring reading redis.conf.
|
||||
1 file changed, 3 insertions(+)
|
||||
|
||||
antirez in commit b25c2451:
|
||||
Regression test for issue #4391.
|
||||
1 file changed, 4 insertions(+)
|
||||
|
||||
antirez in commit 1847b987:
|
||||
More robust object -> double conversion.
|
||||
1 file changed, 8 insertions(+), 4 deletions(-)
|
||||
|
||||
antirez in commit c94cd1d8:
|
||||
Limit statement in RM_BlockClient() to 80 cols.
|
||||
1 file changed, 5 insertions(+), 4 deletions(-)
|
||||
|
||||
Dvir Volk in commit 193e4acc:
|
||||
Added safety net preventing redis from crashing if a module decide to block in MULTI
|
||||
1 file changed, 8 insertions(+), 5 deletions(-)
|
||||
|
||||
Dvir Volk in commit d131921c:
|
||||
Renamed GetCtxFlags to GetContextFlags
|
||||
3 files changed, 11 insertions(+), 11 deletions(-)
|
||||
|
||||
Dvir Volk in commit 2e71edcc:
|
||||
Added support for module context flags with RM_GetCtxFlags
|
||||
3 files changed, 177 insertions(+)
|
||||
|
||||
================================================================================
|
||||
Redis 4.0.2 Released Thu Sep 21 15:47:53 CEST 2017
|
||||
================================================================================
|
||||
|
||||
Upgrade urgency HIGH: Several potentially critical bugs fixed.
|
||||
|
||||
Hello, this release addresses several significant bugs in Redis 4.0:
|
||||
|
||||
1. A number of bugs were fixed in the area of PSYNC2 replication in the
|
||||
specific area of restarting an instance with an RDB file having the
|
||||
repliacation meta-data to continue without a full resynchronization. The
|
||||
old code allowed several inconsistencies under certain conditions, like
|
||||
starting a master with an RDB file generated by a slave, and later using
|
||||
such master to connect previous slaves having the same replication history.
|
||||
Because of other bugs, sometimes the replication resulted in a full
|
||||
synchronization even if actually a partial resynchronization was possible
|
||||
and so forth. Several commits by different authors fix different bugs here.
|
||||
|
||||
2. AOF flush on SHUTDOWN did not cared to really write the AOF buffers
|
||||
(not in the kernel but in the Redis process memory) to disk before exiting.
|
||||
Calling SHUTDOWN during traffic resulted into not every operation to be
|
||||
persisted on disk.
|
||||
|
||||
3. The SLOWLOG could reference values inside string objects stored at keys,
|
||||
creating a race condition during FLUSHALL ASYNC while the DB is reclaimed
|
||||
in another thread.
|
||||
|
||||
There are other smaller bugs addessed in this relase, see the full commit
|
||||
history below for more information.
|
||||
|
||||
A big thank you to all the contributors of this release. Without the
|
||||
help I received, Redis 4.0 would take a much longer time to mature. It's
|
||||
a real pleasure to work together with people around the world, while making
|
||||
Redis better.
|
||||
|
||||
antirez in commit 1c60b7a6:
|
||||
Clarify comment in change fixing #4323.
|
||||
1 file changed, 6 insertions(+), 2 deletions(-)
|
||||
|
||||
zhaozhao.zz in commit 368124e8:
|
||||
Lazyfree: avoid memory leak when free slowlog entry
|
||||
1 file changed, 5 insertions(+), 2 deletions(-)
|
||||
|
||||
antirez in commit 79567b6e:
|
||||
PSYNC2: More refinements related to #4316.
|
||||
2 files changed, 14 insertions(+), 11 deletions(-)
|
||||
|
||||
zhaozhao.zz in commit f1194649:
|
||||
PSYNC2: make persisiting replication info more solid
|
||||
4 files changed, 33 insertions(+), 9 deletions(-)
|
||||
|
||||
antirez in commit 097a5556:
|
||||
PSYNC2: Fix the way replication info is saved/loaded from RDB.
|
||||
4 files changed, 49 insertions(+), 23 deletions(-)
|
||||
|
||||
antirez in commit f1a2cbfd:
|
||||
PSYNC2: Create backlog on slave partial sync as well.
|
||||
1 file changed, 5 insertions(+)
|
||||
|
||||
antirez in commit 0c0b77d1:
|
||||
Add MEMORY DOCTOR to MEMORY HELP.
|
||||
1 file changed, 3 insertions(+), 1 deletion(-)
|
||||
|
||||
Mota in commit fa6bd1b2:
|
||||
redis-benchmark: default value size usage update.
|
||||
1 file changed, 2 insertions(+), 2 deletions(-)
|
||||
|
||||
jybaek in commit ad0ddcf3:
|
||||
Remove Duplicate Processing
|
||||
1 file changed, 1 deletion(-)
|
||||
|
||||
Oran Agra (and also Buğra Gedik) in commit 8651e5d5:
|
||||
Flush append only buffers before existing.
|
||||
1 file changed, 2 insertions(+), 1 deletion(-)
|
||||
|
||||
antirez in commit f2b2897f:
|
||||
Changelog: note that 4.0 CLUSTER NODES output changed.
|
||||
1 file changed, 6 insertions(+)
|
||||
|
||||
Itamar Haber in commit 363be783:
|
||||
Changes command stats iteration to being dict-based
|
||||
1 file changed, 17 insertions(+), 10 deletions(-)
|
||||
|
||||
================================================================================
|
||||
Redis 4.0.1 Released Mon Jul 24 15:51:31 CEST 2017
|
||||
================================================================================
|
||||
@@ -3719,6 +4095,12 @@ non-backward compatible changes introduced in the 4.0 release:
|
||||
Cluster. SO in order to upgrade a Redis Cluster to 4.0, a mass restart of
|
||||
all the instances is needed.
|
||||
|
||||
* Redis Cluster CLUSTER NODES output is now slightly different. Nodes
|
||||
addresses are now in the form host:port@bus-port instead of host:port.
|
||||
Clients should use CLUSTER SLOTS in order to fetch the cluster configuration
|
||||
however if they are still using CLUSTER NODES, they should be modified in
|
||||
order to ignore the @bus-port part.
|
||||
|
||||
* Writable slaves do not propagate writes to their sub-slaves, so writes to
|
||||
writable slaves remain just local.
|
||||
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
This README is just a fast *quick start* document. You can find more detailed documentation at http://redis.io.
|
||||
This README is just a fast *quick start* document. You can find more detailed documentation at [redis.io](https://redis.io).
|
||||
|
||||
What is Redis?
|
||||
--------------
|
||||
|
||||
+3
-3
@@ -606,7 +606,7 @@ slave-priority 100
|
||||
# deletion of the object. It means that the server stops processing new commands
|
||||
# in order to reclaim all the memory associated with an object in a synchronous
|
||||
# way. If the key deleted is associated with a small object, the time needed
|
||||
# in order to execute th DEL command is very small and comparable to most other
|
||||
# in order to execute the DEL command is very small and comparable to most other
|
||||
# O(1) or O(log_N) commands in Redis. However if the key is associated with an
|
||||
# aggregated value containing millions of elements, the server can block for
|
||||
# a long time (even seconds) in order to complete the operation.
|
||||
@@ -621,7 +621,7 @@ slave-priority 100
|
||||
# It's up to the design of the application to understand when it is a good
|
||||
# idea to use one or the other. However the Redis server sometimes has to
|
||||
# delete keys or flush the whole database as a side effect of other operations.
|
||||
# Specifically Redis deletes objects independently of an user call in the
|
||||
# Specifically Redis deletes objects independently of a user call in the
|
||||
# following scenarios:
|
||||
#
|
||||
# 1) On eviction, because of the maxmemory and maxmemory policy configurations,
|
||||
@@ -914,7 +914,7 @@ lua-time-limit 5000
|
||||
# Docker and other containers).
|
||||
#
|
||||
# In order to make Redis Cluster working in such environments, a static
|
||||
# configuration where each node known its public address is needed. The
|
||||
# configuration where each node knows its public address is needed. The
|
||||
# following two options are used for this scope, and are:
|
||||
#
|
||||
# * cluster-announce-ip
|
||||
|
||||
+1
-1
@@ -353,7 +353,7 @@ void listJoin(list *l, list *o) {
|
||||
else
|
||||
l->head = o->head;
|
||||
|
||||
l->tail = o->tail;
|
||||
if (o->tail) l->tail = o->tail;
|
||||
l->len += o->len;
|
||||
|
||||
/* Setup other as an empty list. */
|
||||
|
||||
@@ -243,6 +243,7 @@ int clusterLoadConfig(char *filename) {
|
||||
*p = '\0';
|
||||
direction = p[1]; /* Either '>' or '<' */
|
||||
slot = atoi(argv[j]+1);
|
||||
if (slot < 0 || slot >= CLUSTER_SLOTS) goto fmterr;
|
||||
p += 3;
|
||||
cn = clusterLookupNode(p);
|
||||
if (!cn) {
|
||||
@@ -262,6 +263,8 @@ int clusterLoadConfig(char *filename) {
|
||||
} else {
|
||||
start = stop = atoi(argv[j]);
|
||||
}
|
||||
if (start < 0 || start >= CLUSTER_SLOTS) goto fmterr;
|
||||
if (stop < 0 || stop >= CLUSTER_SLOTS) goto fmterr;
|
||||
while(start <= stop) clusterAddSlot(n, start++);
|
||||
}
|
||||
|
||||
|
||||
+6
-2
@@ -330,13 +330,13 @@ void loadServerConfigFromString(char *config) {
|
||||
}
|
||||
} else if (!strcasecmp(argv[0],"lfu-log-factor") && argc == 2) {
|
||||
server.lfu_log_factor = atoi(argv[1]);
|
||||
if (server.maxmemory_samples < 0) {
|
||||
if (server.lfu_log_factor < 0) {
|
||||
err = "lfu-log-factor must be 0 or greater";
|
||||
goto loaderr;
|
||||
}
|
||||
} else if (!strcasecmp(argv[0],"lfu-decay-time") && argc == 2) {
|
||||
server.lfu_decay_time = atoi(argv[1]);
|
||||
if (server.maxmemory_samples < 1) {
|
||||
if (server.lfu_decay_time < 0) {
|
||||
err = "lfu-decay-time must be 0 or greater";
|
||||
goto loaderr;
|
||||
}
|
||||
@@ -1221,6 +1221,8 @@ void configGetCommand(client *c) {
|
||||
/* Numerical values */
|
||||
config_get_numerical_field("maxmemory",server.maxmemory);
|
||||
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);
|
||||
config_get_numerical_field("timeout",server.maxidletime);
|
||||
config_get_numerical_field("active-defrag-threshold-lower",server.active_defrag_threshold_lower);
|
||||
config_get_numerical_field("active-defrag-threshold-upper",server.active_defrag_threshold_upper);
|
||||
@@ -1992,6 +1994,8 @@ int rewriteConfig(char *path) {
|
||||
rewriteConfigBytesOption(state,"maxmemory",server.maxmemory,CONFIG_DEFAULT_MAXMEMORY);
|
||||
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);
|
||||
rewriteConfigNumericalOption(state,"lfu-decay-time",server.lfu_decay_time,CONFIG_DEFAULT_LFU_DECAY_TIME);
|
||||
rewriteConfigNumericalOption(state,"active-defrag-threshold-lower",server.active_defrag_threshold_lower,CONFIG_DEFAULT_DEFRAG_THRESHOLD_LOWER);
|
||||
rewriteConfigNumericalOption(state,"active-defrag-threshold-upper",server.active_defrag_threshold_upper,CONFIG_DEFAULT_DEFRAG_THRESHOLD_UPPER);
|
||||
rewriteConfigBytesOption(state,"active-defrag-ignore-bytes",server.active_defrag_ignore_bytes,CONFIG_DEFAULT_DEFRAG_IGNORE_BYTES);
|
||||
|
||||
@@ -38,6 +38,15 @@
|
||||
* C-level DB API
|
||||
*----------------------------------------------------------------------------*/
|
||||
|
||||
/* Update LFU when an object is accessed.
|
||||
* Firstly, decrement the counter if the decrement time is reached.
|
||||
* Then logarithmically increment the counter, and update the access time. */
|
||||
void updateLFU(robj *val) {
|
||||
unsigned long counter = LFUDecrAndReturn(val);
|
||||
counter = LFULogIncr(counter);
|
||||
val->lru = (LFUGetTimeInMinutes()<<8) | counter;
|
||||
}
|
||||
|
||||
/* Low level key lookup API, not actually called directly from commands
|
||||
* implementations that should instead rely on lookupKeyRead(),
|
||||
* lookupKeyWrite() and lookupKeyReadWithFlags(). */
|
||||
@@ -54,9 +63,7 @@ robj *lookupKey(redisDb *db, robj *key, int flags) {
|
||||
!(flags & LOOKUP_NOTOUCH))
|
||||
{
|
||||
if (server.maxmemory_policy & MAXMEMORY_FLAG_LFU) {
|
||||
unsigned long ldt = val->lru >> 8;
|
||||
unsigned long counter = LFULogIncr(val->lru & 255);
|
||||
val->lru = (ldt << 8) | counter;
|
||||
updateLFU(val);
|
||||
} else {
|
||||
val->lru = LRU_CLOCK();
|
||||
}
|
||||
@@ -180,6 +187,9 @@ void dbOverwrite(redisDb *db, robj *key, robj *val) {
|
||||
int saved_lru = old->lru;
|
||||
dictReplace(db->dict, key->ptr, val);
|
||||
val->lru = saved_lru;
|
||||
/* LFU should be not only copied but also updated
|
||||
* when a key is overwritten. */
|
||||
updateLFU(val);
|
||||
} else {
|
||||
dictReplace(db->dict, key->ptr, val);
|
||||
}
|
||||
@@ -416,7 +426,9 @@ void flushallCommand(client *c) {
|
||||
/* Normally rdbSave() will reset dirty, but we don't want this here
|
||||
* as otherwise FLUSHALL will not be replicated nor put into the AOF. */
|
||||
int saved_dirty = server.dirty;
|
||||
rdbSave(server.rdb_filename,NULL);
|
||||
rdbSaveInfo rsi, *rsiptr;
|
||||
rsiptr = rdbPopulateSaveInfo(&rsi);
|
||||
rdbSave(server.rdb_filename,rsiptr);
|
||||
server.dirty = saved_dirty;
|
||||
}
|
||||
server.dirty++;
|
||||
|
||||
+12
-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';
|
||||
@@ -335,7 +337,9 @@ void debugCommand(client *c) {
|
||||
if (c->argc >= 3) c->argv[2] = tryObjectEncoding(c->argv[2]);
|
||||
serverAssertWithInfo(c,c->argv[0],1 == 2);
|
||||
} else if (!strcasecmp(c->argv[1]->ptr,"reload")) {
|
||||
if (rdbSave(server.rdb_filename,NULL) != C_OK) {
|
||||
rdbSaveInfo rsi, *rsiptr;
|
||||
rsiptr = rdbPopulateSaveInfo(&rsi);
|
||||
if (rdbSave(server.rdb_filename,rsiptr) != C_OK) {
|
||||
addReply(c,shared.err);
|
||||
return;
|
||||
}
|
||||
@@ -368,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 */
|
||||
@@ -547,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);
|
||||
|
||||
+1
-1
@@ -940,7 +940,7 @@ static unsigned long _dictNextPower(unsigned long size)
|
||||
{
|
||||
unsigned long i = DICT_HT_INITIAL_SIZE;
|
||||
|
||||
if (size >= LONG_MAX) return LONG_MAX;
|
||||
if (size >= LONG_MAX) return LONG_MAX + 1LU;
|
||||
while(1) {
|
||||
if (i >= size)
|
||||
return i;
|
||||
|
||||
+11
-16
@@ -60,8 +60,6 @@ struct evictionPoolEntry {
|
||||
|
||||
static struct evictionPoolEntry *EvictionPoolLRU;
|
||||
|
||||
unsigned long LFUDecrAndReturn(robj *o);
|
||||
|
||||
/* ----------------------------------------------------------------------------
|
||||
* Implementation of eviction, aging and LRU
|
||||
* --------------------------------------------------------------------------*/
|
||||
@@ -302,8 +300,8 @@ unsigned long LFUGetTimeInMinutes(void) {
|
||||
return (server.unixtime/60) & 65535;
|
||||
}
|
||||
|
||||
/* Given an object last decrement time, compute the minimum number of minutes
|
||||
* that elapsed since the last decrement. Handle overflow (ldt greater than
|
||||
/* Given an object last access time, compute the minimum number of minutes
|
||||
* that elapsed since the last access. Handle overflow (ldt greater than
|
||||
* the current 16 bits minutes time) considering the time as wrapping
|
||||
* exactly once. */
|
||||
unsigned long LFUTimeElapsed(unsigned long ldt) {
|
||||
@@ -324,25 +322,22 @@ uint8_t LFULogIncr(uint8_t counter) {
|
||||
return counter;
|
||||
}
|
||||
|
||||
/* If the object decrement time is reached, decrement the LFU counter and
|
||||
* update the decrement time field. Return the object frequency counter.
|
||||
/* If the object decrement time is reached decrement the LFU counter but
|
||||
* do not update LFU fields of the object, we update the access time
|
||||
* and counter in an explicit way when the object is really accessed.
|
||||
* And we will times halve the counter according to the times of
|
||||
* elapsed time than server.lfu_decay_time.
|
||||
* Return the object frequency counter.
|
||||
*
|
||||
* This function is used in order to scan the dataset for the best object
|
||||
* to fit: as we check for the candidate, we incrementally decrement the
|
||||
* counter of the scanned objects if needed. */
|
||||
#define LFU_DECR_INTERVAL 1
|
||||
unsigned long LFUDecrAndReturn(robj *o) {
|
||||
unsigned long ldt = o->lru >> 8;
|
||||
unsigned long counter = o->lru & 255;
|
||||
if (LFUTimeElapsed(ldt) >= server.lfu_decay_time && counter) {
|
||||
if (counter > LFU_INIT_VAL*2) {
|
||||
counter /= 2;
|
||||
if (counter < LFU_INIT_VAL*2) counter = LFU_INIT_VAL*2;
|
||||
} else {
|
||||
counter--;
|
||||
}
|
||||
o->lru = (LFUGetTimeInMinutes()<<8) | counter;
|
||||
}
|
||||
unsigned long num_periods = server.lfu_decay_time ? LFUTimeElapsed(ldt) / server.lfu_decay_time : 0;
|
||||
if (num_periods)
|
||||
counter = (num_periods > counter) ? 0 : counter - num_periods;
|
||||
return counter;
|
||||
}
|
||||
|
||||
|
||||
+11
-8
@@ -103,7 +103,7 @@ void activeExpireCycle(int type) {
|
||||
|
||||
int j, iteration = 0;
|
||||
int dbs_per_call = CRON_DBS_PER_CALL;
|
||||
long long start = ustime(), timelimit;
|
||||
long long start = ustime(), timelimit, elapsed;
|
||||
|
||||
/* When clients are paused the dataset should be static not just from the
|
||||
* POV of clients not being able to write, but also from the POV of
|
||||
@@ -140,7 +140,7 @@ void activeExpireCycle(int type) {
|
||||
if (type == ACTIVE_EXPIRE_CYCLE_FAST)
|
||||
timelimit = ACTIVE_EXPIRE_CYCLE_FAST_DURATION; /* in microseconds. */
|
||||
|
||||
for (j = 0; j < dbs_per_call; j++) {
|
||||
for (j = 0; j < dbs_per_call && timelimit_exit == 0; j++) {
|
||||
int expired;
|
||||
redisDb *db = server.db+(current_db % server.dbnum);
|
||||
|
||||
@@ -155,6 +155,7 @@ void activeExpireCycle(int type) {
|
||||
unsigned long num, slots;
|
||||
long long now, ttl_sum;
|
||||
int ttl_samples;
|
||||
iteration++;
|
||||
|
||||
/* If there is nothing to expire try next DB ASAP. */
|
||||
if ((num = dictSize(db->expires)) == 0) {
|
||||
@@ -207,18 +208,20 @@ void activeExpireCycle(int type) {
|
||||
/* We can't block forever here even if there are many keys to
|
||||
* expire. So after a given amount of milliseconds return to the
|
||||
* caller waiting for the other active expire cycle. */
|
||||
iteration++;
|
||||
if ((iteration & 0xf) == 0) { /* check once every 16 iterations. */
|
||||
long long elapsed = ustime()-start;
|
||||
|
||||
latencyAddSampleIfNeeded("expire-cycle",elapsed/1000);
|
||||
if (elapsed > timelimit) timelimit_exit = 1;
|
||||
elapsed = ustime()-start;
|
||||
if (elapsed > timelimit) {
|
||||
timelimit_exit = 1;
|
||||
break;
|
||||
}
|
||||
}
|
||||
if (timelimit_exit) return;
|
||||
/* We don't repeat the cycle if there are less than 25% of keys
|
||||
* found expired in the current DB. */
|
||||
} while (expired > ACTIVE_EXPIRE_CYCLE_LOOKUPS_PER_LOOP/4);
|
||||
}
|
||||
|
||||
elapsed = ustime()-start;
|
||||
latencyAddSampleIfNeeded("expire-cycle",elapsed/1000);
|
||||
}
|
||||
|
||||
/*-----------------------------------------------------------------------------
|
||||
|
||||
+5
-1
@@ -79,7 +79,11 @@
|
||||
* Unconditionally aligning does not cost very much, so do it if unsure
|
||||
*/
|
||||
#ifndef STRICT_ALIGN
|
||||
# define STRICT_ALIGN !(defined(__i386) || defined (__amd64))
|
||||
# if !(defined(__i386) || defined (__amd64))
|
||||
# define STRICT_ALIGN 1
|
||||
# else
|
||||
# define STRICT_ALIGN 0
|
||||
# endif
|
||||
#endif
|
||||
|
||||
/*
|
||||
|
||||
+123
-30
@@ -442,9 +442,7 @@ void moduleFreeContext(RedisModuleCtx *ctx) {
|
||||
void moduleHandlePropagationAfterCommandCallback(RedisModuleCtx *ctx) {
|
||||
client *c = ctx->client;
|
||||
|
||||
/* We don't want any automatic propagation here since in modules we handle
|
||||
* replication / AOF propagation in explicit ways. */
|
||||
preventCommandPropagation(c);
|
||||
if (c->flags & CLIENT_LUA) return;
|
||||
|
||||
/* Handle the replication of the final EXEC, since whatever a command
|
||||
* emits is always wrappered around MULTI/EXEC. */
|
||||
@@ -615,7 +613,7 @@ int RM_CreateCommand(RedisModuleCtx *ctx, const char *name, RedisModuleCmdFunc c
|
||||
sds cmdname = sdsnew(name);
|
||||
|
||||
/* Check if the command name is busy. */
|
||||
if (lookupCommand((char*)name) != NULL) {
|
||||
if (lookupCommand(cmdname) != NULL) {
|
||||
sdsfree(cmdname);
|
||||
return REDISMODULE_ERR;
|
||||
}
|
||||
@@ -650,7 +648,7 @@ int RM_CreateCommand(RedisModuleCtx *ctx, const char *name, RedisModuleCmdFunc c
|
||||
*
|
||||
* This is an internal function, Redis modules developers don't need
|
||||
* to use it. */
|
||||
void RM_SetModuleAttribs(RedisModuleCtx *ctx, const char *name, int ver, int apiver){
|
||||
void RM_SetModuleAttribs(RedisModuleCtx *ctx, const char *name, int ver, int apiver) {
|
||||
RedisModule *module;
|
||||
|
||||
if (ctx->module != NULL) return;
|
||||
@@ -662,6 +660,15 @@ void RM_SetModuleAttribs(RedisModuleCtx *ctx, const char *name, int ver, int api
|
||||
ctx->module = module;
|
||||
}
|
||||
|
||||
/* Return non-zero if the module name is busy.
|
||||
* Otherwise zero is returned. */
|
||||
int RM_IsModuleNameBusy(const char *name) {
|
||||
sds modulename = sdsnew(name);
|
||||
dictEntry *de = dictFind(modules,modulename);
|
||||
sdsfree(modulename);
|
||||
return de != NULL;
|
||||
}
|
||||
|
||||
/* Return the current UNIX time in milliseconds. */
|
||||
long long RM_Milliseconds(void) {
|
||||
return mstime();
|
||||
@@ -1164,6 +1171,9 @@ int RM_ReplyWithDouble(RedisModuleCtx *ctx, double d) {
|
||||
* in the context of a command execution. EXEC will be handled by the
|
||||
* RedisModuleCommandDispatcher() function. */
|
||||
void moduleReplicateMultiIfNeeded(RedisModuleCtx *ctx) {
|
||||
/* Skip this if client explicitly wrap the command with MULTI, or if
|
||||
* the module command was called by a script. */
|
||||
if (ctx->client->flags & (CLIENT_MULTI|CLIENT_LUA)) return;
|
||||
/* If we already emitted MULTI return ASAP. */
|
||||
if (ctx->flags & REDISMODULE_CTX_MULTI_EMITTED) return;
|
||||
/* If this is a thread safe context, we do not want to wrap commands
|
||||
@@ -1216,6 +1226,7 @@ int RM_Replicate(RedisModuleCtx *ctx, const char *cmdname, const char *fmt, ...)
|
||||
/* Release the argv. */
|
||||
for (j = 0; j < argc; j++) decrRefCount(argv[j]);
|
||||
zfree(argv);
|
||||
server.dirty++;
|
||||
return REDISMODULE_OK;
|
||||
}
|
||||
|
||||
@@ -1234,6 +1245,7 @@ int RM_ReplicateVerbatim(RedisModuleCtx *ctx) {
|
||||
alsoPropagate(ctx->client->cmd,ctx->client->db->id,
|
||||
ctx->client->argv,ctx->client->argc,
|
||||
PROPAGATE_AOF|PROPAGATE_REPL);
|
||||
server.dirty++;
|
||||
return REDISMODULE_OK;
|
||||
}
|
||||
|
||||
@@ -1262,6 +1274,74 @@ int RM_GetSelectedDb(RedisModuleCtx *ctx) {
|
||||
return ctx->client->db->id;
|
||||
}
|
||||
|
||||
|
||||
/* Return the current context's flags. The flags provide information on the
|
||||
* current request context (whether the client is a Lua script or in a MULTI),
|
||||
* and about the Redis instance in general, i.e replication and persistence.
|
||||
*
|
||||
* The available flags are:
|
||||
*
|
||||
* * REDISMODULE_CTX_FLAGS_LUA: The command is running in a Lua script
|
||||
*
|
||||
* * REDISMODULE_CTX_FLAGS_MULTI: The command is running inside a transaction
|
||||
*
|
||||
* * REDISMODULE_CTX_FLAGS_MASTER: The Redis instance is a master
|
||||
*
|
||||
* * REDISMODULE_CTX_FLAGS_SLAVE: The Redis instance is a slave
|
||||
*
|
||||
* * REDISMODULE_CTX_FLAGS_READONLY: The Redis instance is read-only
|
||||
*
|
||||
* * REDISMODULE_CTX_FLAGS_CLUSTER: The Redis instance is in cluster mode
|
||||
*
|
||||
* * REDISMODULE_CTX_FLAGS_AOF: The Redis instance has AOF enabled
|
||||
*
|
||||
* * REDISMODULE_CTX_FLAGS_RDB: The instance has RDB enabled
|
||||
*
|
||||
* * REDISMODULE_CTX_FLAGS_MAXMEMORY: The instance has Maxmemory set
|
||||
*
|
||||
* * REDISMODULE_CTX_FLAGS_EVICT: Maxmemory is set and has an eviction
|
||||
* policy that may delete keys
|
||||
*/
|
||||
int RM_GetContextFlags(RedisModuleCtx *ctx) {
|
||||
|
||||
int flags = 0;
|
||||
/* Client specific flags */
|
||||
if (ctx->client) {
|
||||
if (ctx->client->flags & CLIENT_LUA)
|
||||
flags |= REDISMODULE_CTX_FLAGS_LUA;
|
||||
if (ctx->client->flags & CLIENT_MULTI)
|
||||
flags |= REDISMODULE_CTX_FLAGS_MULTI;
|
||||
}
|
||||
|
||||
if (server.cluster_enabled)
|
||||
flags |= REDISMODULE_CTX_FLAGS_CLUSTER;
|
||||
|
||||
/* Maxmemory and eviction policy */
|
||||
if (server.maxmemory > 0) {
|
||||
flags |= REDISMODULE_CTX_FLAGS_MAXMEMORY;
|
||||
|
||||
if (server.maxmemory_policy != MAXMEMORY_NO_EVICTION)
|
||||
flags |= REDISMODULE_CTX_FLAGS_EVICT;
|
||||
}
|
||||
|
||||
/* Persistence flags */
|
||||
if (server.aof_state != AOF_OFF)
|
||||
flags |= REDISMODULE_CTX_FLAGS_AOF;
|
||||
if (server.saveparamslen > 0)
|
||||
flags |= REDISMODULE_CTX_FLAGS_RDB;
|
||||
|
||||
/* Replication flags */
|
||||
if (server.masterhost == NULL) {
|
||||
flags |= REDISMODULE_CTX_FLAGS_MASTER;
|
||||
} else {
|
||||
flags |= REDISMODULE_CTX_FLAGS_SLAVE;
|
||||
if (server.repl_slave_ro)
|
||||
flags |= REDISMODULE_CTX_FLAGS_READONLY;
|
||||
}
|
||||
|
||||
return flags;
|
||||
}
|
||||
|
||||
/* Change the currently selected DB. Returns an error if the id
|
||||
* is out of range.
|
||||
*
|
||||
@@ -3333,14 +3413,16 @@ void unblockClientFromModule(client *c) {
|
||||
RedisModuleBlockedClient *RM_BlockClient(RedisModuleCtx *ctx, RedisModuleCmdFunc reply_callback, RedisModuleCmdFunc timeout_callback, void (*free_privdata)(void*), long long timeout_ms) {
|
||||
client *c = ctx->client;
|
||||
int islua = c->flags & CLIENT_LUA;
|
||||
int ismulti = c->flags & CLIENT_MULTI;
|
||||
|
||||
c->bpop.module_blocked_handle = zmalloc(sizeof(RedisModuleBlockedClient));
|
||||
RedisModuleBlockedClient *bc = c->bpop.module_blocked_handle;
|
||||
|
||||
/* We need to handle the invalid operation of calling modules blocking
|
||||
* commands from Lua. We actually create an already aborted (client set to
|
||||
* NULL) blocked client handle, and actually reply to Lua with an error. */
|
||||
bc->client = islua ? NULL : c;
|
||||
* commands from Lua or MULTI. We actually create an already aborted
|
||||
* (client set to NULL) blocked client handle, and actually reply with
|
||||
* an error. */
|
||||
bc->client = (islua || ismulti) ? NULL : c;
|
||||
bc->module = ctx->module;
|
||||
bc->reply_callback = reply_callback;
|
||||
bc->timeout_callback = timeout_callback;
|
||||
@@ -3351,9 +3433,11 @@ RedisModuleBlockedClient *RM_BlockClient(RedisModuleCtx *ctx, RedisModuleCmdFunc
|
||||
bc->dbid = c->db->id;
|
||||
c->bpop.timeout = timeout_ms ? (mstime()+timeout_ms) : 0;
|
||||
|
||||
if (islua) {
|
||||
if (islua || ismulti) {
|
||||
c->bpop.module_blocked_handle = NULL;
|
||||
addReplyError(c,"Blocking module command called from Lua script");
|
||||
addReplyError(c, islua ?
|
||||
"Blocking module command called from Lua script" :
|
||||
"Blocking module command called from transaction");
|
||||
} else {
|
||||
blockClient(c,BLOCKED_MODULE);
|
||||
}
|
||||
@@ -3661,6 +3745,28 @@ void moduleFreeModuleStructure(struct RedisModule *module) {
|
||||
zfree(module);
|
||||
}
|
||||
|
||||
void moduleUnregisterCommands(struct RedisModule *module) {
|
||||
/* Unregister all the commands registered by this module. */
|
||||
dictIterator *di = dictGetSafeIterator(server.commands);
|
||||
dictEntry *de;
|
||||
while ((de = dictNext(di)) != NULL) {
|
||||
struct redisCommand *cmd = dictGetVal(de);
|
||||
if (cmd->proc == RedisModuleCommandDispatcher) {
|
||||
RedisModuleCommandProxy *cp =
|
||||
(void*)(unsigned long)cmd->getkeys_proc;
|
||||
sds cmdname = cp->rediscmd->name;
|
||||
if (cp->module == module) {
|
||||
dictDelete(server.commands,cmdname);
|
||||
dictDelete(server.orig_commands,cmdname);
|
||||
sdsfree(cmdname);
|
||||
zfree(cp->rediscmd);
|
||||
zfree(cp);
|
||||
}
|
||||
}
|
||||
}
|
||||
dictReleaseIterator(di);
|
||||
}
|
||||
|
||||
/* Load a module and initialize it. On success C_OK is returned, otherwise
|
||||
* C_ERR is returned. */
|
||||
int moduleLoad(const char *path, void **module_argv, int module_argc) {
|
||||
@@ -3681,7 +3787,10 @@ int moduleLoad(const char *path, void **module_argv, int module_argc) {
|
||||
return C_ERR;
|
||||
}
|
||||
if (onload((void*)&ctx,module_argv,module_argc) == REDISMODULE_ERR) {
|
||||
if (ctx.module) moduleFreeModuleStructure(ctx.module);
|
||||
if (ctx.module) {
|
||||
moduleUnregisterCommands(ctx.module);
|
||||
moduleFreeModuleStructure(ctx.module);
|
||||
}
|
||||
dlclose(handle);
|
||||
serverLog(LL_WARNING,
|
||||
"Module %s initialization failed. Module not loaded",path);
|
||||
@@ -3715,25 +3824,7 @@ int moduleUnload(sds name) {
|
||||
return REDISMODULE_ERR;
|
||||
}
|
||||
|
||||
/* Unregister all the commands registered by this module. */
|
||||
dictIterator *di = dictGetSafeIterator(server.commands);
|
||||
dictEntry *de;
|
||||
while ((de = dictNext(di)) != NULL) {
|
||||
struct redisCommand *cmd = dictGetVal(de);
|
||||
if (cmd->proc == RedisModuleCommandDispatcher) {
|
||||
RedisModuleCommandProxy *cp =
|
||||
(void*)(unsigned long)cmd->getkeys_proc;
|
||||
sds cmdname = cp->rediscmd->name;
|
||||
if (cp->module == module) {
|
||||
dictDelete(server.commands,cmdname);
|
||||
dictDelete(server.orig_commands,cmdname);
|
||||
sdsfree(cmdname);
|
||||
zfree(cp->rediscmd);
|
||||
zfree(cp);
|
||||
}
|
||||
}
|
||||
}
|
||||
dictReleaseIterator(di);
|
||||
moduleUnregisterCommands(module);
|
||||
|
||||
/* Unregister all the hooks. TODO: Yet no hooks support here. */
|
||||
|
||||
@@ -3828,6 +3919,7 @@ void moduleRegisterCoreAPI(void) {
|
||||
REGISTER_API(Strdup);
|
||||
REGISTER_API(CreateCommand);
|
||||
REGISTER_API(SetModuleAttribs);
|
||||
REGISTER_API(IsModuleNameBusy);
|
||||
REGISTER_API(WrongArity);
|
||||
REGISTER_API(ReplyWithLongLong);
|
||||
REGISTER_API(ReplyWithError);
|
||||
@@ -3891,6 +3983,7 @@ void moduleRegisterCoreAPI(void) {
|
||||
REGISTER_API(IsKeysPositionRequest);
|
||||
REGISTER_API(KeyAtPos);
|
||||
REGISTER_API(GetClientId);
|
||||
REGISTER_API(GetContextFlags);
|
||||
REGISTER_API(PoolAlloc);
|
||||
REGISTER_API(CreateDataType);
|
||||
REGISTER_API(ModuleTypeSetValue);
|
||||
|
||||
@@ -121,6 +121,81 @@ int TestStringPrintf(RedisModuleCtx *ctx, RedisModuleString **argv, int argc) {
|
||||
}
|
||||
|
||||
|
||||
/* TEST.CTXFLAGS -- Test GetContextFlags. */
|
||||
int TestCtxFlags(RedisModuleCtx *ctx, RedisModuleString **argv, int argc) {
|
||||
REDISMODULE_NOT_USED(argc);
|
||||
REDISMODULE_NOT_USED(argv);
|
||||
|
||||
RedisModule_AutoMemory(ctx);
|
||||
|
||||
int ok = 1;
|
||||
const char *errString = NULL;
|
||||
|
||||
#define FAIL(msg) \
|
||||
{ \
|
||||
ok = 0; \
|
||||
errString = msg; \
|
||||
goto end; \
|
||||
}
|
||||
|
||||
int flags = RedisModule_GetContextFlags(ctx);
|
||||
if (flags == 0) {
|
||||
FAIL("Got no flags");
|
||||
}
|
||||
|
||||
if (flags & REDISMODULE_CTX_FLAGS_LUA) FAIL("Lua flag was set");
|
||||
if (flags & REDISMODULE_CTX_FLAGS_MULTI) FAIL("Multi flag was set");
|
||||
|
||||
if (flags & REDISMODULE_CTX_FLAGS_AOF) FAIL("AOF Flag was set")
|
||||
/* Enable AOF to test AOF flags */
|
||||
RedisModule_Call(ctx, "config", "ccc", "set", "appendonly", "yes");
|
||||
flags = RedisModule_GetContextFlags(ctx);
|
||||
if (!(flags & REDISMODULE_CTX_FLAGS_AOF))
|
||||
FAIL("AOF Flag not set after config set");
|
||||
|
||||
if (flags & REDISMODULE_CTX_FLAGS_RDB) FAIL("RDB Flag was set");
|
||||
/* Enable RDB to test RDB flags */
|
||||
RedisModule_Call(ctx, "config", "ccc", "set", "save", "900 1");
|
||||
flags = RedisModule_GetContextFlags(ctx);
|
||||
if (!(flags & REDISMODULE_CTX_FLAGS_RDB))
|
||||
FAIL("RDB Flag was not set after config set");
|
||||
|
||||
if (!(flags & REDISMODULE_CTX_FLAGS_MASTER)) FAIL("Master flag was not set");
|
||||
if (flags & REDISMODULE_CTX_FLAGS_SLAVE) FAIL("Slave flag was set");
|
||||
if (flags & REDISMODULE_CTX_FLAGS_READONLY) FAIL("Read-only flag was set");
|
||||
if (flags & REDISMODULE_CTX_FLAGS_CLUSTER) FAIL("Cluster flag was set");
|
||||
|
||||
if (flags & REDISMODULE_CTX_FLAGS_MAXMEMORY) FAIL("Maxmemory flag was set");
|
||||
;
|
||||
RedisModule_Call(ctx, "config", "ccc", "set", "maxmemory", "100000000");
|
||||
flags = RedisModule_GetContextFlags(ctx);
|
||||
if (!(flags & REDISMODULE_CTX_FLAGS_MAXMEMORY))
|
||||
FAIL("Maxmemory flag was not set after config set");
|
||||
|
||||
if (flags & REDISMODULE_CTX_FLAGS_EVICT) FAIL("Eviction flag was set");
|
||||
RedisModule_Call(ctx, "config", "ccc", "set", "maxmemory-policy",
|
||||
"allkeys-lru");
|
||||
flags = RedisModule_GetContextFlags(ctx);
|
||||
if (!(flags & REDISMODULE_CTX_FLAGS_EVICT))
|
||||
FAIL("Eviction flag was not set after config set");
|
||||
|
||||
end:
|
||||
/* Revert config changes */
|
||||
RedisModule_Call(ctx, "config", "ccc", "set", "appendonly", "no");
|
||||
RedisModule_Call(ctx, "config", "ccc", "set", "save", "");
|
||||
RedisModule_Call(ctx, "config", "ccc", "set", "maxmemory", "0");
|
||||
RedisModule_Call(ctx, "config", "ccc", "set", "maxmemory-policy", "noeviction");
|
||||
|
||||
if (!ok) {
|
||||
RedisModule_Log(ctx, "warning", "Failed CTXFLAGS Test. Reason: %s",
|
||||
errString);
|
||||
return RedisModule_ReplyWithSimpleString(ctx, "ERR");
|
||||
}
|
||||
|
||||
return RedisModule_ReplyWithSimpleString(ctx, "OK");
|
||||
}
|
||||
|
||||
|
||||
/* ----------------------------- Test framework ----------------------------- */
|
||||
|
||||
/* Return 1 if the reply matches the specified string, otherwise log errors
|
||||
@@ -188,6 +263,9 @@ int TestIt(RedisModuleCtx *ctx, RedisModuleString **argv, int argc) {
|
||||
T("test.call","");
|
||||
if (!TestAssertStringReply(ctx,reply,"OK",2)) goto fail;
|
||||
|
||||
T("test.ctxflags","");
|
||||
if (!TestAssertStringReply(ctx,reply,"OK",2)) goto fail;
|
||||
|
||||
T("test.string.append","");
|
||||
if (!TestAssertStringReply(ctx,reply,"foobar",6)) goto fail;
|
||||
|
||||
@@ -229,6 +307,10 @@ int RedisModule_OnLoad(RedisModuleCtx *ctx, RedisModuleString **argv, int argc)
|
||||
TestStringPrintf,"write deny-oom",1,1,1) == REDISMODULE_ERR)
|
||||
return REDISMODULE_ERR;
|
||||
|
||||
if (RedisModule_CreateCommand(ctx,"test.ctxflags",
|
||||
TestCtxFlags,"readonly",1,1,1) == REDISMODULE_ERR)
|
||||
return REDISMODULE_ERR;
|
||||
|
||||
if (RedisModule_CreateCommand(ctx,"test.it",
|
||||
TestIt,"readonly",1,1,1) == REDISMODULE_ERR)
|
||||
return REDISMODULE_ERR;
|
||||
|
||||
+38
-13
@@ -558,11 +558,11 @@ int getDoubleFromObject(const robj *o, double *target) {
|
||||
if (sdsEncodedObject(o)) {
|
||||
errno = 0;
|
||||
value = strtod(o->ptr, &eptr);
|
||||
if (isspace(((const char*)o->ptr)[0]) ||
|
||||
eptr[0] != '\0' ||
|
||||
if (sdslen(o->ptr) == 0 ||
|
||||
isspace(((const char*)o->ptr)[0]) ||
|
||||
(size_t)(eptr-(char*)o->ptr) != sdslen(o->ptr) ||
|
||||
(errno == ERANGE &&
|
||||
(value == HUGE_VAL || value == -HUGE_VAL || value == 0)) ||
|
||||
errno == EINVAL ||
|
||||
isnan(value))
|
||||
return C_ERR;
|
||||
} else if (o->encoding == OBJ_ENCODING_INT) {
|
||||
@@ -600,8 +600,12 @@ int getLongDoubleFromObject(robj *o, long double *target) {
|
||||
if (sdsEncodedObject(o)) {
|
||||
errno = 0;
|
||||
value = strtold(o->ptr, &eptr);
|
||||
if (isspace(((char*)o->ptr)[0]) || eptr[0] != '\0' ||
|
||||
errno == ERANGE || isnan(value))
|
||||
if (sdslen(o->ptr) == 0 ||
|
||||
isspace(((const char*)o->ptr)[0]) ||
|
||||
(size_t)(eptr-(char*)o->ptr) != sdslen(o->ptr) ||
|
||||
(errno == ERANGE &&
|
||||
(value == HUGE_VAL || value == -HUGE_VAL || value == 0)) ||
|
||||
isnan(value))
|
||||
return C_ERR;
|
||||
} else if (o->encoding == OBJ_ENCODING_INT) {
|
||||
value = (long)o->ptr;
|
||||
@@ -1008,11 +1012,25 @@ robj *objectCommandLookupOrReply(client *c, robj *key, robj *reply) {
|
||||
}
|
||||
|
||||
/* Object command allows to inspect the internals of an Redis Object.
|
||||
* Usage: OBJECT <refcount|encoding|idletime> <key> */
|
||||
* Usage: OBJECT <refcount|encoding|idletime|freq> <key> */
|
||||
void objectCommand(client *c) {
|
||||
robj *o;
|
||||
|
||||
if (!strcasecmp(c->argv[1]->ptr,"refcount") && c->argc == 3) {
|
||||
if (!strcasecmp(c->argv[1]->ptr,"help") && c->argc == 2) {
|
||||
void *blenp = addDeferredMultiBulkLength(c);
|
||||
int blen = 0;
|
||||
blen++; addReplyStatus(c,
|
||||
"OBJECT <subcommand> key. Subcommands:");
|
||||
blen++; addReplyStatus(c,
|
||||
"refcount -- Return the number of references of the value associated with the specified key.");
|
||||
blen++; addReplyStatus(c,
|
||||
"encoding -- Return the kind of internal representation used in order to store the value associated with a key.");
|
||||
blen++; addReplyStatus(c,
|
||||
"idletime -- Return the idle time of the key, that is the approximated number of seconds elapsed since the last access to the key.");
|
||||
blen++; addReplyStatus(c,
|
||||
"freq -- Return the access frequency index of the key. The returned integer is proportional to the logarithm of the recent access frequency of the key.");
|
||||
setDeferredMultiBulkLength(c,blenp,blen);
|
||||
} else if (!strcasecmp(c->argv[1]->ptr,"refcount") && c->argc == 3) {
|
||||
if ((o = objectCommandLookupOrReply(c,c->argv[2],shared.nullbulk))
|
||||
== NULL) return;
|
||||
addReplyLongLong(c,o->refcount);
|
||||
@@ -1031,13 +1049,18 @@ void objectCommand(client *c) {
|
||||
} else if (!strcasecmp(c->argv[1]->ptr,"freq") && c->argc == 3) {
|
||||
if ((o = objectCommandLookupOrReply(c,c->argv[2],shared.nullbulk))
|
||||
== NULL) return;
|
||||
if (server.maxmemory_policy & MAXMEMORY_FLAG_LRU) {
|
||||
addReplyError(c,"An LRU maxmemory policy is selected, access frequency not tracked. Please note that when switching between policies at runtime LRU and LFU data will take some time to adjust.");
|
||||
if (!(server.maxmemory_policy & MAXMEMORY_FLAG_LFU)) {
|
||||
addReplyError(c,"An LFU maxmemory policy is not selected, access frequency not tracked. Please note that when switching between policies at runtime LRU and LFU data will take some time to adjust.");
|
||||
return;
|
||||
}
|
||||
addReplyLongLong(c,o->lru&255);
|
||||
/* LFUDecrAndReturn should be called
|
||||
* in case of the key has not been accessed for a long time,
|
||||
* because we update the access time only
|
||||
* when the key is read or overwritten. */
|
||||
addReplyLongLong(c,LFUDecrAndReturn(o));
|
||||
} else {
|
||||
addReplyError(c,"Syntax error. Try OBJECT (refcount|encoding|idletime|freq)");
|
||||
addReplyErrorFormat(c, "Unknown subcommand or wrong number of arguments for '%s'. Try OBJECT help",
|
||||
(char *)c->argv[1]->ptr);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1070,7 +1093,7 @@ void memoryCommand(client *c) {
|
||||
if ((o = objectCommandLookupOrReply(c,c->argv[2],shared.nullbulk))
|
||||
== NULL) return;
|
||||
size_t usage = objectComputeSize(o,samples);
|
||||
usage += sdsAllocSize(c->argv[1]->ptr);
|
||||
usage += sdsAllocSize(c->argv[2]->ptr);
|
||||
usage += sizeof(dictEntry);
|
||||
addReplyLongLong(c,usage);
|
||||
} else if (!strcasecmp(c->argv[1]->ptr,"stats") && c->argc == 2) {
|
||||
@@ -1163,7 +1186,9 @@ void memoryCommand(client *c) {
|
||||
/* Nothing to do for other allocators. */
|
||||
#endif
|
||||
} else if (!strcasecmp(c->argv[1]->ptr,"help") && c->argc == 2) {
|
||||
addReplyMultiBulkLen(c,4);
|
||||
addReplyMultiBulkLen(c,5);
|
||||
addReplyBulkCString(c,
|
||||
"MEMORY DOCTOR - Outputs memory problems report");
|
||||
addReplyBulkCString(c,
|
||||
"MEMORY USAGE <key> [SAMPLES <count>] - Estimate memory usage of key");
|
||||
addReplyBulkCString(c,
|
||||
|
||||
+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);
|
||||
|
||||
|
||||
@@ -656,7 +656,7 @@ ssize_t rdbSaveObject(rio *rdb, robj *o) {
|
||||
if ((n = rdbSaveLen(rdb,ql->len)) == -1) return -1;
|
||||
nwritten += n;
|
||||
|
||||
do {
|
||||
while(node) {
|
||||
if (quicklistNodeIsCompressed(node)) {
|
||||
void *data;
|
||||
size_t compress_len = quicklistGetLzf(node, &data);
|
||||
@@ -666,7 +666,8 @@ ssize_t rdbSaveObject(rio *rdb, robj *o) {
|
||||
if ((n = rdbSaveRawString(rdb,node->zl,node->sz)) == -1) return -1;
|
||||
nwritten += n;
|
||||
}
|
||||
} while ((node = node->next));
|
||||
node = node->next;
|
||||
}
|
||||
} else {
|
||||
serverPanic("Unknown list encoding");
|
||||
}
|
||||
@@ -858,16 +859,14 @@ int rdbSaveInfoAuxFields(rio *rdb, int flags, rdbSaveInfo *rsi) {
|
||||
|
||||
/* Handle saving options that generate aux fields. */
|
||||
if (rsi) {
|
||||
if (rsi->repl_stream_db &&
|
||||
rdbSaveAuxFieldStrInt(rdb,"repl-stream-db",rsi->repl_stream_db)
|
||||
== -1)
|
||||
{
|
||||
return -1;
|
||||
}
|
||||
if (rdbSaveAuxFieldStrInt(rdb,"repl-stream-db",rsi->repl_stream_db)
|
||||
== -1) return -1;
|
||||
if (rdbSaveAuxFieldStrStr(rdb,"repl-id",server.replid)
|
||||
== -1) return -1;
|
||||
if (rdbSaveAuxFieldStrInt(rdb,"repl-offset",server.master_repl_offset)
|
||||
== -1) return -1;
|
||||
}
|
||||
if (rdbSaveAuxFieldStrInt(rdb,"aof-preamble",aof_preamble) == -1) return -1;
|
||||
if (rdbSaveAuxFieldStrStr(rdb,"repl-id",server.replid) == -1) return -1;
|
||||
if (rdbSaveAuxFieldStrInt(rdb,"repl-offset",server.master_repl_offset) == -1) return -1;
|
||||
return 1;
|
||||
}
|
||||
|
||||
@@ -944,6 +943,20 @@ int rdbSaveRio(rio *rdb, int *error, int flags, rdbSaveInfo *rsi) {
|
||||
}
|
||||
di = NULL; /* So that we don't release it again on error. */
|
||||
|
||||
/* If we are storing the replication information on disk, persist
|
||||
* the script cache as well: on successful PSYNC after a restart, we need
|
||||
* to be able to process any EVALSHA inside the replication backlog the
|
||||
* master will send us. */
|
||||
if (rsi && dictSize(server.lua_scripts)) {
|
||||
di = dictGetIterator(server.lua_scripts);
|
||||
while((de = dictNext(di)) != NULL) {
|
||||
robj *body = dictGetVal(de);
|
||||
if (rdbSaveAuxField(rdb,"lua",3,body->ptr,sdslen(body->ptr)) == -1)
|
||||
goto werr;
|
||||
}
|
||||
dictReleaseIterator(di);
|
||||
}
|
||||
|
||||
/* EOF opcode */
|
||||
if (rdbSaveType(rdb,RDB_OPCODE_EOF) == -1) goto werr;
|
||||
|
||||
@@ -1590,6 +1603,13 @@ int rdbLoadRio(rio *rdb, rdbSaveInfo *rsi) {
|
||||
}
|
||||
} else if (!strcasecmp(auxkey->ptr,"repl-offset")) {
|
||||
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,auxval) == NULL) {
|
||||
rdbExitReportCorruptRDB(
|
||||
"Can't load Lua script from RDB file! "
|
||||
"BODY: %s", auxval->ptr);
|
||||
}
|
||||
} else {
|
||||
/* We ignore fields we don't understand, as by AUX field
|
||||
* contract. */
|
||||
@@ -1977,7 +1997,9 @@ void saveCommand(client *c) {
|
||||
addReplyError(c,"Background save already in progress");
|
||||
return;
|
||||
}
|
||||
if (rdbSave(server.rdb_filename,NULL) == C_OK) {
|
||||
rdbSaveInfo rsi, *rsiptr;
|
||||
rsiptr = rdbPopulateSaveInfo(&rsi);
|
||||
if (rdbSave(server.rdb_filename,rsiptr) == C_OK) {
|
||||
addReply(c,shared.ok);
|
||||
} else {
|
||||
addReply(c,shared.err);
|
||||
@@ -1999,6 +2021,9 @@ void bgsaveCommand(client *c) {
|
||||
}
|
||||
}
|
||||
|
||||
rdbSaveInfo rsi, *rsiptr;
|
||||
rsiptr = rdbPopulateSaveInfo(&rsi);
|
||||
|
||||
if (server.rdb_child_pid != -1) {
|
||||
addReplyError(c,"Background save already in progress");
|
||||
} else if (server.aof_child_pid != -1) {
|
||||
@@ -2011,9 +2036,58 @@ void bgsaveCommand(client *c) {
|
||||
"Use BGSAVE SCHEDULE in order to schedule a BGSAVE whenever "
|
||||
"possible.");
|
||||
}
|
||||
} else if (rdbSaveBackground(server.rdb_filename,NULL) == C_OK) {
|
||||
} else if (rdbSaveBackground(server.rdb_filename,rsiptr) == C_OK) {
|
||||
addReplyStatus(c,"Background saving started");
|
||||
} else {
|
||||
addReply(c,shared.err);
|
||||
}
|
||||
}
|
||||
|
||||
/* Populate the rdbSaveInfo structure used to persist the replication
|
||||
* information inside the RDB file. Currently the structure explicitly
|
||||
* contains just the currently selected DB from the master stream, however
|
||||
* if the rdbSave*() family functions receive a NULL rsi structure also
|
||||
* the Replication ID/offset is not saved. The function popultes 'rsi'
|
||||
* that is normally stack-allocated in the caller, returns the populated
|
||||
* pointer if the instance has a valid master client, otherwise NULL
|
||||
* is returned, and the RDB saving will not persist any replication related
|
||||
* information. */
|
||||
rdbSaveInfo *rdbPopulateSaveInfo(rdbSaveInfo *rsi) {
|
||||
rdbSaveInfo rsi_init = RDB_SAVE_INFO_INIT;
|
||||
*rsi = rsi_init;
|
||||
|
||||
/* If the instance is a master, we can populate the replication info
|
||||
* only when repl_backlog is not NULL. If the repl_backlog is NULL,
|
||||
* it means that the instance isn't in any replication chains. In this
|
||||
* scenario the replication info is useless, because when a slave
|
||||
* connects to us, the NULL repl_backlog will trigger a full
|
||||
* synchronization, at the same time we will use a new replid and clear
|
||||
* replid2. */
|
||||
if (!server.masterhost && server.repl_backlog) {
|
||||
/* Note that when server.slaveseldb is -1, it means that this master
|
||||
* didn't apply any write commands after a full synchronization.
|
||||
* So we can let repl_stream_db be 0, this allows a restarted slave
|
||||
* to reload replication ID/offset, it's safe because the next write
|
||||
* command must generate a SELECT statement. */
|
||||
rsi->repl_stream_db = server.slaveseldb == -1 ? 0 : server.slaveseldb;
|
||||
return rsi;
|
||||
}
|
||||
|
||||
/* If the instance is a slave we need a connected master
|
||||
* in order to fetch the currently selected DB. */
|
||||
if (server.master) {
|
||||
rsi->repl_stream_db = server.master->db->id;
|
||||
return rsi;
|
||||
}
|
||||
|
||||
/* If we have a cached master we can use it in order to populate the
|
||||
* replication selected DB info inside the RDB file: the slave can
|
||||
* increment the master_repl_offset only from data arriving from the
|
||||
* master, so if we are disconnected the offset in the cached master
|
||||
* is valid. */
|
||||
if (server.cached_master) {
|
||||
rsi->repl_stream_db = server.cached_master->db->id;
|
||||
return rsi;
|
||||
}
|
||||
return NULL;
|
||||
}
|
||||
|
||||
@@ -147,5 +147,6 @@ int rdbLoadBinaryDoubleValue(rio *rdb, double *val);
|
||||
int rdbSaveBinaryFloatValue(rio *rdb, float val);
|
||||
int rdbLoadBinaryFloatValue(rio *rdb, float *val);
|
||||
int rdbLoadRio(rio *rdb, rdbSaveInfo *rsi);
|
||||
rdbSaveInfo *rdbPopulateSaveInfo(rdbSaveInfo *rsi);
|
||||
|
||||
#endif
|
||||
|
||||
@@ -572,8 +572,8 @@ usage:
|
||||
" -a <password> Password for Redis Auth\n"
|
||||
" -c <clients> Number of parallel connections (default 50)\n"
|
||||
" -n <requests> Total number of requests (default 100000)\n"
|
||||
" -d <size> Data size of SET/GET value in bytes (default 2)\n"
|
||||
" --dbnum <db> SELECT the specified db number (default 0)\n"
|
||||
" -d <size> Data size of SET/GET value in bytes (default 3)\n"
|
||||
" --dbnum <db> SELECT the specified db number (default 0)\n"
|
||||
" -k <boolean> 1=keep alive 0=reconnect (default 1)\n"
|
||||
" -r <keyspacelen> Use random keys for SET/GET/INCR, random values for SADD\n"
|
||||
" Using this option the benchmark will expand the string __rand_int__\n"
|
||||
|
||||
@@ -193,12 +193,12 @@ int redis_check_rdb(char *rdbfilename, FILE *fp) {
|
||||
buf[9] = '\0';
|
||||
if (memcmp(buf,"REDIS",5) != 0) {
|
||||
rdbCheckError("Wrong signature trying to load DB from file");
|
||||
return 1;
|
||||
goto err;
|
||||
}
|
||||
rdbver = atoi(buf+5);
|
||||
if (rdbver < 1 || rdbver > RDB_VERSION) {
|
||||
rdbCheckError("Can't handle RDB format version %d",rdbver);
|
||||
return 1;
|
||||
goto err;
|
||||
}
|
||||
|
||||
startLoading(fp);
|
||||
@@ -270,7 +270,7 @@ int redis_check_rdb(char *rdbfilename, FILE *fp) {
|
||||
} else {
|
||||
if (!rdbIsObjectType(type)) {
|
||||
rdbCheckError("Invalid object type: %d", type);
|
||||
return 1;
|
||||
goto err;
|
||||
}
|
||||
rdbstate.key_type = type;
|
||||
}
|
||||
@@ -307,6 +307,7 @@ int redis_check_rdb(char *rdbfilename, FILE *fp) {
|
||||
rdbCheckInfo("RDB file was saved with checksum disabled: no check performed.");
|
||||
} else if (cksum != expected) {
|
||||
rdbCheckError("RDB CRC error");
|
||||
goto err;
|
||||
} else {
|
||||
rdbCheckInfo("Checksum OK");
|
||||
}
|
||||
@@ -321,6 +322,8 @@ eoferr: /* unexpected end of file is handled here with a fatal exit */
|
||||
} else {
|
||||
rdbCheckError("Unexpected EOF reading RDB file");
|
||||
}
|
||||
err:
|
||||
if (closefile) fclose(fp);
|
||||
return 1;
|
||||
}
|
||||
|
||||
|
||||
+226
-3
@@ -107,6 +107,7 @@ static struct config {
|
||||
char *pattern;
|
||||
char *rdb_filename;
|
||||
int bigkeys;
|
||||
int hotkeys;
|
||||
int stdinarg; /* get last arg from stdin. (-x option) */
|
||||
char *auth;
|
||||
int output; /* output mode, see OUTPUT_* defines */
|
||||
@@ -198,6 +199,92 @@ static sds getDotfilePath(char *envoverride, char *dotfilename) {
|
||||
return dotPath;
|
||||
}
|
||||
|
||||
/* URL-style percent decoding. */
|
||||
#define isHexChar(c) (isdigit(c) || (c >= 'a' && c <= 'f'))
|
||||
#define decodeHexChar(c) (isdigit(c) ? c - '0' : c - 'a' + 10)
|
||||
#define decodeHex(h, l) ((decodeHexChar(h) << 4) + decodeHexChar(l))
|
||||
|
||||
static sds percentDecode(const char *pe, size_t len) {
|
||||
const char *end = pe + len;
|
||||
sds ret = sdsempty();
|
||||
const char *curr = pe;
|
||||
|
||||
while (curr < end) {
|
||||
if (*curr == '%') {
|
||||
if ((end - curr) < 2) {
|
||||
fprintf(stderr, "Incomplete URI encoding\n");
|
||||
exit(1);
|
||||
}
|
||||
|
||||
char h = tolower(*(++curr));
|
||||
char l = tolower(*(++curr));
|
||||
if (!isHexChar(h) || !isHexChar(l)) {
|
||||
fprintf(stderr, "Illegal character in URI encoding\n");
|
||||
exit(1);
|
||||
}
|
||||
char c = decodeHex(h, l);
|
||||
ret = sdscatlen(ret, &c, 1);
|
||||
curr++;
|
||||
} else {
|
||||
ret = sdscatlen(ret, curr++, 1);
|
||||
}
|
||||
}
|
||||
|
||||
return ret;
|
||||
}
|
||||
|
||||
/* Parse a URI and extract the server connection information.
|
||||
* URI scheme is based on the the provisional specification[1] excluding support
|
||||
* for query parameters. Valid URIs are:
|
||||
* scheme: "redis://"
|
||||
* authority: [<username> ":"] <password> "@"] [<hostname> [":" <port>]]
|
||||
* path: ["/" [<db>]]
|
||||
*
|
||||
* [1]: https://www.iana.org/assignments/uri-schemes/prov/redis */
|
||||
static void parseRedisUri(const char *uri) {
|
||||
|
||||
const char *scheme = "redis://";
|
||||
const char *curr = uri;
|
||||
const char *end = uri + strlen(uri);
|
||||
const char *userinfo, *username, *port, *host, *path;
|
||||
|
||||
/* URI must start with a valid scheme. */
|
||||
if (strncasecmp(scheme, curr, strlen(scheme))) {
|
||||
fprintf(stderr,"Invalid URI scheme\n");
|
||||
exit(1);
|
||||
}
|
||||
curr += strlen(scheme);
|
||||
if (curr == end) return;
|
||||
|
||||
/* Extract user info. */
|
||||
if ((userinfo = strchr(curr,'@'))) {
|
||||
if ((username = strchr(curr, ':')) && username < userinfo) {
|
||||
/* If provided, username is ignored. */
|
||||
curr = username + 1;
|
||||
}
|
||||
|
||||
config.auth = percentDecode(curr, userinfo - curr);
|
||||
curr = userinfo + 1;
|
||||
}
|
||||
if (curr == end) return;
|
||||
|
||||
/* Extract host and port. */
|
||||
path = strchr(curr, '/');
|
||||
if (*curr != '/') {
|
||||
host = path ? path - 1 : end;
|
||||
if ((port = strchr(curr, ':'))) {
|
||||
config.hostport = atoi(port + 1);
|
||||
host = port - 1;
|
||||
}
|
||||
config.hostip = sdsnewlen(curr, host - curr + 1);
|
||||
}
|
||||
curr = path ? path + 1 : end;
|
||||
if (curr == end) return;
|
||||
|
||||
/* Extract database number. */
|
||||
config.dbnum = atoi(curr);
|
||||
}
|
||||
|
||||
/*------------------------------------------------------------------------------
|
||||
* Help functions
|
||||
*--------------------------------------------------------------------------- */
|
||||
@@ -624,7 +711,7 @@ int isColorTerm(void) {
|
||||
return t != NULL && strstr(t,"xterm") != NULL;
|
||||
}
|
||||
|
||||
/* Helpe function for sdsCatColorizedLdbReply() appending colorize strings
|
||||
/* Helper function for sdsCatColorizedLdbReply() appending colorize strings
|
||||
* to an SDS string. */
|
||||
sds sdscatcolor(sds o, char *s, size_t len, char *color) {
|
||||
if (!isColorTerm()) return sdscatlen(o,s,len);
|
||||
@@ -632,7 +719,6 @@ sds sdscatcolor(sds o, char *s, size_t len, char *color) {
|
||||
int bold = strstr(color,"bold") != NULL;
|
||||
int ccode = 37; /* Defaults to white. */
|
||||
if (strstr(color,"red")) ccode = 31;
|
||||
else if (strstr(color,"red")) ccode = 31;
|
||||
else if (strstr(color,"green")) ccode = 32;
|
||||
else if (strstr(color,"yellow")) ccode = 33;
|
||||
else if (strstr(color,"blue")) ccode = 34;
|
||||
@@ -1003,6 +1089,8 @@ static int parseOptions(int argc, char **argv) {
|
||||
config.dbnum = atoi(argv[++i]);
|
||||
} else if (!strcmp(argv[i],"-a") && !lastarg) {
|
||||
config.auth = argv[++i];
|
||||
} else if (!strcmp(argv[i],"-u") && !lastarg) {
|
||||
parseRedisUri(argv[++i]);
|
||||
} else if (!strcmp(argv[i],"--raw")) {
|
||||
config.output = OUTPUT_RAW;
|
||||
} else if (!strcmp(argv[i],"--no-raw")) {
|
||||
@@ -1042,6 +1130,8 @@ static int parseOptions(int argc, char **argv) {
|
||||
config.pipe_timeout = atoi(argv[++i]);
|
||||
} else if (!strcmp(argv[i],"--bigkeys")) {
|
||||
config.bigkeys = 1;
|
||||
} else if (!strcmp(argv[i],"--hotkeys")) {
|
||||
config.hotkeys = 1;
|
||||
} else if (!strcmp(argv[i],"--eval") && !lastarg) {
|
||||
config.eval = argv[++i];
|
||||
} else if (!strcmp(argv[i],"--ldb")) {
|
||||
@@ -1110,6 +1200,7 @@ static void usage(void) {
|
||||
" -p <port> Server port (default: 6379).\n"
|
||||
" -s <socket> Server socket (overrides hostname and port).\n"
|
||||
" -a <password> Password to use when connecting to the server.\n"
|
||||
" -u <uri> Server URI.\n"
|
||||
" -r <repeat> Execute specified command N times.\n"
|
||||
" -i <interval> When -r is used, waits <interval> seconds per command.\n"
|
||||
" It is possible to specify sub-second times like -i 0.1.\n"
|
||||
@@ -1141,6 +1232,8 @@ static void usage(void) {
|
||||
" no reply is received within <n> seconds.\n"
|
||||
" Default timeout: %d. Use 0 to wait forever.\n"
|
||||
" --bigkeys Sample Redis keys looking for big keys.\n"
|
||||
" --hotkeys Sample Redis keys looking for hot keys.\n"
|
||||
" only works when maxmemory-policy is *lfu.\n"
|
||||
" --scan List all keys using the SCAN command.\n"
|
||||
" --pattern <pat> Useful with --scan to specify a SCAN pattern.\n"
|
||||
" --intrinsic-latency <sec> Run a test to measure intrinsic system latency.\n"
|
||||
@@ -2254,6 +2347,129 @@ static void findBigKeys(void) {
|
||||
exit(0);
|
||||
}
|
||||
|
||||
static void getKeyFreqs(redisReply *keys, unsigned long long *freqs) {
|
||||
redisReply *reply;
|
||||
unsigned int i;
|
||||
|
||||
/* Pipeline OBJECT freq commands */
|
||||
for(i=0;i<keys->elements;i++) {
|
||||
redisAppendCommand(context, "OBJECT freq %s", keys->element[i]->str);
|
||||
}
|
||||
|
||||
/* Retrieve freqs */
|
||||
for(i=0;i<keys->elements;i++) {
|
||||
if(redisGetReply(context, (void**)&reply)!=REDIS_OK) {
|
||||
fprintf(stderr, "Error getting freq for key '%s' (%d: %s)\n",
|
||||
keys->element[i]->str, context->err, context->errstr);
|
||||
exit(1);
|
||||
} else if(reply->type != REDIS_REPLY_INTEGER) {
|
||||
if(reply->type == REDIS_REPLY_ERROR) {
|
||||
fprintf(stderr, "Error: %s\n", reply->str);
|
||||
exit(1);
|
||||
} else {
|
||||
fprintf(stderr, "Warning: OBJECT freq on '%s' failed (may have been deleted)\n", keys->element[i]->str);
|
||||
freqs[i] = 0;
|
||||
}
|
||||
} else {
|
||||
freqs[i] = reply->integer;
|
||||
}
|
||||
freeReplyObject(reply);
|
||||
}
|
||||
}
|
||||
|
||||
#define HOTKEYS_SAMPLE 16
|
||||
static void findHotKeys(void) {
|
||||
redisReply *keys, *reply;
|
||||
unsigned long long counters[HOTKEYS_SAMPLE] = {0};
|
||||
sds hotkeys[HOTKEYS_SAMPLE] = {NULL};
|
||||
unsigned long long sampled = 0, total_keys, *freqs = NULL, it = 0;
|
||||
unsigned int arrsize = 0, i, k;
|
||||
double pct;
|
||||
|
||||
/* Total keys pre scanning */
|
||||
total_keys = getDbSize();
|
||||
|
||||
/* Status message */
|
||||
printf("\n# Scanning the entire keyspace to find hot keys as well as\n");
|
||||
printf("# average sizes per key type. You can use -i 0.1 to sleep 0.1 sec\n");
|
||||
printf("# per 100 SCAN commands (not usually needed).\n\n");
|
||||
|
||||
/* SCAN loop */
|
||||
do {
|
||||
/* Calculate approximate percentage completion */
|
||||
pct = 100 * (double)sampled/total_keys;
|
||||
|
||||
/* Grab some keys and point to the keys array */
|
||||
reply = sendScan(&it);
|
||||
keys = reply->element[1];
|
||||
|
||||
/* Reallocate our freqs array if we need to */
|
||||
if(keys->elements > arrsize) {
|
||||
freqs = zrealloc(freqs, sizeof(unsigned long long)*keys->elements);
|
||||
|
||||
if(!freqs) {
|
||||
fprintf(stderr, "Failed to allocate storage for keys!\n");
|
||||
exit(1);
|
||||
}
|
||||
|
||||
arrsize = keys->elements;
|
||||
}
|
||||
|
||||
getKeyFreqs(keys, freqs);
|
||||
|
||||
/* Now update our stats */
|
||||
for(i=0;i<keys->elements;i++) {
|
||||
sampled++;
|
||||
/* Update overall progress */
|
||||
if(sampled % 1000000 == 0) {
|
||||
printf("[%05.2f%%] Sampled %llu keys so far\n", pct, sampled);
|
||||
}
|
||||
|
||||
/* Use eviction pool here */
|
||||
k = 0;
|
||||
while (k < HOTKEYS_SAMPLE && freqs[i] > counters[k]) k++;
|
||||
if (k == 0) continue;
|
||||
k--;
|
||||
if (k == 0 || counters[k] == 0) {
|
||||
sdsfree(hotkeys[k]);
|
||||
} else {
|
||||
sdsfree(hotkeys[0]);
|
||||
memmove(counters,counters+1,sizeof(counters[0])*k);
|
||||
memmove(hotkeys,hotkeys+1,sizeof(hotkeys[0])*k);
|
||||
}
|
||||
counters[k] = freqs[i];
|
||||
hotkeys[k] = sdsnew(keys->element[i]->str);
|
||||
printf(
|
||||
"[%05.2f%%] Hot key '%s' found so far with counter %llu\n",
|
||||
pct, keys->element[i]->str, freqs[i]);
|
||||
}
|
||||
|
||||
/* Sleep if we've been directed to do so */
|
||||
if(sampled && (sampled %100) == 0 && config.interval) {
|
||||
usleep(config.interval);
|
||||
}
|
||||
|
||||
freeReplyObject(reply);
|
||||
} while(it != 0);
|
||||
|
||||
if (freqs) zfree(freqs);
|
||||
|
||||
/* We're done */
|
||||
printf("\n-------- summary -------\n\n");
|
||||
|
||||
printf("Sampled %llu keys in the keyspace!\n", sampled);
|
||||
|
||||
for (i=1; i<= HOTKEYS_SAMPLE; i++) {
|
||||
k = HOTKEYS_SAMPLE - i;
|
||||
if(counters[k]>0) {
|
||||
printf("hot key found with counter: %llu\tkeyname: %s\n", counters[k], hotkeys[k]);
|
||||
sdsfree(hotkeys[k]);
|
||||
}
|
||||
}
|
||||
|
||||
exit(0);
|
||||
}
|
||||
|
||||
/*------------------------------------------------------------------------------
|
||||
* Stats mode
|
||||
*--------------------------------------------------------------------------- */
|
||||
@@ -2364,7 +2580,7 @@ static void statMode(void) {
|
||||
sprintf(buf,"%ld",aux);
|
||||
printf("%-8s",buf);
|
||||
|
||||
/* Requets */
|
||||
/* Requests */
|
||||
aux = getLongInfoField(reply->str,"total_commands_processed");
|
||||
sprintf(buf,"%ld (+%ld)",aux,requests == 0 ? 0 : aux-requests);
|
||||
printf("%-19s",buf);
|
||||
@@ -2631,6 +2847,7 @@ int main(int argc, char **argv) {
|
||||
config.pipe_mode = 0;
|
||||
config.pipe_timeout = REDIS_CLI_DEFAULT_PIPE_TIMEOUT;
|
||||
config.bigkeys = 0;
|
||||
config.hotkeys = 0;
|
||||
config.stdinarg = 0;
|
||||
config.auth = NULL;
|
||||
config.eval = NULL;
|
||||
@@ -2691,6 +2908,12 @@ int main(int argc, char **argv) {
|
||||
findBigKeys();
|
||||
}
|
||||
|
||||
/* Find hot keys */
|
||||
if (config.hotkeys) {
|
||||
if (cliConnect(0) == REDIS_ERR) exit(1);
|
||||
findHotKeys();
|
||||
}
|
||||
|
||||
/* Stat mode */
|
||||
if (config.stat_mode) {
|
||||
if (cliConnect(0) == REDIS_ERR) exit(1);
|
||||
|
||||
+30
-1
@@ -58,6 +58,30 @@
|
||||
#define REDISMODULE_HASH_CFIELDS (1<<2)
|
||||
#define REDISMODULE_HASH_EXISTS (1<<3)
|
||||
|
||||
/* Context Flags: Info about the current context returned by RM_GetContextFlags */
|
||||
|
||||
/* The command is running in the context of a Lua script */
|
||||
#define REDISMODULE_CTX_FLAGS_LUA 0x0001
|
||||
/* The command is running inside a Redis transaction */
|
||||
#define REDISMODULE_CTX_FLAGS_MULTI 0x0002
|
||||
/* The instance is a master */
|
||||
#define REDISMODULE_CTX_FLAGS_MASTER 0x0004
|
||||
/* The instance is a slave */
|
||||
#define REDISMODULE_CTX_FLAGS_SLAVE 0x0008
|
||||
/* The instance is read-only (usually meaning it's a slave as well) */
|
||||
#define REDISMODULE_CTX_FLAGS_READONLY 0x0010
|
||||
/* The instance is running in cluster mode */
|
||||
#define REDISMODULE_CTX_FLAGS_CLUSTER 0x0020
|
||||
/* The instance has AOF enabled */
|
||||
#define REDISMODULE_CTX_FLAGS_AOF 0x0040 //
|
||||
/* The instance has RDB enabled */
|
||||
#define REDISMODULE_CTX_FLAGS_RDB 0x0080 //
|
||||
/* The instance has Maxmemory set */
|
||||
#define REDISMODULE_CTX_FLAGS_MAXMEMORY 0x0100
|
||||
/* Maxmemory is set and has an eviction policy that may delete keys */
|
||||
#define REDISMODULE_CTX_FLAGS_EVICT 0x0200
|
||||
|
||||
|
||||
/* A special pointer that we can use between the core and the module to signal
|
||||
* field deletion, and that is impossible to be a valid pointer. */
|
||||
#define REDISMODULE_HASH_DELETE ((RedisModuleString*)(long)1)
|
||||
@@ -119,7 +143,8 @@ void *REDISMODULE_API_FUNC(RedisModule_Calloc)(size_t nmemb, size_t size);
|
||||
char *REDISMODULE_API_FUNC(RedisModule_Strdup)(const char *str);
|
||||
int REDISMODULE_API_FUNC(RedisModule_GetApi)(const char *, void *);
|
||||
int REDISMODULE_API_FUNC(RedisModule_CreateCommand)(RedisModuleCtx *ctx, const char *name, RedisModuleCmdFunc cmdfunc, const char *strflags, int firstkey, int lastkey, int keystep);
|
||||
int REDISMODULE_API_FUNC(RedisModule_SetModuleAttribs)(RedisModuleCtx *ctx, const char *name, int ver, int apiver);
|
||||
void REDISMODULE_API_FUNC(RedisModule_SetModuleAttribs)(RedisModuleCtx *ctx, const char *name, int ver, int apiver);
|
||||
int REDISMODULE_API_FUNC(RedisModule_IsModuleNameBusy)(const char *name);
|
||||
int REDISMODULE_API_FUNC(RedisModule_WrongArity)(RedisModuleCtx *ctx);
|
||||
int REDISMODULE_API_FUNC(RedisModule_ReplyWithLongLong)(RedisModuleCtx *ctx, long long ll);
|
||||
int REDISMODULE_API_FUNC(RedisModule_GetSelectedDb)(RedisModuleCtx *ctx);
|
||||
@@ -183,6 +208,7 @@ int REDISMODULE_API_FUNC(RedisModule_HashGet)(RedisModuleKey *key, int flags, ..
|
||||
int REDISMODULE_API_FUNC(RedisModule_IsKeysPositionRequest)(RedisModuleCtx *ctx);
|
||||
void REDISMODULE_API_FUNC(RedisModule_KeyAtPos)(RedisModuleCtx *ctx, int pos);
|
||||
unsigned long long REDISMODULE_API_FUNC(RedisModule_GetClientId)(RedisModuleCtx *ctx);
|
||||
int REDISMODULE_API_FUNC(RedisModule_GetContextFlags)(RedisModuleCtx *ctx);
|
||||
void *REDISMODULE_API_FUNC(RedisModule_PoolAlloc)(RedisModuleCtx *ctx, size_t bytes);
|
||||
RedisModuleType *REDISMODULE_API_FUNC(RedisModule_CreateDataType)(RedisModuleCtx *ctx, const char *name, int encver, RedisModuleTypeMethods *typemethods);
|
||||
int REDISMODULE_API_FUNC(RedisModule_ModuleTypeSetValue)(RedisModuleKey *key, RedisModuleType *mt, void *value);
|
||||
@@ -238,6 +264,7 @@ static int RedisModule_Init(RedisModuleCtx *ctx, const char *name, int ver, int
|
||||
REDISMODULE_GET_API(Strdup);
|
||||
REDISMODULE_GET_API(CreateCommand);
|
||||
REDISMODULE_GET_API(SetModuleAttribs);
|
||||
REDISMODULE_GET_API(IsModuleNameBusy);
|
||||
REDISMODULE_GET_API(WrongArity);
|
||||
REDISMODULE_GET_API(ReplyWithLongLong);
|
||||
REDISMODULE_GET_API(ReplyWithError);
|
||||
@@ -302,6 +329,7 @@ static int RedisModule_Init(RedisModuleCtx *ctx, const char *name, int ver, int
|
||||
REDISMODULE_GET_API(IsKeysPositionRequest);
|
||||
REDISMODULE_GET_API(KeyAtPos);
|
||||
REDISMODULE_GET_API(GetClientId);
|
||||
REDISMODULE_GET_API(GetContextFlags);
|
||||
REDISMODULE_GET_API(PoolAlloc);
|
||||
REDISMODULE_GET_API(CreateDataType);
|
||||
REDISMODULE_GET_API(ModuleTypeSetValue);
|
||||
@@ -344,6 +372,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;
|
||||
RedisModule_SetModuleAttribs(ctx,name,ver,apiver);
|
||||
return REDISMODULE_OK;
|
||||
}
|
||||
|
||||
+35
-12
@@ -569,18 +569,19 @@ int startBgsaveForReplication(int mincapa) {
|
||||
serverLog(LL_NOTICE,"Starting BGSAVE for SYNC with target: %s",
|
||||
socket_target ? "slaves sockets" : "disk");
|
||||
|
||||
rdbSaveInfo rsi = RDB_SAVE_INFO_INIT;
|
||||
/* If we are saving for a chained slave (that is, if we are,
|
||||
* in turn, a slave of another instance), make sure after
|
||||
* loadig the RDB, our slaves select the right DB: we'll just
|
||||
* send the replication stream we receive from our master, so
|
||||
* no way to send SELECT commands. */
|
||||
if (server.master) rsi.repl_stream_db = server.master->db->id;
|
||||
|
||||
if (socket_target)
|
||||
retval = rdbSaveToSlavesSockets(&rsi);
|
||||
else
|
||||
retval = rdbSaveBackground(server.rdb_filename,&rsi);
|
||||
rdbSaveInfo rsi, *rsiptr;
|
||||
rsiptr = rdbPopulateSaveInfo(&rsi);
|
||||
/* Only do rdbSave* when rsiptr is not NULL,
|
||||
* otherwise slave will miss repl-stream-db. */
|
||||
if (rsiptr) {
|
||||
if (socket_target)
|
||||
retval = rdbSaveToSlavesSockets(rsiptr);
|
||||
else
|
||||
retval = rdbSaveBackground(server.rdb_filename,rsiptr);
|
||||
} else {
|
||||
serverLog(LL_WARNING,"BGSAVE for replication: replication information not available, can't generate the RDB file right now. Try later.");
|
||||
retval = C_ERR;
|
||||
}
|
||||
|
||||
/* If we failed to BGSAVE, remove the slaves waiting for a full
|
||||
* resynchorinization from the list of salves, inform them with
|
||||
@@ -1531,6 +1532,11 @@ int slaveTryPartialResynchronization(int fd, int read_reply) {
|
||||
/* Setup the replication to continue. */
|
||||
sdsfree(reply);
|
||||
replicationResurrectCachedMaster(fd);
|
||||
|
||||
/* If this instance was restarted and we read the metadata to
|
||||
* PSYNC from the persistence file, our replication backlog could
|
||||
* be still not initialized. Create it. */
|
||||
if (server.repl_backlog == NULL) createReplicationBacklog();
|
||||
return PSYNC_CONTINUE;
|
||||
}
|
||||
|
||||
@@ -2607,6 +2613,23 @@ void replicationCron(void) {
|
||||
time_t idle = server.unixtime - server.repl_no_slaves_since;
|
||||
|
||||
if (idle > server.repl_backlog_time_limit) {
|
||||
/* When we free the backlog, we always use a new
|
||||
* replication ID and clear the ID2. This is needed
|
||||
* because when there is no backlog, the master_repl_offset
|
||||
* is not updated, but we would still retain our replication
|
||||
* ID, leading to the following problem:
|
||||
*
|
||||
* 1. We are a master instance.
|
||||
* 2. Our slave is promoted to master. It's repl-id-2 will
|
||||
* be the same as our repl-id.
|
||||
* 3. We, yet as master, receive some updates, that will not
|
||||
* increment the master_repl_offset.
|
||||
* 4. Later we are turned into a slave, connecto to the new
|
||||
* master that will accept our PSYNC request by second
|
||||
* replication ID, but there will be data inconsistency
|
||||
* because we received writes. */
|
||||
changeReplicationId();
|
||||
clearReplicationId2();
|
||||
freeReplicationBacklog();
|
||||
serverLog(LL_NOTICE,
|
||||
"Replication backlog freed after %d seconds "
|
||||
|
||||
+59
-39
@@ -358,6 +358,13 @@ int luaRedisGenericCommand(lua_State *lua, int raise_error) {
|
||||
static size_t cached_objects_len[LUA_CMD_OBJCACHE_SIZE];
|
||||
static int inuse = 0; /* Recursive calls detection. */
|
||||
|
||||
/* Reflect MULTI state */
|
||||
if (server.lua_multi_emitted || (server.lua_caller->flags & CLIENT_MULTI)) {
|
||||
c->flags |= CLIENT_MULTI;
|
||||
} else {
|
||||
c->flags &= ~CLIENT_MULTI;
|
||||
}
|
||||
|
||||
/* By using Lua debug hooks it is possible to trigger a recursive call
|
||||
* to luaRedisGenericCommand(), which normally should never happen.
|
||||
* To make this function reentrant is futile and makes it slower, but
|
||||
@@ -535,6 +542,7 @@ int luaRedisGenericCommand(lua_State *lua, int raise_error) {
|
||||
* a Lua script in the context of AOF and slaves. */
|
||||
if (server.lua_replicate_commands &&
|
||||
!server.lua_multi_emitted &&
|
||||
!(server.lua_caller->flags & CLIENT_MULTI) &&
|
||||
server.lua_write_dirty &&
|
||||
server.lua_repl != PROPAGATE_NONE)
|
||||
{
|
||||
@@ -1133,18 +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>
|
||||
*
|
||||
* 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();
|
||||
* The function increments the reference count of the 'body' object as a
|
||||
* side effect of a successful call.
|
||||
*
|
||||
* 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;
|
||||
|
||||
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);
|
||||
@@ -1152,30 +1180,35 @@ int luaCreateFunction(client *c, lua_State *lua, char *funcname, robj *body) {
|
||||
funcdef = sdscatlen(funcdef,"\nend",4);
|
||||
|
||||
if (luaL_loadbuffer(lua,funcdef,sdslen(funcdef),"@user_script")) {
|
||||
addReplyErrorFormat(c,"Error compiling script (new function): %s\n",
|
||||
lua_tostring(lua,-1));
|
||||
if (c != NULL) {
|
||||
addReplyErrorFormat(c,
|
||||
"Error compiling script (new function): %s\n",
|
||||
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)) {
|
||||
addReplyErrorFormat(c,"Error running script (new function): %s\n",
|
||||
lua_tostring(lua,-1));
|
||||
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,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. */
|
||||
@@ -1274,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 */
|
||||
@@ -1438,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) {
|
||||
|
||||
@@ -248,16 +248,23 @@ sds sdsMakeRoomFor(sds s, size_t addlen) {
|
||||
sds sdsRemoveFreeSpace(sds s) {
|
||||
void *sh, *newsh;
|
||||
char type, oldtype = s[-1] & SDS_TYPE_MASK;
|
||||
int hdrlen;
|
||||
int hdrlen, oldhdrlen = sdsHdrSize(oldtype);
|
||||
size_t len = sdslen(s);
|
||||
sh = (char*)s-sdsHdrSize(oldtype);
|
||||
sh = (char*)s-oldhdrlen;
|
||||
|
||||
/* Check what would be the minimum SDS header that is just good enough to
|
||||
* fit this string. */
|
||||
type = sdsReqType(len);
|
||||
hdrlen = sdsHdrSize(type);
|
||||
if (oldtype==type) {
|
||||
newsh = s_realloc(sh, hdrlen+len+1);
|
||||
|
||||
/* If the type is the same, or at least a large enough type is still
|
||||
* required, we just realloc(), letting the allocator to do the copy
|
||||
* only if really needed. Otherwise if the change is huge, we manually
|
||||
* reallocate the string to use the different header type. */
|
||||
if (oldtype==type || type > SDS_TYPE_8) {
|
||||
newsh = s_realloc(sh, oldhdrlen+len+1);
|
||||
if (newsh == NULL) return NULL;
|
||||
s = (char*)newsh+hdrlen;
|
||||
s = (char*)newsh+oldhdrlen;
|
||||
} else {
|
||||
newsh = s_malloc(hdrlen+len+1);
|
||||
if (newsh == NULL) return NULL;
|
||||
|
||||
+65
-24
@@ -276,7 +276,7 @@ struct redisCommand redisCommandTable[] = {
|
||||
{"readonly",readonlyCommand,1,"F",0,NULL,0,0,0,0,0},
|
||||
{"readwrite",readwriteCommand,1,"F",0,NULL,0,0,0,0,0},
|
||||
{"dump",dumpCommand,2,"r",0,NULL,1,1,1,0,0},
|
||||
{"object",objectCommand,3,"r",0,NULL,2,2,2,0,0},
|
||||
{"object",objectCommand,-2,"r",0,NULL,2,2,2,0,0},
|
||||
{"memory",memoryCommand,-2,"r",0,NULL,0,0,0,0,0},
|
||||
{"client",clientCommand,-2,"as",0,NULL,0,0,0,0,0},
|
||||
{"eval",evalCommand,-3,"s",0,evalGetKeys,0,0,0,0,0},
|
||||
@@ -908,12 +908,15 @@ void databasesCron(void) {
|
||||
/* Rehash */
|
||||
if (server.activerehashing) {
|
||||
for (j = 0; j < dbs_per_call; j++) {
|
||||
int work_done = incrementallyRehash(rehash_db % server.dbnum);
|
||||
rehash_db++;
|
||||
int work_done = incrementallyRehash(rehash_db);
|
||||
if (work_done) {
|
||||
/* If the function did some work, stop here, we'll do
|
||||
* more at the next cron loop. */
|
||||
break;
|
||||
} else {
|
||||
/* If this db didn't need rehash, we'll try the next one. */
|
||||
rehash_db++;
|
||||
rehash_db %= server.dbnum;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1092,7 +1095,9 @@ int serverCron(struct aeEventLoop *eventLoop, long long id, void *clientData) {
|
||||
{
|
||||
serverLog(LL_NOTICE,"%d changes in %d seconds. Saving...",
|
||||
sp->changes, (int)sp->seconds);
|
||||
rdbSaveBackground(server.rdb_filename,NULL);
|
||||
rdbSaveInfo rsi, *rsiptr;
|
||||
rsiptr = rdbPopulateSaveInfo(&rsi);
|
||||
rdbSaveBackground(server.rdb_filename,rsiptr);
|
||||
break;
|
||||
}
|
||||
}
|
||||
@@ -1164,7 +1169,9 @@ int serverCron(struct aeEventLoop *eventLoop, long long id, void *clientData) {
|
||||
(server.unixtime-server.lastbgsave_try > CONFIG_BGSAVE_RETRY_DELAY ||
|
||||
server.lastbgsave_status == C_OK))
|
||||
{
|
||||
if (rdbSaveBackground(server.rdb_filename,NULL) == C_OK)
|
||||
rdbSaveInfo rsi, *rsiptr;
|
||||
rsiptr = rdbPopulateSaveInfo(&rsi);
|
||||
if (rdbSaveBackground(server.rdb_filename,rsiptr) == C_OK)
|
||||
server.rdb_bgsave_scheduled = 0;
|
||||
}
|
||||
|
||||
@@ -1542,16 +1549,29 @@ int restartServer(int flags, mstime_t delay) {
|
||||
|
||||
/* Check if we still have accesses to the executable that started this
|
||||
* server instance. */
|
||||
if (access(server.executable,X_OK) == -1) return C_ERR;
|
||||
if (access(server.executable,X_OK) == -1) {
|
||||
serverLog(LL_WARNING,"Can't restart: this process has no "
|
||||
"permissions to execute %s", server.executable);
|
||||
return C_ERR;
|
||||
}
|
||||
|
||||
/* Config rewriting. */
|
||||
if (flags & RESTART_SERVER_CONFIG_REWRITE &&
|
||||
server.configfile &&
|
||||
rewriteConfig(server.configfile) == -1) return C_ERR;
|
||||
rewriteConfig(server.configfile) == -1)
|
||||
{
|
||||
serverLog(LL_WARNING,"Can't restart: configuration rewrite process "
|
||||
"failed");
|
||||
return C_ERR;
|
||||
}
|
||||
|
||||
/* Perform a proper shutdown. */
|
||||
if (flags & RESTART_SERVER_GRACEFULLY &&
|
||||
prepareForShutdown(SHUTDOWN_NOFLAGS) != C_OK) return C_ERR;
|
||||
prepareForShutdown(SHUTDOWN_NOFLAGS) != C_OK)
|
||||
{
|
||||
serverLog(LL_WARNING,"Can't restart: error preparing for shutdown");
|
||||
return C_ERR;
|
||||
}
|
||||
|
||||
/* Close all file descriptors, with the exception of stdin, stdout, strerr
|
||||
* which are useful if we restart a Redis server which is not daemonized. */
|
||||
@@ -1563,6 +1583,8 @@ int restartServer(int flags, mstime_t delay) {
|
||||
|
||||
/* Execute the server with the original command line. */
|
||||
if (delay) usleep(delay*1000);
|
||||
zfree(server.exec_argv[0]);
|
||||
server.exec_argv[0] = zstrdup(server.executable);
|
||||
execve(server.executable,server.exec_argv,environ);
|
||||
|
||||
/* If an error occurred here, there is nothing we can do, but exit. */
|
||||
@@ -1987,15 +2009,18 @@ void populateCommandTable(void) {
|
||||
}
|
||||
|
||||
void resetCommandTableStats(void) {
|
||||
int numcommands = sizeof(redisCommandTable)/sizeof(struct redisCommand);
|
||||
int j;
|
||||
|
||||
for (j = 0; j < numcommands; j++) {
|
||||
struct redisCommand *c = redisCommandTable+j;
|
||||
struct redisCommand *c;
|
||||
dictEntry *de;
|
||||
dictIterator *di;
|
||||
|
||||
di = dictGetSafeIterator(server.commands);
|
||||
while((de = dictNext(di)) != NULL) {
|
||||
c = (struct redisCommand *) dictGetVal(de);
|
||||
c->microseconds = 0;
|
||||
c->calls = 0;
|
||||
}
|
||||
dictReleaseIterator(di);
|
||||
|
||||
}
|
||||
|
||||
/* ========================== Redis OP Array API ============================ */
|
||||
@@ -2255,8 +2280,9 @@ void call(client *c, int flags) {
|
||||
propagate_flags &= ~PROPAGATE_AOF;
|
||||
|
||||
/* Call propagate() only if at least one of AOF / replication
|
||||
* propagation is needed. */
|
||||
if (propagate_flags != PROPAGATE_NONE)
|
||||
* propagation is needed. Note that modules commands handle replication
|
||||
* in an explicit way, so we never replicate them automatically. */
|
||||
if (propagate_flags != PROPAGATE_NONE && !(c->cmd->flags & CMD_MODULE))
|
||||
propagate(c->cmd,c->db->id,c->argv,c->argc,propagate_flags);
|
||||
}
|
||||
|
||||
@@ -2533,8 +2559,9 @@ int prepareForShutdown(int flags) {
|
||||
"There is a child rewriting the AOF. Killing it!");
|
||||
kill(server.aof_child_pid,SIGUSR1);
|
||||
}
|
||||
/* Append only file: fsync() the AOF and exit */
|
||||
/* Append only file: flush buffers and fsync() the AOF at exit */
|
||||
serverLog(LL_NOTICE,"Calling fsync() on the AOF file.");
|
||||
flushAppendOnlyFile(1);
|
||||
aof_fsync(server.aof_fd);
|
||||
}
|
||||
|
||||
@@ -2542,7 +2569,9 @@ int prepareForShutdown(int flags) {
|
||||
if ((server.saveparamslen > 0 && !nosave) || save) {
|
||||
serverLog(LL_NOTICE,"Saving the final RDB snapshot before exiting.");
|
||||
/* Snapshotting. Perform a SYNC SAVE and exit */
|
||||
if (rdbSave(server.rdb_filename,NULL) != C_OK) {
|
||||
rdbSaveInfo rsi, *rsiptr;
|
||||
rsiptr = rdbPopulateSaveInfo(&rsi);
|
||||
if (rdbSave(server.rdb_filename,rsiptr) != C_OK) {
|
||||
/* Ooops.. error saving! The best we can do is to continue
|
||||
* operating. Note that if there was a background saving process,
|
||||
* in the next cron() Redis will be notified that the background
|
||||
@@ -2794,7 +2823,7 @@ void bytesToHuman(char *s, unsigned long long n) {
|
||||
sds genRedisInfoString(char *section) {
|
||||
sds info = sdsempty();
|
||||
time_t uptime = server.unixtime-server.stat_starttime;
|
||||
int j, numcommands;
|
||||
int j;
|
||||
struct rusage self_ru, c_ru;
|
||||
unsigned long lol, bib;
|
||||
int allsections = 0, defsections = 0;
|
||||
@@ -3255,20 +3284,24 @@ sds genRedisInfoString(char *section) {
|
||||
(float)c_ru.ru_utime.tv_sec+(float)c_ru.ru_utime.tv_usec/1000000);
|
||||
}
|
||||
|
||||
/* cmdtime */
|
||||
/* Command statistics */
|
||||
if (allsections || !strcasecmp(section,"commandstats")) {
|
||||
if (sections++) info = sdscat(info,"\r\n");
|
||||
info = sdscatprintf(info, "# Commandstats\r\n");
|
||||
numcommands = sizeof(redisCommandTable)/sizeof(struct redisCommand);
|
||||
for (j = 0; j < numcommands; j++) {
|
||||
struct redisCommand *c = redisCommandTable+j;
|
||||
|
||||
struct redisCommand *c;
|
||||
dictEntry *de;
|
||||
dictIterator *di;
|
||||
di = dictGetSafeIterator(server.commands);
|
||||
while((de = dictNext(di)) != NULL) {
|
||||
c = (struct redisCommand *) dictGetVal(de);
|
||||
if (!c->calls) continue;
|
||||
info = sdscatprintf(info,
|
||||
"cmdstat_%s:calls=%lld,usec=%lld,usec_per_call=%.2f\r\n",
|
||||
c->name, c->calls, c->microseconds,
|
||||
(c->calls == 0) ? 0 : ((float)c->microseconds/c->calls));
|
||||
}
|
||||
dictReleaseIterator(di);
|
||||
}
|
||||
|
||||
/* Cluster */
|
||||
@@ -3518,13 +3551,21 @@ void loadDataFromDisk(void) {
|
||||
(float)(ustime()-start)/1000000);
|
||||
|
||||
/* Restore the replication ID / offset from the RDB file. */
|
||||
if (rsi.repl_id_is_set && rsi.repl_offset != -1) {
|
||||
if (server.masterhost &&
|
||||
rsi.repl_id_is_set &&
|
||||
rsi.repl_offset != -1 &&
|
||||
/* Note that older implementations may save a repl_stream_db
|
||||
* of -1 inside the RDB file in a wrong way, see more information
|
||||
* in function rdbPopulateSaveInfo. */
|
||||
rsi.repl_stream_db != -1)
|
||||
{
|
||||
memcpy(server.replid,rsi.repl_id,sizeof(server.replid));
|
||||
server.master_repl_offset = rsi.repl_offset;
|
||||
/* If we are a slave, create a cached master from this
|
||||
* information, in order to allow partial resynchronizations
|
||||
* with masters. */
|
||||
if (server.masterhost) replicationCacheMasterUsingMyself();
|
||||
replicationCacheMasterUsingMyself();
|
||||
selectDb(server.cached_master,rsi.repl_stream_db);
|
||||
}
|
||||
} else if (errno != ENOENT) {
|
||||
serverLog(LL_WARNING,"Fatal error loading the DB: %s. Exiting.",strerror(errno));
|
||||
|
||||
+5
-3
@@ -586,7 +586,7 @@ typedef struct redisObject {
|
||||
unsigned encoding:4;
|
||||
unsigned lru:LRU_BITS; /* LRU time (relative to global lru_clock) or
|
||||
* LFU data (least significant 8 bits frequency
|
||||
* and most significant 16 bits decreas time). */
|
||||
* and most significant 16 bits access time). */
|
||||
int refcount;
|
||||
void *ptr;
|
||||
} robj;
|
||||
@@ -1118,8 +1118,8 @@ struct redisServer {
|
||||
unsigned long long maxmemory; /* Max number of memory bytes to use */
|
||||
int maxmemory_policy; /* Policy for key eviction */
|
||||
int maxmemory_samples; /* Pricision of random sampling */
|
||||
unsigned int lfu_log_factor; /* LFU logarithmic counter factor. */
|
||||
unsigned int lfu_decay_time; /* LFU counter decay factor. */
|
||||
int lfu_log_factor; /* LFU logarithmic counter factor. */
|
||||
int lfu_decay_time; /* LFU counter decay factor. */
|
||||
/* 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,6 +1781,7 @@ void scriptingInit(int setup);
|
||||
int ldbRemoveChild(pid_t pid);
|
||||
void ldbKillForkedSessions(void);
|
||||
int ldbPendingChildren(void);
|
||||
sds luaCreateFunction(client *c, lua_State *lua, robj *body);
|
||||
|
||||
/* Blocked clients */
|
||||
void processUnblockedClients(void);
|
||||
@@ -1802,6 +1803,7 @@ void evictionPoolAlloc(void);
|
||||
#define LFU_INIT_VAL 5
|
||||
unsigned long LFUGetTimeInMinutes(void);
|
||||
uint8_t LFULogIncr(uint8_t value);
|
||||
unsigned long LFUDecrAndReturn(robj *o);
|
||||
|
||||
/* Keys hashing / comparison functions for dict.c hash tables. */
|
||||
uint64_t dictSdsHash(const void *key);
|
||||
|
||||
+5
-1
@@ -39,7 +39,11 @@
|
||||
#include <errno.h> /* errno program_invocation_name program_invocation_short_name */
|
||||
|
||||
#if !defined(HAVE_SETPROCTITLE)
|
||||
#define HAVE_SETPROCTITLE (defined __NetBSD__ || defined __FreeBSD__ || defined __OpenBSD__)
|
||||
#if (defined __NetBSD__ || defined __FreeBSD__ || defined __OpenBSD__)
|
||||
#define HAVE_SETPROCTITLE 1
|
||||
#else
|
||||
#define HAVE_SETPROCTITLE 0
|
||||
#endif
|
||||
#endif
|
||||
|
||||
|
||||
|
||||
+9
-2
@@ -72,9 +72,16 @@ slowlogEntry *slowlogCreateEntry(client *c, robj **argv, int argc, long long dur
|
||||
(unsigned long)
|
||||
sdslen(argv[j]->ptr) - SLOWLOG_ENTRY_MAX_STRING);
|
||||
se->argv[j] = createObject(OBJ_STRING,s);
|
||||
} else {
|
||||
} else if (argv[j]->refcount == OBJ_SHARED_REFCOUNT) {
|
||||
se->argv[j] = argv[j];
|
||||
incrRefCount(argv[j]);
|
||||
} else {
|
||||
/* Here we need to dupliacate the string objects composing the
|
||||
* argument vector of the command, because those may otherwise
|
||||
* end shared with string objects stored into keys. Having
|
||||
* shared objects between any part of Redis, and the data
|
||||
* structure holding the data, is a problem: FLUSHALL ASYNC
|
||||
* may release the shared string object and create a race. */
|
||||
se->argv[j] = dupStringObject(argv[j]);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+2
-2
@@ -287,8 +287,8 @@ int hashTypeDelete(robj *o, sds field) {
|
||||
if (fptr != NULL) {
|
||||
fptr = ziplistFind(fptr, (unsigned char*)field, sdslen(field), 1);
|
||||
if (fptr != NULL) {
|
||||
zl = ziplistDelete(zl,&fptr);
|
||||
zl = ziplistDelete(zl,&fptr);
|
||||
zl = ziplistDelete(zl,&fptr); /* Delete the key. */
|
||||
zl = ziplistDelete(zl,&fptr); /* Delete the value. */
|
||||
o->ptr = zl;
|
||||
deleted = 1;
|
||||
}
|
||||
|
||||
+1
-1
@@ -1 +1 @@
|
||||
#define REDIS_VERSION "4.0.1"
|
||||
#define REDIS_VERSION "4.0.6"
|
||||
|
||||
+1
-1
@@ -318,7 +318,7 @@ proc end_tests {} {
|
||||
puts "GOOD! No errors."
|
||||
exit 0
|
||||
} else {
|
||||
puts "WARNING $::failed tests faield."
|
||||
puts "WARNING $::failed test(s) failed."
|
||||
exit 1
|
||||
}
|
||||
}
|
||||
|
||||
@@ -10,7 +10,7 @@ start_server {} {
|
||||
# Config
|
||||
set debug_msg 0 ; # Enable additional debug messages
|
||||
|
||||
set no_exit 0; ; # Do not exit at end of the test
|
||||
set no_exit 0 ; # Do not exit at end of the test
|
||||
|
||||
set duration 20 ; # Total test seconds
|
||||
|
||||
@@ -175,6 +175,69 @@ start_server {} {
|
||||
assert {$sync_count == $new_sync_count}
|
||||
}
|
||||
|
||||
test "PSYNC2: Slave RDB restart with EVALSHA in backlog issue #4483" {
|
||||
# Pick a random slave
|
||||
set slave_id [expr {($master_id+1)%5}]
|
||||
set sync_count [status $R($master_id) sync_full]
|
||||
|
||||
# Make sure to replicate the first EVAL while the salve is online
|
||||
# so that it's part of the scripts the master believes it's safe
|
||||
# to propagate as EVALSHA.
|
||||
$R($master_id) EVAL {return redis.call("incr","__mycounter")} 0
|
||||
$R($master_id) EVALSHA e6e0b547500efcec21eddb619ac3724081afee89 0
|
||||
|
||||
# Wait for the two to sync
|
||||
wait_for_condition 50 1000 {
|
||||
[$R($master_id) debug digest] == [$R($slave_id) debug digest]
|
||||
} else {
|
||||
fail "Slave not reconnecting"
|
||||
}
|
||||
|
||||
# Prevent the slave from receiving master updates, and at
|
||||
# the same time send a new script several times to the
|
||||
# master, so that we'll end with EVALSHA into the backlog.
|
||||
$R($slave_id) slaveof 127.0.0.1 0
|
||||
|
||||
$R($master_id) EVALSHA e6e0b547500efcec21eddb619ac3724081afee89 0
|
||||
$R($master_id) EVALSHA e6e0b547500efcec21eddb619ac3724081afee89 0
|
||||
$R($master_id) EVALSHA e6e0b547500efcec21eddb619ac3724081afee89 0
|
||||
|
||||
catch {
|
||||
$R($slave_id) config rewrite
|
||||
$R($slave_id) debug restart
|
||||
}
|
||||
|
||||
# Reconfigure the slave correctly again, when it's back online.
|
||||
set retry 50
|
||||
while {$retry} {
|
||||
if {[catch {
|
||||
$R($slave_id) slaveof $master_host $master_port
|
||||
}]} {
|
||||
after 1000
|
||||
} else {
|
||||
break
|
||||
}
|
||||
incr retry -1
|
||||
}
|
||||
|
||||
# The master should be back at 4 slaves eventually
|
||||
wait_for_condition 50 1000 {
|
||||
[status $R($master_id) connected_slaves] == 4
|
||||
} else {
|
||||
fail "Slave not reconnecting"
|
||||
}
|
||||
set new_sync_count [status $R($master_id) sync_full]
|
||||
assert {$sync_count == $new_sync_count}
|
||||
|
||||
# However if the slave started with the full state of the
|
||||
# scripting engine, we should now have the same digest.
|
||||
wait_for_condition 50 1000 {
|
||||
[$R($master_id) debug digest] == [$R($slave_id) debug digest]
|
||||
} else {
|
||||
fail "Debug digest mismatch between master and slave in post-restart handshake"
|
||||
}
|
||||
}
|
||||
|
||||
if {$no_exit} {
|
||||
while 1 { puts -nonewline .; flush stdout; after 1000}
|
||||
}
|
||||
|
||||
@@ -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"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -47,4 +47,18 @@ start_server {tags {"latency-monitor"}} {
|
||||
assert {[r latency reset] > 0}
|
||||
assert {[r latency latest] eq {}}
|
||||
}
|
||||
|
||||
test {LATENCY of expire events are correctly collected} {
|
||||
r config set latency-monitor-threshold 20
|
||||
r eval {
|
||||
local i = 0
|
||||
while (i < 1000000) do
|
||||
redis.call('sadd','mybigkey',i)
|
||||
i = i+1
|
||||
end
|
||||
} 0
|
||||
r pexpire mybigkey 1
|
||||
after 500
|
||||
assert_match {*expire-cycle*} [r latency latest]
|
||||
}
|
||||
}
|
||||
|
||||
@@ -144,4 +144,11 @@ start_server {tags {"incr"}} {
|
||||
r set foo 1
|
||||
roundFloat [r incrbyfloat foo -1.1]
|
||||
} {-0.1}
|
||||
|
||||
test {string to double with null terminator} {
|
||||
r set foo 1
|
||||
r setrange foo 2 2
|
||||
catch {r incrbyfloat foo 1} err
|
||||
format $err
|
||||
} {ERR*valid*}
|
||||
}
|
||||
|
||||
@@ -696,6 +696,10 @@ start_server {tags {"zset"}} {
|
||||
}
|
||||
}
|
||||
|
||||
test "ZSET commands don't accept the empty strings as valid score" {
|
||||
assert_error "*not*float*" {r zadd myzset "" abc}
|
||||
}
|
||||
|
||||
proc stressers {encoding} {
|
||||
if {$encoding == "ziplist"} {
|
||||
# Little extra to allow proper fuzzing in the sorting stresser
|
||||
|
||||
Reference in New Issue
Block a user