Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
a767d84a72 | ||
|
|
ee91248505 | ||
|
|
3cf3beff2c | ||
|
|
54eb66495f | ||
|
|
77f91e09cf | ||
|
|
d60953bf25 | ||
|
|
89c68ba3f7 | ||
|
|
fcda82930e | ||
|
|
cefd33925c | ||
|
|
10202ba1fd | ||
|
|
97816fd63e | ||
|
|
3b792f5100 | ||
|
|
da2906e507 | ||
|
|
18264d641b | ||
|
|
918c9aa58c | ||
|
|
8cc6698567 | ||
|
|
e9c9e4c2af | ||
|
|
ee4696b150 | ||
|
|
c9e370c6b8 | ||
|
|
a3ca53e4a7 | ||
|
|
7a62eb96ef | ||
|
|
0efb93d0c0 |
+114
@@ -11,6 +11,120 @@ CRITICAL: There is a critical bug affecting MOST USERS. Upgrade ASAP.
|
||||
SECURITY: There are security fixes in the release.
|
||||
--------------------------------------------------------------------------------
|
||||
|
||||
================================================================================
|
||||
Redis 5.0.10 Released Mon Oct 26 09:21:49 IST 2020
|
||||
================================================================================
|
||||
|
||||
Upgrade urgency: SECURITY if you use an affected platform (see below).
|
||||
Otherwise the upgrade urgency is MODERATE.
|
||||
|
||||
This release fixes a potential heap overflow when using a heap allocator other
|
||||
than jemalloc or glibc's malloc. See:
|
||||
https://github.com/redis/redis/pull/7963
|
||||
|
||||
Other fixes in this release:
|
||||
|
||||
* Avoid case of Lua scripts being consistently aborted due to OOM
|
||||
* XPENDING will not update consumer's seen-time
|
||||
* A blocked XREADGROUP didn't propagated the XSETID to replicas / AOF
|
||||
* UNLINK support for streams
|
||||
* RESTORE ABSTTL won't store expired keys into the DB
|
||||
* Hide AUTH from MONITOR
|
||||
* Cluster: reduce spurious PFAIL/FAIL states upon delayed PONG receival
|
||||
* Cluster: Fix case of clusters mixing accidentally by gossip
|
||||
* Cluster: Allow blocked XREAD on a cluster replica
|
||||
* Cluster: Optimize memory usage CLUSTER SLOTS command
|
||||
* RedisModule_ValueLength support for stream data type
|
||||
* Minor fixes in redis-check-rdb and redis-cli
|
||||
* Fix redis-check-rdb support for modules aux data
|
||||
* Add fsync in replica when full RDB payload was received
|
||||
|
||||
Full list of commits:
|
||||
|
||||
Yossi Gottlieb in commit ce0d74d8f:
|
||||
Fix wrong zmalloc_size() assumption. (#7963)
|
||||
1 file changed, 3 deletions(-)
|
||||
|
||||
Yossi Gottlieb in commit 066699240:
|
||||
Backport Lua 5.2.2 stack overflow fix. (#7733)
|
||||
1 file changed, 1 insertion(+), 1 deletion(-)
|
||||
|
||||
WuYunlong in commit 8a90c7ef3:
|
||||
Add fsync to readSyncBulkPayload(). (#7839)
|
||||
1 file changed, 11 insertions(+)
|
||||
|
||||
Ariel Shtul in commit f0df2bb3c:
|
||||
Fix redis-check-rdb support for modules aux data (#7826)
|
||||
3 files changed, 21 insertions(+), 1 deletion(-)
|
||||
|
||||
hwware in commit 7add2a412:
|
||||
fix memory leak in sentinel connection sharing
|
||||
1 file changed, 1 insertion(+)
|
||||
|
||||
Oran Agra in commit 315e648f8:
|
||||
Allow blocked XREAD on a cluster replica (#7881)
|
||||
3 files changed, 100 insertions(+), 2 deletions(-)
|
||||
|
||||
guybe7 in commit 4967ee94e:
|
||||
Modules: Invalidate saved_oparray after use (#7688)
|
||||
1 file changed, 2 insertions(+)
|
||||
|
||||
antirez in commit 065003e8f:
|
||||
Modules: remove spurious call from moduleHandleBlockedClients().
|
||||
1 file changed, 1 deletion(-)
|
||||
|
||||
Angus Pearson in commit 6cdf62928:
|
||||
Fix broken interval and repeat bahaviour in redis-cli (incluing cluster mode)
|
||||
1 file changed, 11 insertions(+), 6 deletions(-)
|
||||
|
||||
antirez in commit cb6a4971c:
|
||||
Cluster: introduce data_received field.
|
||||
2 files changed, 27 insertions(+), 10 deletions(-)
|
||||
|
||||
Madelyn Olson in commit 83f4de865:
|
||||
Hide AUTH from monitor
|
||||
1 file changed, 1 insertion(+), 1 deletion(-)
|
||||
|
||||
Guy Benoish in commit 3ba08d185:
|
||||
Support streams in general module API functions
|
||||
3 files changed, 11 insertions(+), 1 deletion(-)
|
||||
|
||||
Itamar Haber in commit 109c0635c:
|
||||
Expands lazyfree's effort estimate to include Streams (#5794)
|
||||
1 file changed, 24 insertions(+)
|
||||
|
||||
huangzhw in commit 235210d5b:
|
||||
defrag.c activeDefragSdsListAndDict when defrag sdsele, We can't use (#7492)
|
||||
1 file changed, 1 insertion(+), 1 deletion(-)
|
||||
|
||||
Oran Agra in commit fdd3162fe:
|
||||
RESTORE ABSTTL skip expired keys - leak (#7511)
|
||||
1 file changed, 1 insertion(+)
|
||||
|
||||
Oran Agra in commit 6139d6d18:
|
||||
RESTORE ABSTTL won't store expired keys into the db (#7472)
|
||||
4 files changed, 45 insertions(+), 15 deletions(-)
|
||||
|
||||
Liu Zhen in commit 0f502c58d:
|
||||
fix clusters mixing accidentally by gossip
|
||||
1 file changed, 10 insertions(+), 2 deletions(-)
|
||||
|
||||
Guy Benoish in commit 37fd50718:
|
||||
XPENDING should not update consumer's seen-time
|
||||
4 files changed, 30 insertions(+), 18 deletions(-)
|
||||
|
||||
antirez in commit a3ca53e4a:
|
||||
Also use propagate() in streamPropagateGroupID().
|
||||
1 file changed, 11 insertions(+), 1 deletion(-)
|
||||
|
||||
yanhui13 in commit 7a62eb96e:
|
||||
optimize the output of cluster slots
|
||||
1 file changed, 7 insertions(+), 4 deletions(-)
|
||||
|
||||
srzhao in commit 0efb93d0c:
|
||||
Check OOM at script start to get stable lua OOM state.
|
||||
3 files changed, 11 insertions(+), 4 deletions(-)
|
||||
|
||||
================================================================================
|
||||
Redis 5.0.9 Released Thu Apr 17 12:41:00 CET 2020
|
||||
================================================================================
|
||||
|
||||
Vendored
+1
-1
@@ -274,7 +274,7 @@ int luaD_precall (lua_State *L, StkId func, int nresults) {
|
||||
CallInfo *ci;
|
||||
StkId st, base;
|
||||
Proto *p = cl->p;
|
||||
luaD_checkstack(L, p->maxstacksize);
|
||||
luaD_checkstack(L, p->maxstacksize + p->numparams);
|
||||
func = restorestack(L, funcr);
|
||||
if (!p->is_vararg) { /* no varargs? */
|
||||
base = func + 1;
|
||||
|
||||
+1
-1
@@ -452,7 +452,7 @@ void handleClientsBlockedOnKeys(void) {
|
||||
if (group) {
|
||||
consumer = streamLookupConsumer(group,
|
||||
receiver->bpop.xread_consumer->ptr,
|
||||
1);
|
||||
SLC_NONE);
|
||||
noack = receiver->bpop.xread_group_noack;
|
||||
}
|
||||
|
||||
|
||||
+74
-23
@@ -721,6 +721,7 @@ clusterNode *createClusterNode(char *nodename, int flags) {
|
||||
node->slaves = NULL;
|
||||
node->slaveof = NULL;
|
||||
node->ping_sent = node->pong_received = 0;
|
||||
node->data_received = 0;
|
||||
node->fail_time = 0;
|
||||
node->link = NULL;
|
||||
memset(node->ip,0,sizeof(node->ip));
|
||||
@@ -1434,7 +1435,10 @@ void clusterProcessGossipSection(clusterMsg *hdr, clusterLink *link) {
|
||||
}
|
||||
} else {
|
||||
/* If it's not in NOADDR state and we don't have it, we
|
||||
* start a handshake process against this IP/PORT pairs.
|
||||
* add it to our trusted dict with exact nodeid and flag.
|
||||
* Note that we cannot simply start a handshake against
|
||||
* this IP/PORT pairs, since IP/PORT can be reused already,
|
||||
* otherwise we risk joining another cluster.
|
||||
*
|
||||
* Note that we require that the sender of this gossip message
|
||||
* is a well known node in our cluster, otherwise we risk
|
||||
@@ -1443,7 +1447,12 @@ void clusterProcessGossipSection(clusterMsg *hdr, clusterLink *link) {
|
||||
!(flags & CLUSTER_NODE_NOADDR) &&
|
||||
!clusterBlacklistExists(g->nodename))
|
||||
{
|
||||
clusterStartHandshake(g->ip,ntohs(g->port),ntohs(g->cport));
|
||||
clusterNode *node;
|
||||
node = createClusterNode(g->nodename, flags);
|
||||
memcpy(node->ip,g->ip,NET_IP_STR_LEN);
|
||||
node->port = ntohs(g->port);
|
||||
node->cport = ntohs(g->cport);
|
||||
clusterAddNode(node);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1650,6 +1659,7 @@ int clusterProcessPacket(clusterLink *link) {
|
||||
clusterMsg *hdr = (clusterMsg*) link->rcvbuf;
|
||||
uint32_t totlen = ntohl(hdr->totlen);
|
||||
uint16_t type = ntohs(hdr->type);
|
||||
mstime_t now = mstime();
|
||||
|
||||
if (type < CLUSTERMSG_TYPE_COUNT)
|
||||
server.cluster->stats_bus_messages_received[type]++;
|
||||
@@ -1713,6 +1723,13 @@ int clusterProcessPacket(clusterLink *link) {
|
||||
|
||||
/* Check if the sender is a known node. */
|
||||
sender = clusterLookupNode(hdr->sender);
|
||||
|
||||
/* Update the last time we saw any data from this node. We
|
||||
* use this in order to avoid detecting a timeout from a node that
|
||||
* is just sending a lot of data in the cluster bus, for instance
|
||||
* because of Pub/Sub. */
|
||||
if (sender) sender->data_received = now;
|
||||
|
||||
if (sender && !nodeInHandshake(sender)) {
|
||||
/* Update our curretEpoch if we see a newer epoch in the cluster. */
|
||||
senderCurrentEpoch = ntohu64(hdr->currentEpoch);
|
||||
@@ -1727,7 +1744,7 @@ int clusterProcessPacket(clusterLink *link) {
|
||||
}
|
||||
/* Update the replication offset info for this node. */
|
||||
sender->repl_offset = ntohu64(hdr->offset);
|
||||
sender->repl_offset_time = mstime();
|
||||
sender->repl_offset_time = now;
|
||||
/* If we are a slave performing a manual failover and our master
|
||||
* sent its offset while already paused, populate the MF state. */
|
||||
if (server.cluster->mf_end &&
|
||||
@@ -1841,7 +1858,7 @@ int clusterProcessPacket(clusterLink *link) {
|
||||
* address. */
|
||||
serverLog(LL_DEBUG,"PONG contains mismatching sender ID. About node %.40s added %d ms ago, having flags %d",
|
||||
link->node->name,
|
||||
(int)(mstime()-(link->node->ctime)),
|
||||
(int)(now-(link->node->ctime)),
|
||||
link->node->flags);
|
||||
link->node->flags |= CLUSTER_NODE_NOADDR;
|
||||
link->node->ip[0] = '\0';
|
||||
@@ -1876,7 +1893,7 @@ int clusterProcessPacket(clusterLink *link) {
|
||||
|
||||
/* Update our info about the node */
|
||||
if (link->node && type == CLUSTERMSG_TYPE_PONG) {
|
||||
link->node->pong_received = mstime();
|
||||
link->node->pong_received = now;
|
||||
link->node->ping_sent = 0;
|
||||
|
||||
/* The PFAIL condition can be reversed without external
|
||||
@@ -2023,7 +2040,7 @@ int clusterProcessPacket(clusterLink *link) {
|
||||
"FAIL message received from %.40s about %.40s",
|
||||
hdr->sender, hdr->data.fail.about.nodename);
|
||||
failing->flags |= CLUSTER_NODE_FAIL;
|
||||
failing->fail_time = mstime();
|
||||
failing->fail_time = now;
|
||||
failing->flags &= ~CLUSTER_NODE_PFAIL;
|
||||
clusterDoBeforeSleep(CLUSTER_TODO_SAVE_CONFIG|
|
||||
CLUSTER_TODO_UPDATE_STATE);
|
||||
@@ -2076,9 +2093,9 @@ int clusterProcessPacket(clusterLink *link) {
|
||||
/* Manual failover requested from slaves. Initialize the state
|
||||
* accordingly. */
|
||||
resetManualFailover();
|
||||
server.cluster->mf_end = mstime() + CLUSTER_MF_TIMEOUT;
|
||||
server.cluster->mf_end = now + CLUSTER_MF_TIMEOUT;
|
||||
server.cluster->mf_slave = sender;
|
||||
pauseClients(mstime()+(CLUSTER_MF_TIMEOUT*2));
|
||||
pauseClients(now+(CLUSTER_MF_TIMEOUT*CLUSTER_MF_PAUSE_MULT));
|
||||
serverLog(LL_WARNING,"Manual failover requested by replica %.40s.",
|
||||
sender->name);
|
||||
} else if (type == CLUSTERMSG_TYPE_UPDATE) {
|
||||
@@ -3486,7 +3503,6 @@ void clusterCron(void) {
|
||||
while((de = dictNext(di)) != NULL) {
|
||||
clusterNode *node = dictGetVal(de);
|
||||
now = mstime(); /* Use an updated time at every iteration. */
|
||||
mstime_t delay;
|
||||
|
||||
if (node->flags &
|
||||
(CLUSTER_NODE_MYSELF|CLUSTER_NODE_NOADDR|CLUSTER_NODE_HANDSHAKE))
|
||||
@@ -3510,7 +3526,7 @@ void clusterCron(void) {
|
||||
this_slaves = okslaves;
|
||||
}
|
||||
|
||||
/* If we are waiting for the PONG more than half the cluster
|
||||
/* If we are not receiving any data for more than half the cluster
|
||||
* timeout, reconnect the link: maybe there is a connection
|
||||
* issue even if the node is alive. */
|
||||
if (node->link && /* is connected */
|
||||
@@ -3519,7 +3535,9 @@ void clusterCron(void) {
|
||||
node->ping_sent && /* we already sent a ping */
|
||||
node->pong_received < node->ping_sent && /* still waiting pong */
|
||||
/* and we are waiting for the pong more than timeout/2 */
|
||||
now - node->ping_sent > server.cluster_node_timeout/2)
|
||||
now - node->ping_sent > server.cluster_node_timeout/2 &&
|
||||
/* and in such interval we are not seeing any traffic at all. */
|
||||
now - node->data_received > server.cluster_node_timeout/2)
|
||||
{
|
||||
/* Disconnect the link, it will be reconnected automatically. */
|
||||
freeClusterLink(node->link);
|
||||
@@ -3554,7 +3572,13 @@ void clusterCron(void) {
|
||||
/* Compute the delay of the PONG. Note that if we already received
|
||||
* the PONG, then node->ping_sent is zero, so can't reach this
|
||||
* code at all. */
|
||||
delay = now - node->ping_sent;
|
||||
mstime_t delay = now - node->ping_sent;
|
||||
|
||||
/* We consider every incoming data as proof of liveness, since
|
||||
* our cluster bus link is also used for data: under heavy data
|
||||
* load pong delays are possible. */
|
||||
mstime_t data_delay = now - node->data_received;
|
||||
if (data_delay < delay) delay = data_delay;
|
||||
|
||||
if (delay > server.cluster_node_timeout) {
|
||||
/* Timeout reached. Set the node as possibly failing if it is
|
||||
@@ -4148,11 +4172,17 @@ void clusterReplyMultiBulkSlots(client *c) {
|
||||
while((de = dictNext(di)) != NULL) {
|
||||
clusterNode *node = dictGetVal(de);
|
||||
int j = 0, start = -1;
|
||||
int i, nested_elements = 0;
|
||||
|
||||
/* Skip slaves (that are iterated when producing the output of their
|
||||
* master) and masters not serving any slot. */
|
||||
if (!nodeIsMaster(node) || node->numslots == 0) continue;
|
||||
|
||||
for(i = 0; i < node->numslaves; i++) {
|
||||
if (nodeFailed(node->slaves[i])) continue;
|
||||
nested_elements++;
|
||||
}
|
||||
|
||||
for (j = 0; j < CLUSTER_SLOTS; j++) {
|
||||
int bit, i;
|
||||
|
||||
@@ -4160,8 +4190,7 @@ void clusterReplyMultiBulkSlots(client *c) {
|
||||
if (start == -1) start = j;
|
||||
}
|
||||
if (start != -1 && (!bit || j == CLUSTER_SLOTS-1)) {
|
||||
int nested_elements = 3; /* slots (2) + master addr (1). */
|
||||
void *nested_replylen = addDeferredMultiBulkLength(c);
|
||||
addReplyMultiBulkLen(c, nested_elements + 3); /* slots (2) + master addr (1). */
|
||||
|
||||
if (bit && j == CLUSTER_SLOTS-1) j++;
|
||||
|
||||
@@ -4191,9 +4220,7 @@ void clusterReplyMultiBulkSlots(client *c) {
|
||||
addReplyBulkCString(c, node->slaves[i]->ip);
|
||||
addReplyLongLong(c, node->slaves[i]->port);
|
||||
addReplyBulkCBuffer(c, node->slaves[i]->name, CLUSTER_NAMELEN);
|
||||
nested_elements++;
|
||||
}
|
||||
setDeferredMultiBulkLength(c, nested_replylen, nested_elements);
|
||||
num_masters++;
|
||||
}
|
||||
}
|
||||
@@ -4907,7 +4934,8 @@ void restoreCommand(client *c) {
|
||||
}
|
||||
|
||||
/* Make sure this key does not already exist here... */
|
||||
if (!replace && lookupKeyWrite(c->db,c->argv[1]) != NULL) {
|
||||
robj *key = c->argv[1];
|
||||
if (!replace && lookupKeyWrite(c->db,key) != NULL) {
|
||||
addReply(c,shared.busykeyerr);
|
||||
return;
|
||||
}
|
||||
@@ -4929,23 +4957,37 @@ void restoreCommand(client *c) {
|
||||
|
||||
rioInitWithBuffer(&payload,c->argv[3]->ptr);
|
||||
if (((type = rdbLoadObjectType(&payload)) == -1) ||
|
||||
((obj = rdbLoadObject(type,&payload,c->argv[1])) == NULL))
|
||||
((obj = rdbLoadObject(type,&payload,key)) == NULL))
|
||||
{
|
||||
addReplyError(c,"Bad data format");
|
||||
return;
|
||||
}
|
||||
|
||||
/* Remove the old key if needed. */
|
||||
if (replace) dbDelete(c->db,c->argv[1]);
|
||||
int deleted = 0;
|
||||
if (replace)
|
||||
deleted = dbDelete(c->db,key);
|
||||
|
||||
if (ttl && !absttl) ttl+=mstime();
|
||||
if (ttl && checkAlreadyExpired(ttl)) {
|
||||
if (deleted) {
|
||||
rewriteClientCommandVector(c,2,shared.del,key);
|
||||
signalModifiedKey(c->db,key);
|
||||
notifyKeyspaceEvent(NOTIFY_GENERIC,"del",key,c->db->id);
|
||||
server.dirty++;
|
||||
}
|
||||
decrRefCount(obj);
|
||||
addReply(c, shared.ok);
|
||||
return;
|
||||
}
|
||||
|
||||
/* Create the key and set the TTL if any */
|
||||
dbAdd(c->db,c->argv[1],obj);
|
||||
dbAdd(c->db,key,obj);
|
||||
if (ttl) {
|
||||
if (!absttl) ttl+=mstime();
|
||||
setExpire(c,c->db,c->argv[1],ttl);
|
||||
setExpire(c,c->db,key,ttl);
|
||||
}
|
||||
objectSetLRUOrLFU(obj,lfu_freq,lru_idle,lru_clock);
|
||||
signalModifiedKey(c->db,c->argv[1]);
|
||||
signalModifiedKey(c->db,key);
|
||||
addReply(c,shared.ok);
|
||||
server.dirty++;
|
||||
}
|
||||
@@ -5699,6 +5741,15 @@ int clusterRedirectBlockedClientIfNeeded(client *c) {
|
||||
int slot = keyHashSlot((char*)key->ptr, sdslen(key->ptr));
|
||||
clusterNode *node = server.cluster->slots[slot];
|
||||
|
||||
/* if the client is read-only and attempting to access key that our
|
||||
* replica can handle, allow it. */
|
||||
if ((c->flags & CLIENT_READONLY) &&
|
||||
(c->lastcmd->flags & CMD_READONLY) &&
|
||||
nodeIsSlave(myself) && myself->slaveof == node)
|
||||
{
|
||||
node = myself;
|
||||
}
|
||||
|
||||
/* We send an error and unblock the client if:
|
||||
* 1) The slot is unassigned, emitting a cluster down error.
|
||||
* 2) The slot is not handled by this node, nor being imported. */
|
||||
|
||||
@@ -128,6 +128,7 @@ typedef struct clusterNode {
|
||||
tables. */
|
||||
mstime_t ping_sent; /* Unix time we sent latest ping */
|
||||
mstime_t pong_received; /* Unix time we received the pong */
|
||||
mstime_t data_received; /* Unix time we received any data */
|
||||
mstime_t fail_time; /* Unix time when FAIL flag was set */
|
||||
mstime_t voted_time; /* Last time we voted for a slave of this master */
|
||||
mstime_t repl_offset_time; /* Unix time we received offset for this node */
|
||||
|
||||
+1
-1
@@ -355,7 +355,7 @@ long activeDefragSdsListAndDict(list *l, dict *d, int dict_val_type) {
|
||||
sdsele = ln->value;
|
||||
if ((newsds = activeDefragSds(sdsele))) {
|
||||
/* When defragging an sds value, we need to update the dict key */
|
||||
uint64_t hash = dictGetHash(d, sdsele);
|
||||
uint64_t hash = dictGetHash(d, newsds);
|
||||
replaceSateliteDictKeyPtrAndOrDefragDictEntry(d, sdsele, newsds, hash, &defragged);
|
||||
ln->value = newsds;
|
||||
defragged++;
|
||||
|
||||
+11
-7
@@ -391,6 +391,16 @@ void flushSlaveKeysWithExpireList(void) {
|
||||
}
|
||||
}
|
||||
|
||||
int checkAlreadyExpired(long long when) {
|
||||
/* EXPIRE with negative TTL, or EXPIREAT with a timestamp into the past
|
||||
* should never be executed as a DEL when load the AOF or in the context
|
||||
* of a slave instance.
|
||||
*
|
||||
* Instead we add the already expired key to the database with expire time
|
||||
* (possibly in the past) and wait for an explicit DEL from the master. */
|
||||
return (when <= mstime() && !server.loading && !server.masterhost);
|
||||
}
|
||||
|
||||
/*-----------------------------------------------------------------------------
|
||||
* Expires Commands
|
||||
*----------------------------------------------------------------------------*/
|
||||
@@ -418,13 +428,7 @@ void expireGenericCommand(client *c, long long basetime, int unit) {
|
||||
return;
|
||||
}
|
||||
|
||||
/* EXPIRE with negative TTL, or EXPIREAT with a timestamp into the past
|
||||
* should never be executed as a DEL when load the AOF or in the context
|
||||
* of a slave instance.
|
||||
*
|
||||
* Instead we take the other branch of the IF statement setting an expire
|
||||
* (possibly in the past) and wait for an explicit DEL from the master. */
|
||||
if (when <= mstime() && !server.loading && !server.masterhost) {
|
||||
if (checkAlreadyExpired(when)) {
|
||||
robj *aux;
|
||||
|
||||
int deleted = server.lazyfree_lazy_expire ? dbAsyncDelete(c->db,key) :
|
||||
|
||||
@@ -41,6 +41,30 @@ size_t lazyfreeGetFreeEffort(robj *obj) {
|
||||
} else if (obj->type == OBJ_HASH && obj->encoding == OBJ_ENCODING_HT) {
|
||||
dict *ht = obj->ptr;
|
||||
return dictSize(ht);
|
||||
} else if (obj->type == OBJ_STREAM) {
|
||||
size_t effort = 0;
|
||||
stream *s = obj->ptr;
|
||||
|
||||
/* Make a best effort estimate to maintain constant runtime. Every macro
|
||||
* node in the Stream is one allocation. */
|
||||
effort += s->rax->numnodes;
|
||||
|
||||
/* Every consumer group is an allocation and so are the entries in its
|
||||
* PEL. We use size of the first group's PEL as an estimate for all
|
||||
* others. */
|
||||
if (s->cgroups) {
|
||||
raxIterator ri;
|
||||
streamCG *cg;
|
||||
raxStart(&ri,s->cgroups);
|
||||
raxSeek(&ri,"^",NULL,0);
|
||||
/* There must be at least one group so the following should always
|
||||
* work. */
|
||||
serverAssert(raxNext(&ri));
|
||||
cg = ri.data;
|
||||
effort += raxSize(s->cgroups)*(1+raxSize(cg->pel));
|
||||
raxStop(&ri);
|
||||
}
|
||||
return effort;
|
||||
} else {
|
||||
return 1; /* Everything else is a single allocation. */
|
||||
}
|
||||
|
||||
+6
-2
@@ -471,7 +471,8 @@ int moduleDelKeyIfEmpty(RedisModuleKey *key) {
|
||||
case OBJ_LIST: isempty = listTypeLength(o) == 0; break;
|
||||
case OBJ_SET: isempty = setTypeSize(o) == 0; break;
|
||||
case OBJ_ZSET: isempty = zsetLength(o) == 0; break;
|
||||
case OBJ_HASH : isempty = hashTypeLength(o) == 0; break;
|
||||
case OBJ_HASH: isempty = hashTypeLength(o) == 0; break;
|
||||
case OBJ_STREAM: isempty = streamLength(o) == 0; break;
|
||||
default: isempty = 0;
|
||||
}
|
||||
|
||||
@@ -542,6 +543,8 @@ void moduleHandlePropagationAfterCommandCallback(RedisModuleCtx *ctx) {
|
||||
redisOpArrayFree(&server.also_propagate);
|
||||
/* Restore the previous oparray in case of nexted use of the API. */
|
||||
server.also_propagate = ctx->saved_oparray;
|
||||
/* We're done with saved_oparray, let's invalidate it. */
|
||||
redisOpArrayInit(&ctx->saved_oparray);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1684,6 +1687,7 @@ int RM_KeyType(RedisModuleKey *key) {
|
||||
case OBJ_ZSET: return REDISMODULE_KEYTYPE_ZSET;
|
||||
case OBJ_HASH: return REDISMODULE_KEYTYPE_HASH;
|
||||
case OBJ_MODULE: return REDISMODULE_KEYTYPE_MODULE;
|
||||
/* case OBJ_STREAM: return REDISMODULE_KEYTYPE_STREAM; - don't wanna add new API to 5.0 */
|
||||
default: return 0;
|
||||
}
|
||||
}
|
||||
@@ -1701,6 +1705,7 @@ size_t RM_ValueLength(RedisModuleKey *key) {
|
||||
case OBJ_SET: return setTypeSize(key->value);
|
||||
case OBJ_ZSET: return zsetLength(key->value);
|
||||
case OBJ_HASH: return hashTypeLength(key->value);
|
||||
case OBJ_STREAM: return streamLength(key->value);
|
||||
default: return 0;
|
||||
}
|
||||
}
|
||||
@@ -3893,7 +3898,6 @@ void moduleHandleBlockedClients(void) {
|
||||
ctx.client = bc->client;
|
||||
ctx.blocked_client = bc;
|
||||
bc->reply_callback(&ctx,(void**)c->argv,c->argc);
|
||||
moduleHandlePropagationAfterCommandCallback(&ctx);
|
||||
moduleFreeContext(&ctx);
|
||||
}
|
||||
|
||||
|
||||
@@ -1097,6 +1097,8 @@ ssize_t rdbSaveSingleModuleAux(rio *rdb, int when, moduleType *mt) {
|
||||
/* Save a module-specific aux value. */
|
||||
RedisModuleIO io;
|
||||
int retval = rdbSaveType(rdb, RDB_OPCODE_MODULE_AUX);
|
||||
if (retval == -1) return -1;
|
||||
io.bytes += retval;
|
||||
|
||||
/* Write the "module" identifier as prefix, so that we'll be able
|
||||
* to call the right module during loading. */
|
||||
@@ -1768,8 +1770,8 @@ robj *rdbLoadObject(int rdbtype, rio *rdb, robj *key) {
|
||||
rdbExitReportCorruptRDB(
|
||||
"Error reading the consumer name from Stream group");
|
||||
}
|
||||
streamConsumer *consumer = streamLookupConsumer(cgroup,cname,
|
||||
1);
|
||||
streamConsumer *consumer =
|
||||
streamLookupConsumer(cgroup,cname,SLC_NONE);
|
||||
sdsfree(cname);
|
||||
consumer->seen_time = rdbLoadMillisecondTime(rdb,RDB_VERSION);
|
||||
|
||||
|
||||
@@ -146,6 +146,7 @@ robj *rdbLoadObject(int type, rio *rdb, robj *key);
|
||||
void backgroundSaveDoneHandler(int exitcode, int bysignal);
|
||||
int rdbSaveKeyValuePair(rio *rdb, robj *key, robj *val, long long expiretime);
|
||||
ssize_t rdbSaveSingleModuleAux(rio *rdb, int when, moduleType *mt);
|
||||
robj *rdbLoadCheckModuleValue(rio *rdb, char *modulename);
|
||||
robj *rdbLoadStringObject(rio *rdb);
|
||||
ssize_t rdbSaveStringObject(rio *rdb, robj *obj);
|
||||
ssize_t rdbSaveRawString(rio *rdb, unsigned char *s, size_t len);
|
||||
|
||||
+18
-1
@@ -58,6 +58,7 @@ struct {
|
||||
#define RDB_CHECK_DOING_CHECK_SUM 5
|
||||
#define RDB_CHECK_DOING_READ_LEN 6
|
||||
#define RDB_CHECK_DOING_READ_AUX 7
|
||||
#define RDB_CHECK_DOING_READ_MODULE_AUX 8
|
||||
|
||||
char *rdb_check_doing_string[] = {
|
||||
"start",
|
||||
@@ -67,7 +68,8 @@ char *rdb_check_doing_string[] = {
|
||||
"read-object-value",
|
||||
"check-sum",
|
||||
"read-len",
|
||||
"read-aux"
|
||||
"read-aux",
|
||||
"read-module-aux"
|
||||
};
|
||||
|
||||
char *rdb_type_string[] = {
|
||||
@@ -270,6 +272,21 @@ int redis_check_rdb(char *rdbfilename, FILE *fp) {
|
||||
decrRefCount(auxkey);
|
||||
decrRefCount(auxval);
|
||||
continue; /* Read type again. */
|
||||
} else if (type == RDB_OPCODE_MODULE_AUX) {
|
||||
/* AUX: Auxiliary data for modules. */
|
||||
uint64_t moduleid, when_opcode, when;
|
||||
rdbstate.doing = RDB_CHECK_DOING_READ_MODULE_AUX;
|
||||
if ((moduleid = rdbLoadLen(&rdb,NULL)) == RDB_LENERR) goto eoferr;
|
||||
if ((when_opcode = rdbLoadLen(&rdb,NULL)) == RDB_LENERR) goto eoferr;
|
||||
if ((when = rdbLoadLen(&rdb,NULL)) == RDB_LENERR) goto eoferr;
|
||||
|
||||
char name[10];
|
||||
moduleTypeNameByID(name,moduleid);
|
||||
rdbCheckInfo("MODULE AUX for: %s", name);
|
||||
|
||||
robj *o = rdbLoadCheckModuleValue(&rdb,name);
|
||||
decrRefCount(o);
|
||||
continue; /* Read type again. */
|
||||
} else {
|
||||
if (!rdbIsObjectType(type)) {
|
||||
rdbCheckError("Invalid object type: %d", type);
|
||||
|
||||
+11
-6
@@ -1144,7 +1144,7 @@ static int cliSendCommand(int argc, char **argv, long repeat) {
|
||||
for (j = 0; j < argc; j++)
|
||||
argvlen[j] = sdslen(argv[j]);
|
||||
|
||||
while(repeat-- > 0) {
|
||||
while(repeat < 0 || repeat-- > 0) {
|
||||
redisAppendCommandArgv(context,argc,(const char**)argv,argvlen);
|
||||
while (config.monitor_mode) {
|
||||
if (cliReadReply(output_raw) != REDIS_OK) exit(1);
|
||||
@@ -1179,6 +1179,11 @@ static int cliSendCommand(int argc, char **argv, long repeat) {
|
||||
cliSelect();
|
||||
}
|
||||
}
|
||||
if (config.cluster_reissue_command){
|
||||
/* If we need to reissue the command, break to prevent a
|
||||
further 'repeat' number of dud interations */
|
||||
break;
|
||||
}
|
||||
if (config.interval) usleep(config.interval);
|
||||
fflush(stdout); /* Make it grep friendly */
|
||||
}
|
||||
@@ -1589,12 +1594,12 @@ static int issueCommandRepeat(int argc, char **argv, long repeat) {
|
||||
cliPrintContextError();
|
||||
return REDIS_ERR;
|
||||
}
|
||||
}
|
||||
/* Issue the command again if we got redirected in cluster mode */
|
||||
if (config.cluster_mode && config.cluster_reissue_command) {
|
||||
}
|
||||
/* Issue the command again if we got redirected in cluster mode */
|
||||
if (config.cluster_mode && config.cluster_reissue_command) {
|
||||
cliConnect(CC_FORCE);
|
||||
} else {
|
||||
break;
|
||||
} else {
|
||||
break;
|
||||
}
|
||||
}
|
||||
return REDIS_OK;
|
||||
|
||||
@@ -1292,6 +1292,17 @@ void readSyncBulkPayload(aeEventLoop *el, int fd, void *privdata, int mask) {
|
||||
rdbRemoveTempFile(server.rdb_child_pid);
|
||||
}
|
||||
|
||||
/* Make sure the new file (also used for persistence) is fully synced
|
||||
* (not covered by earlier calls to rdb_fsync_range). */
|
||||
if (fsync(server.repl_transfer_fd) == -1) {
|
||||
serverLog(LL_WARNING,
|
||||
"Failed trying to sync the temp DB to disk in "
|
||||
"MASTER <-> REPLICA synchronization: %s",
|
||||
strerror(errno));
|
||||
cancelReplicationHandshake();
|
||||
return;
|
||||
}
|
||||
|
||||
if (rename(server.repl_transfer_tmpfile,server.rdb_filename) == -1) {
|
||||
serverLog(LL_WARNING,"Failed trying to rename the temp DB into dump.rdb in MASTER <-> REPLICA synchronization: %s", strerror(errno));
|
||||
cancelReplicationHandshake();
|
||||
|
||||
+3
-4
@@ -521,12 +521,11 @@ int luaRedisGenericCommand(lua_State *lua, int raise_error) {
|
||||
!server.loading && /* Don't care about mem if loading. */
|
||||
!server.masterhost && /* Slave must execute the script. */
|
||||
server.lua_write_dirty == 0 && /* Script had no side effects so far. */
|
||||
server.lua_oom && /* Detected OOM when script start. */
|
||||
(cmd->flags & CMD_DENYOOM))
|
||||
{
|
||||
if (getMaxmemoryState(NULL,NULL,NULL,NULL) != C_OK) {
|
||||
luaPushError(lua, shared.oomerr->ptr);
|
||||
goto cleanup;
|
||||
}
|
||||
luaPushError(lua, shared.oomerr->ptr);
|
||||
goto cleanup;
|
||||
}
|
||||
|
||||
if (cmd->flags & CMD_RANDOM) server.lua_random_dirty = 1;
|
||||
|
||||
@@ -1061,6 +1061,7 @@ int sentinelTryConnectionSharing(sentinelRedisInstance *ri) {
|
||||
releaseInstanceLink(ri->link,NULL);
|
||||
ri->link = match->link;
|
||||
match->link->refcount++;
|
||||
dictReleaseIterator(di);
|
||||
return C_OK;
|
||||
}
|
||||
dictReleaseIterator(di);
|
||||
|
||||
+8
-1
@@ -236,7 +236,7 @@ struct redisCommand redisCommandTable[] = {
|
||||
{"keys",keysCommand,2,"rS",0,NULL,0,0,0,0,0},
|
||||
{"scan",scanCommand,-2,"rR",0,NULL,0,0,0,0,0},
|
||||
{"dbsize",dbsizeCommand,1,"rF",0,NULL,0,0,0,0,0},
|
||||
{"auth",authCommand,2,"sltF",0,NULL,0,0,0,0,0},
|
||||
{"auth",authCommand,2,"sltMF",0,NULL,0,0,0,0,0},
|
||||
{"ping",pingCommand,-1,"tF",0,NULL,0,0,0,0,0},
|
||||
{"echo",echoCommand,2,"F",0,NULL,0,0,0,0,0},
|
||||
{"save",saveCommand,1,"as",0,NULL,0,0,0,0,0},
|
||||
@@ -2667,6 +2667,13 @@ int processCommand(client *c) {
|
||||
addReply(c, shared.oomerr);
|
||||
return C_OK;
|
||||
}
|
||||
|
||||
/* Save out_of_memory result at script start, otherwise if we check OOM
|
||||
* untill first write within script, memory used by lua stack and
|
||||
* arguments might interfere. */
|
||||
if (c->cmd->proc == evalCommand || c->cmd->proc == evalShaCommand) {
|
||||
server.lua_oom = out_of_memory;
|
||||
}
|
||||
}
|
||||
|
||||
/* Don't accept write commands if there are problems persisting on disk
|
||||
|
||||
@@ -1277,6 +1277,7 @@ struct redisServer {
|
||||
execution. */
|
||||
int lua_kill; /* Kill the script if true. */
|
||||
int lua_always_replicate_commands; /* Default replication type. */
|
||||
int lua_oom; /* OOM detected when script start? */
|
||||
/* Lazy free */
|
||||
int lazyfree_lazy_eviction;
|
||||
int lazyfree_lazy_expire;
|
||||
@@ -1833,6 +1834,7 @@ void propagateExpire(redisDb *db, robj *key, int lazy);
|
||||
int expireIfNeeded(redisDb *db, robj *key);
|
||||
long long getExpire(redisDb *db, robj *key);
|
||||
void setExpire(client *c, redisDb *db, robj *key, long long when);
|
||||
int checkAlreadyExpired(long long when);
|
||||
robj *lookupKey(redisDb *db, robj *key, int flags);
|
||||
robj *lookupKeyRead(redisDb *db, robj *key);
|
||||
robj *lookupKeyWrite(redisDb *db, robj *key);
|
||||
|
||||
+7
-1
@@ -96,15 +96,21 @@ typedef struct sreamPropInfo {
|
||||
/* Prototypes of exported APIs. */
|
||||
struct client;
|
||||
|
||||
/* Flags for streamLookupConsumer */
|
||||
#define SLC_NONE 0
|
||||
#define SLC_NOCREAT (1<<0) /* Do not create the consumer if it doesn't exist */
|
||||
#define SLC_NOREFRESH (1<<1) /* Do not update consumer's seen-time */
|
||||
|
||||
stream *streamNew(void);
|
||||
void freeStream(stream *s);
|
||||
unsigned long streamLength(const robj *subject);
|
||||
size_t streamReplyWithRange(client *c, stream *s, streamID *start, streamID *end, size_t count, int rev, streamCG *group, streamConsumer *consumer, int flags, streamPropInfo *spi);
|
||||
void streamIteratorStart(streamIterator *si, stream *s, streamID *start, streamID *end, int rev);
|
||||
int streamIteratorGetID(streamIterator *si, streamID *id, int64_t *numfields);
|
||||
void streamIteratorGetField(streamIterator *si, unsigned char **fieldptr, unsigned char **valueptr, int64_t *fieldlen, int64_t *valuelen);
|
||||
void streamIteratorStop(streamIterator *si);
|
||||
streamCG *streamLookupCG(stream *s, sds groupname);
|
||||
streamConsumer *streamLookupConsumer(streamCG *cg, sds name, int create);
|
||||
streamConsumer *streamLookupConsumer(streamCG *cg, sds name, int flags);
|
||||
streamCG *streamCreateCG(stream *s, char *name, size_t namelen, streamID *id);
|
||||
streamNACK *streamCreateNACK(streamConsumer *consumer);
|
||||
void streamDecodeID(void *buf, streamID *id);
|
||||
|
||||
+37
-14
@@ -82,6 +82,12 @@ void streamIncrID(streamID *id) {
|
||||
}
|
||||
}
|
||||
|
||||
/* Return the length of a stream. */
|
||||
unsigned long streamLength(const robj *subject) {
|
||||
stream *s = subject->ptr;
|
||||
return s->length;
|
||||
}
|
||||
|
||||
/* Generate the next stream item ID given the previous one. If the current
|
||||
* milliseconds Unix time is greater than the previous one, just use this
|
||||
* as time part and start with sequence part of zero. Otherwise we use the
|
||||
@@ -842,6 +848,11 @@ void streamPropagateXCLAIM(client *c, robj *key, streamCG *group, robj *groupnam
|
||||
argv[11] = createStringObject("JUSTID",6);
|
||||
argv[12] = createStringObject("LASTID",6);
|
||||
argv[13] = createObjectFromStreamID(&group->last_id);
|
||||
|
||||
/* We use progagate() because this code path is not always called from
|
||||
* the command execution context. Moreover this will just alter the
|
||||
* consumer group state, and we don't need MULTI/EXEC wrapping because
|
||||
* there is no message state cross-message atomicity required. */
|
||||
propagate(server.xclaimCommand,c->db->id,argv,14,PROPAGATE_AOF|PROPAGATE_REPL);
|
||||
decrRefCount(argv[0]);
|
||||
decrRefCount(argv[3]);
|
||||
@@ -869,7 +880,12 @@ void streamPropagateGroupID(client *c, robj *key, streamCG *group, robj *groupna
|
||||
argv[2] = key;
|
||||
argv[3] = groupname;
|
||||
argv[4] = createObjectFromStreamID(&group->last_id);
|
||||
alsoPropagate(server.xgroupCommand,c->db->id,argv,5,PROPAGATE_AOF|PROPAGATE_REPL);
|
||||
|
||||
/* We use progagate() because this code path is not always called from
|
||||
* the command execution context. Moreover this will just alter the
|
||||
* consumer group state, and we don't need MULTI/EXEC wrapping because
|
||||
* there is no message state cross-message atomicity required. */
|
||||
propagate(server.xgroupCommand,c->db->id,argv,5,PROPAGATE_AOF|PROPAGATE_REPL);
|
||||
decrRefCount(argv[0]);
|
||||
decrRefCount(argv[1]);
|
||||
decrRefCount(argv[4]);
|
||||
@@ -1564,7 +1580,8 @@ void xreadCommand(client *c) {
|
||||
addReplyBulk(c,c->argv[streams_arg+i]);
|
||||
streamConsumer *consumer = NULL;
|
||||
if (groups) consumer = streamLookupConsumer(groups[i],
|
||||
consumername->ptr,1);
|
||||
consumername->ptr,
|
||||
SLC_NONE);
|
||||
streamPropInfo spi = {c->argv[i+streams_arg],groupname};
|
||||
int flags = 0;
|
||||
if (noack) flags |= STREAM_RWR_NOACK;
|
||||
@@ -1697,7 +1714,9 @@ streamCG *streamLookupCG(stream *s, sds groupname) {
|
||||
* consumer does not exist it is automatically created as a side effect
|
||||
* of calling this function, otherwise its last seen time is updated and
|
||||
* the existing consumer reference returned. */
|
||||
streamConsumer *streamLookupConsumer(streamCG *cg, sds name, int create) {
|
||||
streamConsumer *streamLookupConsumer(streamCG *cg, sds name, int flags) {
|
||||
int create = !(flags & SLC_NOCREAT);
|
||||
int refresh = !(flags & SLC_NOREFRESH);
|
||||
streamConsumer *consumer = raxFind(cg->consumers,(unsigned char*)name,
|
||||
sdslen(name));
|
||||
if (consumer == raxNotFound) {
|
||||
@@ -1708,7 +1727,7 @@ streamConsumer *streamLookupConsumer(streamCG *cg, sds name, int create) {
|
||||
raxInsert(cg->consumers,(unsigned char*)name,sdslen(name),
|
||||
consumer,NULL);
|
||||
}
|
||||
consumer->seen_time = mstime();
|
||||
if (refresh) consumer->seen_time = mstime();
|
||||
return consumer;
|
||||
}
|
||||
|
||||
@@ -1716,7 +1735,8 @@ streamConsumer *streamLookupConsumer(streamCG *cg, sds name, int create) {
|
||||
* may have pending messages: they are removed from the PEL, and the number
|
||||
* of pending messages "lost" is returned. */
|
||||
uint64_t streamDelConsumer(streamCG *cg, sds name) {
|
||||
streamConsumer *consumer = streamLookupConsumer(cg,name,0);
|
||||
streamConsumer *consumer =
|
||||
streamLookupConsumer(cg,name,SLC_NOCREAT|SLC_NOREFRESH);
|
||||
if (consumer == NULL) return 0;
|
||||
|
||||
uint64_t retval = raxSize(consumer->pel);
|
||||
@@ -2047,15 +2067,18 @@ void xpendingCommand(client *c) {
|
||||
}
|
||||
/* XPENDING <key> <group> <start> <stop> <count> [<consumer>] variant. */
|
||||
else {
|
||||
streamConsumer *consumer = consumername ?
|
||||
streamLookupConsumer(group,consumername->ptr,0):
|
||||
NULL;
|
||||
streamConsumer *consumer = NULL;
|
||||
if (consumername) {
|
||||
consumer = streamLookupConsumer(group,
|
||||
consumername->ptr,
|
||||
SLC_NOCREAT|SLC_NOREFRESH);
|
||||
|
||||
/* If a consumer name was mentioned but it does not exist, we can
|
||||
* just return an empty array. */
|
||||
if (consumername && consumer == NULL) {
|
||||
addReplyMultiBulkLen(c,0);
|
||||
return;
|
||||
/* If a consumer name was mentioned but it does not exist, we can
|
||||
* just return an empty array. */
|
||||
if (consumer == NULL) {
|
||||
addReplyMultiBulkLen(c,0);
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
||||
rax *pel = consumer ? consumer->pel : group->pel;
|
||||
@@ -2317,7 +2340,7 @@ void xclaimCommand(client *c) {
|
||||
raxRemove(nack->consumer->pel,buf,sizeof(buf),NULL);
|
||||
/* Update the consumer and idle time. */
|
||||
if (consumer == NULL)
|
||||
consumer = streamLookupConsumer(group,c->argv[3]->ptr,1);
|
||||
consumer = streamLookupConsumer(group,c->argv[3]->ptr,SLC_NONE);
|
||||
nack->consumer = consumer;
|
||||
nack->delivery_time = deliverytime;
|
||||
/* Set the delivery attempts counter if given, otherwise
|
||||
|
||||
+1
-1
@@ -1 +1 @@
|
||||
#define REDIS_VERSION "5.0.9"
|
||||
#define REDIS_VERSION "5.0.10"
|
||||
|
||||
@@ -177,9 +177,6 @@ void *zrealloc(void *ptr, size_t size) {
|
||||
size_t zmalloc_size(void *ptr) {
|
||||
void *realptr = (char*)ptr-PREFIX_SIZE;
|
||||
size_t size = *((size_t*)realptr);
|
||||
/* Assume at least that all the allocations are padded at sizeof(long) by
|
||||
* the underlying allocator. */
|
||||
if (size&(sizeof(long)-1)) size += sizeof(long)-(size&(sizeof(long)-1));
|
||||
return size+PREFIX_SIZE;
|
||||
}
|
||||
size_t zmalloc_usable(void *ptr) {
|
||||
|
||||
@@ -0,0 +1,70 @@
|
||||
# Check basic transactions on a replica.
|
||||
|
||||
source "../tests/includes/init-tests.tcl"
|
||||
|
||||
test "Create a primary with a replica" {
|
||||
create_cluster 1 1
|
||||
}
|
||||
|
||||
test "Cluster should start ok" {
|
||||
assert_cluster_state ok
|
||||
}
|
||||
|
||||
set primary [Rn 0]
|
||||
set replica [Rn 1]
|
||||
|
||||
test "Cant read from replica without READONLY" {
|
||||
$primary SET a 1
|
||||
catch {$replica GET a} err
|
||||
assert {[string range $err 0 4] eq {MOVED}}
|
||||
}
|
||||
|
||||
test "Can read from replica after READONLY" {
|
||||
$replica READONLY
|
||||
assert {[$replica GET a] eq {1}}
|
||||
}
|
||||
|
||||
test "Can preform HSET primary and HGET from replica" {
|
||||
$primary HSET h a 1
|
||||
$primary HSET h b 2
|
||||
$primary HSET h c 3
|
||||
assert {[$replica HGET h a] eq {1}}
|
||||
assert {[$replica HGET h b] eq {2}}
|
||||
assert {[$replica HGET h c] eq {3}}
|
||||
}
|
||||
|
||||
# didn't cherry pick b120366d4 to 5.0 yet
|
||||
#test "Can MULTI-EXEC transaction of HGET operations from replica" {
|
||||
# $replica MULTI
|
||||
# assert {[$replica HGET h a] eq {QUEUED}}
|
||||
# assert {[$replica HGET h b] eq {QUEUED}}
|
||||
# assert {[$replica HGET h c] eq {QUEUED}}
|
||||
# assert {[$replica EXEC] eq {1 2 3}}
|
||||
#}
|
||||
|
||||
test "MULTI-EXEC with write operations is MOVED" {
|
||||
$replica MULTI
|
||||
catch {$replica HSET h b 4} err
|
||||
assert {[string range $err 0 4] eq {MOVED}}
|
||||
catch {$replica exec} err
|
||||
assert {[string range $err 0 8] eq {EXECABORT}}
|
||||
}
|
||||
|
||||
test "read-only blocking operations from replica" {
|
||||
set rd [redis_deferring_client redis 1]
|
||||
$rd readonly
|
||||
$rd read
|
||||
$rd XREAD BLOCK 0 STREAMS k 0
|
||||
|
||||
wait_for_condition 1000 50 {
|
||||
[RI 1 blocked_clients] eq {1}
|
||||
} else {
|
||||
fail "client wasn't blocked"
|
||||
}
|
||||
|
||||
$primary XADD k * foo bar
|
||||
set res [$rd read]
|
||||
set res [lindex [lindex [lindex [lindex $res 0] 1] 0] 1]
|
||||
assert {$res eq {foo bar}}
|
||||
$rd close
|
||||
}
|
||||
+21
-2
@@ -334,10 +334,16 @@ proc S {n args} {
|
||||
[dict get $s link] {*}$args
|
||||
}
|
||||
|
||||
# Returns a Redis instance by index.
|
||||
# Example:
|
||||
# [Rn 0] info
|
||||
proc Rn {n} {
|
||||
return [dict get [lindex $::redis_instances $n] link]
|
||||
}
|
||||
|
||||
# Like R but to chat with Redis instances.
|
||||
proc R {n args} {
|
||||
set r [lindex $::redis_instances $n]
|
||||
[dict get $r link] {*}$args
|
||||
[Rn $n] {*}$args
|
||||
}
|
||||
|
||||
proc get_info_field {info field} {
|
||||
@@ -509,3 +515,16 @@ proc restart_instance {type id} {
|
||||
}
|
||||
}
|
||||
|
||||
proc redis_deferring_client {type id} {
|
||||
set port [get_instance_attrib $type $id port]
|
||||
set host [get_instance_attrib $type $id host]
|
||||
set client [redis $host $port 1]
|
||||
return $client
|
||||
}
|
||||
|
||||
proc redis_client {type id} {
|
||||
set port [get_instance_attrib $type $id port]
|
||||
set host [get_instance_attrib $type $id host]
|
||||
set client [redis $host $port 0]
|
||||
return $client
|
||||
}
|
||||
|
||||
+12
-1
@@ -36,7 +36,18 @@ start_server {tags {"dump"}} {
|
||||
assert {$ttl >= 2900 && $ttl <= 3100}
|
||||
r get foo
|
||||
} {bar}
|
||||
|
||||
|
||||
test {RESTORE with ABSTTL in the past} {
|
||||
r set foo bar
|
||||
set encoded [r dump foo]
|
||||
set now [clock milliseconds]
|
||||
r debug set-active-expire 0
|
||||
r restore foo [expr $now-3000] $encoded absttl REPLACE
|
||||
catch {r debug object foo} e
|
||||
r debug set-active-expire 1
|
||||
set e
|
||||
} {ERR no such key}
|
||||
|
||||
test {RESTORE can set LRU} {
|
||||
r set foo bar
|
||||
set encoded [r dump foo]
|
||||
|
||||
Reference in New Issue
Block a user