Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1fab07e078 | ||
|
|
8ebae5d630 | ||
|
|
21c3d77118 | ||
|
|
60a28fad8a | ||
|
|
e42baed4c3 | ||
|
|
aa67aec84e | ||
|
|
93959bc09f | ||
|
|
2e92d0f04a | ||
|
|
adcb470130 | ||
|
|
1b71fea998 | ||
|
|
2b5cf6bf78 | ||
|
|
7e78ab4b6f | ||
|
|
3468cd3664 | ||
|
|
d1b5c5defd | ||
|
|
66899a42fc | ||
|
|
76b18c7a0e | ||
|
|
ca804a1022 | ||
|
|
d15d9fecd2 | ||
|
|
c2717911db | ||
|
|
1641f41cfc | ||
|
|
b37b2b5c14 | ||
|
|
c43c970344 | ||
|
|
b64c861171 | ||
|
|
47bbaa17b0 | ||
|
|
0595420b1e | ||
|
|
2d34ec60bf | ||
|
|
2d7d75adb3 |
+84
-13
@@ -12,6 +12,88 @@ 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 3.0.0 ] Release date: 1 Apr 2015
|
||||
|
||||
>> What's new in Redis 3.0 compared to Redis 2.8?
|
||||
|
||||
* Redis Cluster: a distributed implementation of a subset of Redis.
|
||||
* New "embedded string" object encoding resulting in less cache
|
||||
misses. Big speed gain under certain work loads.
|
||||
* AOF child -> parent final data transmission to minimize latency due
|
||||
to "last write" during AOF rewrites.
|
||||
* Much improved LRU approximation algorithm for keys eviction.
|
||||
* WAIT command to block waiting for a write to be transmitted to
|
||||
the specified number of slaves.
|
||||
* MIGRATE connection caching. Much faster keys migraitons.
|
||||
* MIGARTE new options COPY and REPLACE.
|
||||
* CLIENT PAUSE command: stop processing client requests for a
|
||||
specified amount of time.
|
||||
* BITCOUNT performance improvements.
|
||||
* CONFIG SET accepts memory values in different units (for example
|
||||
you can use "CONFIG SET maxmemory 1gb").
|
||||
* Redis log format slightly changed reporting in each line the role of the
|
||||
instance (master/slave) or if it's a saving child log.
|
||||
* INCR performance improvements.
|
||||
|
||||
>> Refactoring changes (no new features nor bug fixes)
|
||||
|
||||
* Blocking operations full refactoring (blocked.c)
|
||||
* Client output buffer memory tracking refactored.
|
||||
|
||||
Changes between RC6 and 3.0.0 stable:
|
||||
|
||||
>> General changes
|
||||
|
||||
* Fixes to diskless replication. (Oran Agra)
|
||||
* Test for BLPOP replication on role change. (Salvatore Sanfilippo)
|
||||
* prepareClientToWrite() error handling improvements. (Salvatore Sanfilippo)
|
||||
* Remove dict.c no longer used function. (Salvatore Sanfilippo)
|
||||
|
||||
>> Cluster changes
|
||||
|
||||
None
|
||||
|
||||
>> Sentinel changes
|
||||
|
||||
None
|
||||
|
||||
--[ Redis 3.0.0 RC6 (version 2.9.106) ] Release date: 24 mar 2015
|
||||
|
||||
Upgrade urgency: HIGH because of bugs related to Redis Custer and replication.
|
||||
|
||||
This is the 6th release candidate of Redis 3.0.0. This release fixes important
|
||||
issues discovered during stress testing, and implements safest behavior
|
||||
for blocking operations during clients reshardings, and a new much needed
|
||||
functionality of Redis Cluster manual failovers.
|
||||
|
||||
In order to fix certain bugs quite a bit of refactoring was needed which
|
||||
is usually non advisabble in a Release Candidate, but needed in order to
|
||||
end with a clean fix.
|
||||
|
||||
>> General changes
|
||||
|
||||
* [FIX] Redis (non clustered & clustered) replication bug involving blocking
|
||||
operations: see issue #2473. (Salvatore Sanfilippo)
|
||||
|
||||
>> Cluster changes
|
||||
|
||||
* [FIX] clientsArePaused() fix crashing the old master during manual failover.
|
||||
(Salvatore Sanfilippo)
|
||||
* [FIX] Lua scripts replication in Redis Cluster was totally broken.
|
||||
(Salvatore Sanfilippo)
|
||||
* [FIX] Redirect clients blocked into list operations when the hash slot
|
||||
they are blocked into is migrated to another instance or the cluster
|
||||
state turns into "fail". (Salvatore Sanfilippo)
|
||||
|
||||
* [NEW] TAKEOVER option for CLUSTER FAILOVER implemented. It is now possible
|
||||
to fix a cluster manually in the minority side of the partition, for
|
||||
example in order to allow for multi DC setups & recovery.
|
||||
(Salvatore Sanfilippo)
|
||||
|
||||
>> Sentinel changes
|
||||
|
||||
No changes in Sentinel.
|
||||
|
||||
--[ Redis 3.0.0 RC5 (version 2.9.105) ] Release date: 20 mar 2015
|
||||
|
||||
Upgrade urgency: Moderate for Redis Cluster users, low otherwise.
|
||||
@@ -29,6 +111,7 @@ process of finishing the documentation for Redis Cluster).
|
||||
* [FIX] Fix for backtrace generation issue. (Mariano Pérez Rodríguez, Matt Stancliff, Salvatore Sanfilippo)
|
||||
|
||||
* [NEW] Redis-cli --latency-dist backported from unstable.
|
||||
(Salvatore Sanfilippo)
|
||||
|
||||
>> Cluster changes
|
||||
|
||||
@@ -492,22 +575,10 @@ This is the second beta of Redis 3.0.0.
|
||||
|
||||
This is the first beta of Redis 3.0.0.
|
||||
|
||||
The following is a list of improvements in Redis 3.0, compared to Redis 2.8.
|
||||
|
||||
* [NEW] Redis Cluster: a distributed implementation of a subset of Redis.
|
||||
* [NEW] New "embedded string" object encoding resulting in less cache
|
||||
misses. Big speed gain under certain work loads.
|
||||
* [NEW] WAIT command to block waiting for a write to be transmitted to
|
||||
the specified number of slaves.
|
||||
* [NEW] MIGRATE connection caching. Much faster keys migraitons.
|
||||
* [NEW] MIGARTE new options COPY and REPLACE.
|
||||
* [NEW] CLIENT PAUSE command: stop processing client requests for a
|
||||
specified amount of time.
|
||||
|
||||
Migrating from 2.8 to 3.0
|
||||
=========================
|
||||
|
||||
Redis 3.0 is mostly a strict subset of 2.8, you should not have any problem
|
||||
Redis 2.8 is mostly a strict subset of 3.0, you should not have any problem
|
||||
upgrading your application from 2.8 to 3.0. However this is a list of small
|
||||
non-backward compatible changes introduced in the 3.0 release:
|
||||
|
||||
|
||||
+26
-1
@@ -59,6 +59,8 @@
|
||||
* When implementing a new type of blocking opeation, the implementation
|
||||
* should modify unblockClient() and replyToBlockedClientTimedOut() in order
|
||||
* to handle the btype-specific behavior of this two functions.
|
||||
* If the blocking operation waits for certain keys to change state, the
|
||||
* clusterRedirectBlockedClientIfNeeded() function should also be updated.
|
||||
*/
|
||||
|
||||
#include "redis.h"
|
||||
@@ -114,7 +116,6 @@ void processUnblockedClients(void) {
|
||||
c = ln->value;
|
||||
listDelNode(server.unblocked_clients,ln);
|
||||
c->flags &= ~REDIS_UNBLOCKED;
|
||||
c->btype = REDIS_BLOCKED_NONE;
|
||||
|
||||
/* Process remaining data in the input buffer. */
|
||||
if (c->querybuf && sdslen(c->querybuf) > 0) {
|
||||
@@ -156,3 +157,27 @@ void replyToBlockedClientTimedOut(redisClient *c) {
|
||||
}
|
||||
}
|
||||
|
||||
/* Mass-unblock clients because something changed in the instance that makes
|
||||
* blocking no longer safe. For example clients blocked in list operations
|
||||
* in an instance which turns from master to slave is unsafe, so this function
|
||||
* is called when a master turns into a slave.
|
||||
*
|
||||
* The semantics is to send an -UNBLOCKED error to the client, disconnecting
|
||||
* it at the same time. */
|
||||
void disconnectAllBlockedClients(void) {
|
||||
listNode *ln;
|
||||
listIter li;
|
||||
|
||||
listRewind(server.clients,&li);
|
||||
while((ln = listNext(&li))) {
|
||||
redisClient *c = listNodeValue(ln);
|
||||
|
||||
if (c->flags & REDIS_BLOCKED) {
|
||||
addReplySds(c,sdsnew(
|
||||
"-UNBLOCKED force unblock from blocking operation, "
|
||||
"instance state changed (master -> slave?)\r\n"));
|
||||
unblockClient(c);
|
||||
c->flags |= REDIS_CLOSE_AFTER_REPLY;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+287
-131
@@ -74,27 +74,13 @@ void clusterCloseAllSlots(void);
|
||||
void clusterSetNodeAsMaster(clusterNode *n);
|
||||
void clusterDelNode(clusterNode *delnode);
|
||||
sds representRedisNodeFlags(sds ci, uint16_t flags);
|
||||
uint64_t clusterGetMaxEpoch(void);
|
||||
int clusterBumpConfigEpochWithoutConsensus(void);
|
||||
|
||||
/* -----------------------------------------------------------------------------
|
||||
* Initialization
|
||||
* -------------------------------------------------------------------------- */
|
||||
|
||||
/* Return the greatest configEpoch found in the cluster. */
|
||||
uint64_t clusterGetMaxEpoch(void) {
|
||||
uint64_t max = 0;
|
||||
dictIterator *di;
|
||||
dictEntry *de;
|
||||
|
||||
di = dictGetSafeIterator(server.cluster->nodes);
|
||||
while((de = dictNext(di)) != NULL) {
|
||||
clusterNode *node = dictGetVal(de);
|
||||
if (node->configEpoch > max) max = node->configEpoch;
|
||||
}
|
||||
dictReleaseIterator(di);
|
||||
if (max < server.cluster->currentEpoch) max = server.cluster->currentEpoch;
|
||||
return max;
|
||||
}
|
||||
|
||||
/* Load the cluster config from 'filename'.
|
||||
*
|
||||
* If the file does not exist or is zero-length (this may happen because
|
||||
@@ -927,6 +913,138 @@ void clusterRenameNode(clusterNode *node, char *newname) {
|
||||
clusterAddNode(node);
|
||||
}
|
||||
|
||||
/* -----------------------------------------------------------------------------
|
||||
* CLUSTER config epoch handling
|
||||
* -------------------------------------------------------------------------- */
|
||||
|
||||
/* Return the greatest configEpoch found in the cluster, or the current
|
||||
* epoch if greater than any node configEpoch. */
|
||||
uint64_t clusterGetMaxEpoch(void) {
|
||||
uint64_t max = 0;
|
||||
dictIterator *di;
|
||||
dictEntry *de;
|
||||
|
||||
di = dictGetSafeIterator(server.cluster->nodes);
|
||||
while((de = dictNext(di)) != NULL) {
|
||||
clusterNode *node = dictGetVal(de);
|
||||
if (node->configEpoch > max) max = node->configEpoch;
|
||||
}
|
||||
dictReleaseIterator(di);
|
||||
if (max < server.cluster->currentEpoch) max = server.cluster->currentEpoch;
|
||||
return max;
|
||||
}
|
||||
|
||||
/* If this node epoch is zero or is not already the greatest across the
|
||||
* cluster (from the POV of the local configuration), this function will:
|
||||
*
|
||||
* 1) Generate a new config epoch increment the current epoch.
|
||||
* 2) Assign the new epoch to this node, WITHOUT any consensus.
|
||||
* 3) Persist the configuration on disk before sending packets with the
|
||||
* new configuration.
|
||||
*
|
||||
* If the new config epoch is generated and assigend, REDIS_OK is returned,
|
||||
* otherwise REDIS_ERR is returned (since the node has already the greatest
|
||||
* configuration around) and no operation is performed.
|
||||
*
|
||||
* Important note: this function violates the principle that config epochs
|
||||
* should be generated with consensus and should be unique across the cluster.
|
||||
* However Redis Cluster uses this auto-generated new config epochs in two
|
||||
* cases:
|
||||
*
|
||||
* 1) When slots are closed after importing. Otherwise resharding would be
|
||||
* too exansive.
|
||||
* 2) When CLUSTER FAILOVER is called with options that force a slave to
|
||||
* failover its master even if there is not master majority able to
|
||||
* create a new configuration epoch.
|
||||
*
|
||||
* Redis Cluster does not explode using this function, even in the case of
|
||||
* a collision between this node and another node, generating the same
|
||||
* configuration epoch unilaterally, because the config epoch conflict
|
||||
* resolution algorithm will eventually move colliding nodes to different
|
||||
* config epochs. However usign this function may violate the "last failover
|
||||
* wins" rule, so should only be used with care. */
|
||||
int clusterBumpConfigEpochWithoutConsensus(void) {
|
||||
uint64_t maxEpoch = clusterGetMaxEpoch();
|
||||
|
||||
if (myself->configEpoch == 0 ||
|
||||
myself->configEpoch != maxEpoch)
|
||||
{
|
||||
server.cluster->currentEpoch++;
|
||||
myself->configEpoch = server.cluster->currentEpoch;
|
||||
clusterDoBeforeSleep(CLUSTER_TODO_SAVE_CONFIG|
|
||||
CLUSTER_TODO_FSYNC_CONFIG);
|
||||
redisLog(REDIS_WARNING,
|
||||
"New configEpoch set to %llu",
|
||||
(unsigned long long) myself->configEpoch);
|
||||
return REDIS_OK;
|
||||
} else {
|
||||
return REDIS_ERR;
|
||||
}
|
||||
}
|
||||
|
||||
/* This function is called when this node is a master, and we receive from
|
||||
* another master a configuration epoch that is equal to our configuration
|
||||
* epoch.
|
||||
*
|
||||
* BACKGROUND
|
||||
*
|
||||
* It is not possible that different slaves get the same config
|
||||
* epoch during a failover election, because the slaves need to get voted
|
||||
* by a majority. However when we perform a manual resharding of the cluster
|
||||
* the node will assign a configuration epoch to itself without to ask
|
||||
* for agreement. Usually resharding happens when the cluster is working well
|
||||
* and is supervised by the sysadmin, however it is possible for a failover
|
||||
* to happen exactly while the node we are resharding a slot to assigns itself
|
||||
* a new configuration epoch, but before it is able to propagate it.
|
||||
*
|
||||
* So technically it is possible in this condition that two nodes end with
|
||||
* the same configuration epoch.
|
||||
*
|
||||
* Another possibility is that there are bugs in the implementation causing
|
||||
* this to happen.
|
||||
*
|
||||
* Moreover when a new cluster is created, all the nodes start with the same
|
||||
* configEpoch. This collision resolution code allows nodes to automatically
|
||||
* end with a different configEpoch at startup automatically.
|
||||
*
|
||||
* In all the cases, we want a mechanism that resolves this issue automatically
|
||||
* as a safeguard. The same configuration epoch for masters serving different
|
||||
* set of slots is not harmful, but it is if the nodes end serving the same
|
||||
* slots for some reason (manual errors or software bugs) without a proper
|
||||
* failover procedure.
|
||||
*
|
||||
* In general we want a system that eventually always ends with different
|
||||
* masters having different configuration epochs whatever happened, since
|
||||
* nothign is worse than a split-brain condition in a distributed system.
|
||||
*
|
||||
* BEHAVIOR
|
||||
*
|
||||
* When this function gets called, what happens is that if this node
|
||||
* has the lexicographically smaller Node ID compared to the other node
|
||||
* with the conflicting epoch (the 'sender' node), it will assign itself
|
||||
* the greatest configuration epoch currently detected among nodes plus 1.
|
||||
*
|
||||
* This means that even if there are multiple nodes colliding, the node
|
||||
* with the greatest Node ID never moves forward, so eventually all the nodes
|
||||
* end with a different configuration epoch.
|
||||
*/
|
||||
void clusterHandleConfigEpochCollision(clusterNode *sender) {
|
||||
/* Prerequisites: nodes have the same configEpoch and are both masters. */
|
||||
if (sender->configEpoch != myself->configEpoch ||
|
||||
!nodeIsMaster(sender) || !nodeIsMaster(myself)) return;
|
||||
/* Don't act if the colliding node has a smaller Node ID. */
|
||||
if (memcmp(sender->name,myself->name,REDIS_CLUSTER_NAMELEN) <= 0) return;
|
||||
/* Get the next ID available at the best of this node knowledge. */
|
||||
server.cluster->currentEpoch++;
|
||||
myself->configEpoch = server.cluster->currentEpoch;
|
||||
clusterSaveConfigOrDie(1);
|
||||
redisLog(REDIS_VERBOSE,
|
||||
"WARNING: configEpoch collision with node %.40s."
|
||||
" configEpoch set to %llu",
|
||||
sender->name,
|
||||
(unsigned long long) myself->configEpoch);
|
||||
}
|
||||
|
||||
/* -----------------------------------------------------------------------------
|
||||
* CLUSTER nodes blacklist
|
||||
*
|
||||
@@ -1399,69 +1517,6 @@ void clusterUpdateSlotsConfigWith(clusterNode *sender, uint64_t senderConfigEpoc
|
||||
}
|
||||
}
|
||||
|
||||
/* This function is called when this node is a master, and we receive from
|
||||
* another master a configuration epoch that is equal to our configuration
|
||||
* epoch.
|
||||
*
|
||||
* BACKGROUND
|
||||
*
|
||||
* It is not possible that different slaves get the same config
|
||||
* epoch during a failover election, because the slaves need to get voted
|
||||
* by a majority. However when we perform a manual resharding of the cluster
|
||||
* the node will assign a configuration epoch to itself without to ask
|
||||
* for agreement. Usually resharding happens when the cluster is working well
|
||||
* and is supervised by the sysadmin, however it is possible for a failover
|
||||
* to happen exactly while the node we are resharding a slot to assigns itself
|
||||
* a new configuration epoch, but before it is able to propagate it.
|
||||
*
|
||||
* So technically it is possible in this condition that two nodes end with
|
||||
* the same configuration epoch.
|
||||
*
|
||||
* Another possibility is that there are bugs in the implementation causing
|
||||
* this to happen.
|
||||
*
|
||||
* Moreover when a new cluster is created, all the nodes start with the same
|
||||
* configEpoch. This collision resolution code allows nodes to automatically
|
||||
* end with a different configEpoch at startup automatically.
|
||||
*
|
||||
* In all the cases, we want a mechanism that resolves this issue automatically
|
||||
* as a safeguard. The same configuration epoch for masters serving different
|
||||
* set of slots is not harmful, but it is if the nodes end serving the same
|
||||
* slots for some reason (manual errors or software bugs) without a proper
|
||||
* failover procedure.
|
||||
*
|
||||
* In general we want a system that eventually always ends with different
|
||||
* masters having different configuration epochs whatever happened, since
|
||||
* nothign is worse than a split-brain condition in a distributed system.
|
||||
*
|
||||
* BEHAVIOR
|
||||
*
|
||||
* When this function gets called, what happens is that if this node
|
||||
* has the lexicographically smaller Node ID compared to the other node
|
||||
* with the conflicting epoch (the 'sender' node), it will assign itself
|
||||
* the greatest configuration epoch currently detected among nodes plus 1.
|
||||
*
|
||||
* This means that even if there are multiple nodes colliding, the node
|
||||
* with the greatest Node ID never moves forward, so eventually all the nodes
|
||||
* end with a different configuration epoch.
|
||||
*/
|
||||
void clusterHandleConfigEpochCollision(clusterNode *sender) {
|
||||
/* Prerequisites: nodes have the same configEpoch and are both masters. */
|
||||
if (sender->configEpoch != myself->configEpoch ||
|
||||
!nodeIsMaster(sender) || !nodeIsMaster(myself)) return;
|
||||
/* Don't act if the colliding node has a smaller Node ID. */
|
||||
if (memcmp(sender->name,myself->name,REDIS_CLUSTER_NAMELEN) <= 0) return;
|
||||
/* Get the next ID available at the best of this node knowledge. */
|
||||
server.cluster->currentEpoch++;
|
||||
myself->configEpoch = server.cluster->currentEpoch;
|
||||
clusterSaveConfigOrDie(1);
|
||||
redisLog(REDIS_VERBOSE,
|
||||
"WARNING: configEpoch collision with node %.40s."
|
||||
" configEpoch set to %llu",
|
||||
sender->name,
|
||||
(unsigned long long) myself->configEpoch);
|
||||
}
|
||||
|
||||
/* When this function is called, there is a packet to process starting
|
||||
* at node->rcvbuf. Releasing the buffer is up to the caller, so this
|
||||
* function should just handle the higher level stuff of processing the
|
||||
@@ -2582,6 +2637,42 @@ void clusterLogCantFailover(int reason) {
|
||||
redisLog(REDIS_WARNING,"Currently unable to failover: %s", msg);
|
||||
}
|
||||
|
||||
/* This function implements the final part of automatic and manual failovers,
|
||||
* where the slave grabs its master's hash slots, and propagates the new
|
||||
* configuration.
|
||||
*
|
||||
* Note that it's up to the caller to be sure that the node got a new
|
||||
* configuration epoch already. */
|
||||
void clusterFailoverReplaceYourMaster(void) {
|
||||
int j;
|
||||
clusterNode *oldmaster = myself->slaveof;
|
||||
|
||||
if (nodeIsMaster(myself) || oldmaster == NULL) return;
|
||||
|
||||
/* 1) Turn this node into a master. */
|
||||
clusterSetNodeAsMaster(myself);
|
||||
replicationUnsetMaster();
|
||||
|
||||
/* 2) Claim all the slots assigned to our master. */
|
||||
for (j = 0; j < REDIS_CLUSTER_SLOTS; j++) {
|
||||
if (clusterNodeGetSlotBit(oldmaster,j)) {
|
||||
clusterDelSlot(j);
|
||||
clusterAddSlot(myself,j);
|
||||
}
|
||||
}
|
||||
|
||||
/* 3) Update state and save config. */
|
||||
clusterUpdateState();
|
||||
clusterSaveConfigOrDie(1);
|
||||
|
||||
/* 4) Pong all the other nodes so that they can update the state
|
||||
* accordingly and detect that we switched to master role. */
|
||||
clusterBroadcastPong(CLUSTER_BROADCAST_ALL);
|
||||
|
||||
/* 5) If there was a manual failover in progress, clear the state. */
|
||||
resetManualFailover();
|
||||
}
|
||||
|
||||
/* This function is called if we are a slave node and our master serving
|
||||
* a non-zero amount of hash slots is in FAIL state.
|
||||
*
|
||||
@@ -2596,7 +2687,6 @@ void clusterHandleSlaveFailover(void) {
|
||||
int needed_quorum = (server.cluster->size / 2) + 1;
|
||||
int manual_failover = server.cluster->mf_end != 0 &&
|
||||
server.cluster->mf_can_start;
|
||||
int j;
|
||||
mstime_t auth_timeout, auth_retry_time;
|
||||
|
||||
server.cluster->todo_before_sleep &= ~CLUSTER_TODO_HANDLE_FAILOVER;
|
||||
@@ -2738,26 +2828,12 @@ void clusterHandleSlaveFailover(void) {
|
||||
|
||||
/* Check if we reached the quorum. */
|
||||
if (server.cluster->failover_auth_count >= needed_quorum) {
|
||||
clusterNode *oldmaster = myself->slaveof;
|
||||
/* We have the quorum, we can finally failover the master. */
|
||||
|
||||
redisLog(REDIS_WARNING,
|
||||
"Failover election won: I'm the new master.");
|
||||
/* We have the quorum, perform all the steps to correctly promote
|
||||
* this slave to a master.
|
||||
*
|
||||
* 1) Turn this node into a master. */
|
||||
clusterSetNodeAsMaster(myself);
|
||||
replicationUnsetMaster();
|
||||
|
||||
/* 2) Claim all the slots assigned to our master. */
|
||||
for (j = 0; j < REDIS_CLUSTER_SLOTS; j++) {
|
||||
if (clusterNodeGetSlotBit(oldmaster,j)) {
|
||||
clusterDelSlot(j);
|
||||
clusterAddSlot(myself,j);
|
||||
}
|
||||
}
|
||||
|
||||
/* 3) Update my configEpoch to the epoch of the election. */
|
||||
/* Update my configEpoch to the epoch of the election. */
|
||||
if (myself->configEpoch < server.cluster->failover_auth_epoch) {
|
||||
myself->configEpoch = server.cluster->failover_auth_epoch;
|
||||
redisLog(REDIS_WARNING,
|
||||
@@ -2765,16 +2841,8 @@ void clusterHandleSlaveFailover(void) {
|
||||
(unsigned long long) myself->configEpoch);
|
||||
}
|
||||
|
||||
/* 4) Update state and save config. */
|
||||
clusterUpdateState();
|
||||
clusterSaveConfigOrDie(1);
|
||||
|
||||
/* 5) Pong all the other nodes so that they can update the state
|
||||
* accordingly and detect that we switched to master role. */
|
||||
clusterBroadcastPong(CLUSTER_BROADCAST_ALL);
|
||||
|
||||
/* 6) If there was a manual failover in progress, clear the state. */
|
||||
resetManualFailover();
|
||||
/* Take responsability for the cluster slots. */
|
||||
clusterFailoverReplaceYourMaster();
|
||||
} else {
|
||||
clusterLogCantFailover(REDIS_CLUSTER_CANT_FAILOVER_WAITING_VOTES);
|
||||
}
|
||||
@@ -3567,7 +3635,7 @@ sds clusterGenNodeDescription(clusterNode *node) {
|
||||
else
|
||||
ci = sdscatlen(ci," - ",3);
|
||||
|
||||
/* Latency from the POV of this node, link status */
|
||||
/* Latency from the POV of this node, config epoch, link status */
|
||||
ci = sdscatprintf(ci,"%lld %lld %llu %s",
|
||||
(long long) node->ping_sent,
|
||||
(long long) node->pong_received,
|
||||
@@ -3902,17 +3970,9 @@ void clusterCommand(redisClient *c) {
|
||||
* failover happens at the same time we close the slot, the
|
||||
* configEpoch collision resolution will fix it assigning
|
||||
* a different epoch to each node. */
|
||||
uint64_t maxEpoch = clusterGetMaxEpoch();
|
||||
|
||||
if (myself->configEpoch == 0 ||
|
||||
myself->configEpoch != maxEpoch)
|
||||
{
|
||||
server.cluster->currentEpoch++;
|
||||
myself->configEpoch = server.cluster->currentEpoch;
|
||||
clusterDoBeforeSleep(CLUSTER_TODO_FSYNC_CONFIG);
|
||||
if (clusterBumpConfigEpochWithoutConsensus() == REDIS_OK) {
|
||||
redisLog(REDIS_WARNING,
|
||||
"configEpoch set to %llu after importing slot %d",
|
||||
(unsigned long long) myself->configEpoch, slot);
|
||||
"configEpoch updated after importing slot %d", slot);
|
||||
}
|
||||
server.cluster->importing_slots_from[slot] = NULL;
|
||||
}
|
||||
@@ -4115,24 +4175,31 @@ void clusterCommand(redisClient *c) {
|
||||
} else if (!strcasecmp(c->argv[1]->ptr,"failover") &&
|
||||
(c->argc == 2 || c->argc == 3))
|
||||
{
|
||||
/* CLUSTER FAILOVER [FORCE] */
|
||||
int force = 0;
|
||||
/* CLUSTER FAILOVER [FORCE|TAKEOVER] */
|
||||
int force = 0, takeover = 0;
|
||||
|
||||
if (c->argc == 3) {
|
||||
if (!strcasecmp(c->argv[2]->ptr,"force")) {
|
||||
force = 1;
|
||||
} else if (!strcasecmp(c->argv[2]->ptr,"takeover")) {
|
||||
takeover = 1;
|
||||
force = 1; /* Takeover also implies force. */
|
||||
} else {
|
||||
addReply(c,shared.syntaxerr);
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
||||
/* Check preconditions. */
|
||||
if (nodeIsMaster(myself)) {
|
||||
addReplyError(c,"You should send CLUSTER FAILOVER to a slave");
|
||||
return;
|
||||
} else if (myself->slaveof == NULL) {
|
||||
addReplyError(c,"I'm a slave but my master is unknown to me");
|
||||
return;
|
||||
} else if (!force &&
|
||||
(myself->slaveof == NULL || nodeFailed(myself->slaveof) ||
|
||||
myself->slaveof->link == NULL))
|
||||
(nodeFailed(myself->slaveof) ||
|
||||
myself->slaveof->link == NULL))
|
||||
{
|
||||
addReplyError(c,"Master is down or failed, "
|
||||
"please use CLUSTER FAILOVER FORCE");
|
||||
@@ -4141,15 +4208,24 @@ void clusterCommand(redisClient *c) {
|
||||
resetManualFailover();
|
||||
server.cluster->mf_end = mstime() + REDIS_CLUSTER_MF_TIMEOUT;
|
||||
|
||||
/* If this is a forced failover, we don't need to talk with our master
|
||||
* to agree about the offset. We just failover taking over it without
|
||||
* coordination. */
|
||||
if (force) {
|
||||
if (takeover) {
|
||||
/* A takeover does not perform any initial check. It just
|
||||
* generates a new configuration epoch for this node without
|
||||
* consensus, claims the master's slots, and broadcast the new
|
||||
* configuration. */
|
||||
redisLog(REDIS_WARNING,"Taking over the master (user request).");
|
||||
clusterBumpConfigEpochWithoutConsensus();
|
||||
clusterFailoverReplaceYourMaster();
|
||||
} else if (force) {
|
||||
/* If this is a forced failover, we don't need to talk with our
|
||||
* master to agree about the offset. We just failover taking over
|
||||
* it without coordination. */
|
||||
redisLog(REDIS_WARNING,"Forced failover user request accepted.");
|
||||
server.cluster->mf_can_start = 1;
|
||||
} else {
|
||||
redisLog(REDIS_WARNING,"Manual failover user request accepted.");
|
||||
clusterSendMFStart(myself->slaveof);
|
||||
}
|
||||
redisLog(REDIS_WARNING,"Manual failover user request accepted.");
|
||||
addReply(c,shared.ok);
|
||||
} else if (!strcasecmp(c->argv[1]->ptr,"set-config-epoch") && c->argc == 3)
|
||||
{
|
||||
@@ -4692,10 +4768,10 @@ void readwriteCommand(redisClient *c) {
|
||||
* belonging to the same slot, but the slot is not stable (in migration or
|
||||
* importing state, likely because a resharding is in progress).
|
||||
*
|
||||
* REDIS_CLUSTER_REDIR_DOWN if the request addresses a slot which is not
|
||||
* bound to any node. In this case the cluster global state should be already
|
||||
* "down" but it is fragile to rely on the update of the global state, so
|
||||
* we also handle it here. */
|
||||
* REDIS_CLUSTER_REDIR_DOWN_UNBOUND if the request addresses a slot which is
|
||||
* not bound to any node. In this case the cluster global state should be
|
||||
* already "down" but it is fragile to rely on the update of the global state,
|
||||
* so we also handle it here. */
|
||||
clusterNode *getNodeByQuery(redisClient *c, struct redisCommand *cmd, robj **argv, int argc, int *hashslot, int *error_code) {
|
||||
clusterNode *n = NULL;
|
||||
robj *firstkey = NULL;
|
||||
@@ -4757,7 +4833,7 @@ clusterNode *getNodeByQuery(redisClient *c, struct redisCommand *cmd, robj **arg
|
||||
if (n == NULL) {
|
||||
getKeysFreeResult(keyindex);
|
||||
if (error_code)
|
||||
*error_code = REDIS_CLUSTER_REDIR_DOWN;
|
||||
*error_code = REDIS_CLUSTER_REDIR_DOWN_UNBOUND;
|
||||
return NULL;
|
||||
}
|
||||
|
||||
@@ -4849,3 +4925,83 @@ clusterNode *getNodeByQuery(redisClient *c, struct redisCommand *cmd, robj **arg
|
||||
if (n != myself && error_code) *error_code = REDIS_CLUSTER_REDIR_MOVED;
|
||||
return n;
|
||||
}
|
||||
|
||||
/* Send the client the right redirection code, according to error_code
|
||||
* that should be set to one of REDIS_CLUSTER_REDIR_* macros.
|
||||
*
|
||||
* If REDIS_CLUSTER_REDIR_ASK or REDIS_CLUSTER_REDIR_MOVED error codes
|
||||
* are used, then the node 'n' should not be NULL, but should be the
|
||||
* node we want to mention in the redirection. Moreover hashslot should
|
||||
* be set to the hash slot that caused the redirection. */
|
||||
void clusterRedirectClient(redisClient *c, clusterNode *n, int hashslot, int error_code) {
|
||||
if (error_code == REDIS_CLUSTER_REDIR_CROSS_SLOT) {
|
||||
addReplySds(c,sdsnew("-CROSSSLOT Keys in request don't hash to the same slot\r\n"));
|
||||
} else if (error_code == REDIS_CLUSTER_REDIR_UNSTABLE) {
|
||||
/* The request spawns mutliple keys in the same slot,
|
||||
* but the slot is not "stable" currently as there is
|
||||
* a migration or import in progress. */
|
||||
addReplySds(c,sdsnew("-TRYAGAIN Multiple keys request during rehashing of slot\r\n"));
|
||||
} else if (error_code == REDIS_CLUSTER_REDIR_DOWN_STATE) {
|
||||
addReplySds(c,sdsnew("-CLUSTERDOWN The cluster is down\r\n"));
|
||||
} else if (error_code == REDIS_CLUSTER_REDIR_DOWN_UNBOUND) {
|
||||
addReplySds(c,sdsnew("-CLUSTERDOWN Hash slot not served\r\n"));
|
||||
} else if (error_code == REDIS_CLUSTER_REDIR_MOVED ||
|
||||
error_code == REDIS_CLUSTER_REDIR_ASK)
|
||||
{
|
||||
addReplySds(c,sdscatprintf(sdsempty(),
|
||||
"-%s %d %s:%d\r\n",
|
||||
(error_code == REDIS_CLUSTER_REDIR_ASK) ? "ASK" : "MOVED",
|
||||
hashslot,n->ip,n->port));
|
||||
} else {
|
||||
redisPanic("getNodeByQuery() unknown error.");
|
||||
}
|
||||
}
|
||||
|
||||
/* This function is called by the function processing clients incrementally
|
||||
* to detect timeouts, in order to handle the following case:
|
||||
*
|
||||
* 1) A client blocks with BLPOP or similar blocking operation.
|
||||
* 2) The master migrates the hash slot elsewhere or turns into a slave.
|
||||
* 3) The client may remain blocked forever (or up to the max timeout time)
|
||||
* waiting for a key change that will never happen.
|
||||
*
|
||||
* If the client is found to be blocked into an hash slot this node no
|
||||
* longer handles, the client is sent a redirection error, and the function
|
||||
* returns 1. Otherwise 0 is returned and no operation is performed. */
|
||||
int clusterRedirectBlockedClientIfNeeded(redisClient *c) {
|
||||
if (c->flags & REDIS_BLOCKED && c->btype == REDIS_BLOCKED_LIST) {
|
||||
dictEntry *de;
|
||||
dictIterator *di;
|
||||
|
||||
/* If the cluster is down, unblock the client with the right error. */
|
||||
if (server.cluster->state == REDIS_CLUSTER_FAIL) {
|
||||
clusterRedirectClient(c,NULL,0,REDIS_CLUSTER_REDIR_DOWN_STATE);
|
||||
return 1;
|
||||
}
|
||||
|
||||
di = dictGetIterator(c->bpop.keys);
|
||||
while((de = dictNext(di)) != NULL) {
|
||||
robj *key = dictGetKey(de);
|
||||
int slot = keyHashSlot((char*)key->ptr, sdslen(key->ptr));
|
||||
clusterNode *node = server.cluster->slots[slot];
|
||||
|
||||
/* 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. */
|
||||
if (node != myself &&
|
||||
server.cluster->importing_slots_from[slot] == NULL)
|
||||
{
|
||||
if (node == NULL) {
|
||||
clusterRedirectClient(c,NULL,0,
|
||||
REDIS_CLUSTER_REDIR_DOWN_UNBOUND);
|
||||
} else {
|
||||
clusterRedirectClient(c,node,slot,
|
||||
REDIS_CLUSTER_REDIR_MOVED);
|
||||
}
|
||||
return 1;
|
||||
}
|
||||
}
|
||||
dictReleaseIterator(di);
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
|
||||
+6
-3
@@ -26,11 +26,12 @@
|
||||
|
||||
/* Redirection errors returned by getNodeByQuery(). */
|
||||
#define REDIS_CLUSTER_REDIR_NONE 0 /* Node can serve the request. */
|
||||
#define REDIS_CLUSTER_REDIR_CROSS_SLOT 1 /* Keys in different slots. */
|
||||
#define REDIS_CLUSTER_REDIR_UNSTABLE 2 /* Keys in slot resharding. */
|
||||
#define REDIS_CLUSTER_REDIR_CROSS_SLOT 1 /* -CROSSSLOT request. */
|
||||
#define REDIS_CLUSTER_REDIR_UNSTABLE 2 /* -TRYAGAIN redirection required */
|
||||
#define REDIS_CLUSTER_REDIR_ASK 3 /* -ASK redirection required. */
|
||||
#define REDIS_CLUSTER_REDIR_MOVED 4 /* -MOVED redirection required. */
|
||||
#define REDIS_CLUSTER_REDIR_DOWN 5 /* -CLUSTERDOWN error. */
|
||||
#define REDIS_CLUSTER_REDIR_DOWN_STATE 5 /* -CLUSTERDOWN, global state. */
|
||||
#define REDIS_CLUSTER_REDIR_DOWN_UNBOUND 6 /* -CLUSTERDOWN, unbound slot. */
|
||||
|
||||
struct clusterNode;
|
||||
|
||||
@@ -249,5 +250,7 @@ typedef struct {
|
||||
|
||||
/* ---------------------- API exported outside cluster.c -------------------- */
|
||||
clusterNode *getNodeByQuery(redisClient *c, struct redisCommand *cmd, robj **argv, int argc, int *hashslot, int *ask);
|
||||
int clusterRedirectBlockedClientIfNeeded(redisClient *c);
|
||||
void clusterRedirectClient(redisClient *c, clusterNode *n, int hashslot, int error_code);
|
||||
|
||||
#endif /* __REDIS_CLUSTER_H */
|
||||
|
||||
+1
-76
@@ -661,81 +661,6 @@ dictEntry *dictGetRandomKey(dict *d)
|
||||
return he;
|
||||
}
|
||||
|
||||
/* XXX: This is going to be removed soon and SPOP internals
|
||||
* reimplemented.
|
||||
*
|
||||
* This is a version of dictGetRandomKey() that is modified in order to
|
||||
* return multiple entries by jumping at a random place of the hash table
|
||||
* and scanning linearly for entries.
|
||||
*
|
||||
* Returned pointers to hash table entries are stored into 'des' that
|
||||
* points to an array of dictEntry pointers. The array must have room for
|
||||
* at least 'count' elements, that is the argument we pass to the function
|
||||
* to tell how many random elements we need.
|
||||
*
|
||||
* The function returns the number of items stored into 'des', that may
|
||||
* be less than 'count' if the hash table has less than 'count' elements
|
||||
* inside.
|
||||
*
|
||||
* Note that this function is not suitable when you need a good distribution
|
||||
* of the returned items, but only when you need to "sample" a given number
|
||||
* of continuous elements to run some kind of algorithm or to produce
|
||||
* statistics. However the function is much faster than dictGetRandomKey()
|
||||
* at producing N elements, and the elements are guaranteed to be non
|
||||
* repeating. */
|
||||
unsigned int dictGetRandomKeys(dict *d, dictEntry **des, unsigned int count) {
|
||||
unsigned int j; /* internal hash table id, 0 or 1. */
|
||||
unsigned int tables; /* 1 or 2 tables? */
|
||||
unsigned int stored = 0, maxsizemask;
|
||||
|
||||
if (dictSize(d) < count) count = dictSize(d);
|
||||
|
||||
/* Try to do a rehashing work proportional to 'count'. */
|
||||
for (j = 0; j < count; j++) {
|
||||
if (dictIsRehashing(d))
|
||||
_dictRehashStep(d);
|
||||
else
|
||||
break;
|
||||
}
|
||||
|
||||
tables = dictIsRehashing(d) ? 2 : 1;
|
||||
maxsizemask = d->ht[0].sizemask;
|
||||
if (tables > 1 && maxsizemask < d->ht[1].sizemask)
|
||||
maxsizemask = d->ht[1].sizemask;
|
||||
|
||||
/* Pick a random point inside the larger table. */
|
||||
unsigned int i = random() & maxsizemask;
|
||||
while(stored < count) {
|
||||
for (j = 0; j < tables; j++) {
|
||||
/* Invariant of the dict.c rehashing: up to the indexes already
|
||||
* visited in ht[0] during the rehashing, there are no populated
|
||||
* buckets, so we can skip ht[0] for indexes between 0 and idx-1. */
|
||||
if (tables == 2 && j == 0 && i < d->rehashidx) {
|
||||
/* Moreover, if we are currently out of range in the second
|
||||
* table, there will be no elements in both tables up to
|
||||
* the current rehashing index, so we jump if possible.
|
||||
* (this happens when going from big to small table). */
|
||||
if (i >= d->ht[1].size) i = d->rehashidx;
|
||||
continue;
|
||||
}
|
||||
if (i >= d->ht[j].size) continue; /* Out of range for this table. */
|
||||
dictEntry *he = d->ht[j].table[i];
|
||||
while (he) {
|
||||
/* Collect all the elements of the buckets found non
|
||||
* empty while iterating. */
|
||||
*des = he;
|
||||
des++;
|
||||
he = he->next;
|
||||
stored++;
|
||||
if (stored == count) return stored;
|
||||
}
|
||||
}
|
||||
i = (i+1) & maxsizemask;
|
||||
}
|
||||
return stored; /* Never reached. */
|
||||
}
|
||||
|
||||
|
||||
/* This function samples the dictionary to return a few keys from random
|
||||
* locations.
|
||||
*
|
||||
@@ -788,7 +713,7 @@ unsigned int dictGetSomeKeys(dict *d, dictEntry **des, unsigned int count) {
|
||||
/* Invariant of the dict.c rehashing: up to the indexes already
|
||||
* visited in ht[0] during the rehashing, there are no populated
|
||||
* buckets, so we can skip ht[0] for indexes between 0 and idx-1. */
|
||||
if (tables == 2 && j == 0 && i < d->rehashidx) {
|
||||
if (tables == 2 && j == 0 && i < (unsigned int) d->rehashidx) {
|
||||
/* Moreover, if we are currently out of range in the second
|
||||
* table, there will be no elements in both tables up to
|
||||
* the current rehashing index, so we jump if possible.
|
||||
|
||||
@@ -165,7 +165,6 @@ dictEntry *dictNext(dictIterator *iter);
|
||||
void dictReleaseIterator(dictIterator *iter);
|
||||
dictEntry *dictGetRandomKey(dict *d);
|
||||
unsigned int dictGetSomeKeys(dict *d, dictEntry **des, unsigned int count);
|
||||
unsigned int dictGetRandomKeys(dict *d, dictEntry **des, unsigned int count);
|
||||
void dictPrintStats(dict *d);
|
||||
unsigned int dictGenHashFunction(const void *key, int len);
|
||||
unsigned int dictGenCaseHashFunction(const unsigned char *buf, int len);
|
||||
|
||||
+40
-9
@@ -135,23 +135,49 @@ redisClient *createClient(int fd) {
|
||||
* returns REDIS_OK, and make sure to install the write handler in our event
|
||||
* loop so that when the socket is writable new data gets written.
|
||||
*
|
||||
* If the client should not receive new data, because it is a fake client,
|
||||
* a master, a slave not yet online, or because the setup of the write handler
|
||||
* failed, the function returns REDIS_ERR.
|
||||
* If the client should not receive new data, because it is a fake client
|
||||
* (used to load AOF in memory), a master or because the setup of the write
|
||||
* handler failed, the function returns REDIS_ERR.
|
||||
*
|
||||
* The function may return REDIS_OK without actually installing the write
|
||||
* event handler in the following cases:
|
||||
*
|
||||
* 1) The event handler should already be installed since the output buffer
|
||||
* already contained something.
|
||||
* 2) The client is a slave but not yet online, so we want to just accumulate
|
||||
* writes in the buffer but not actually sending them yet.
|
||||
*
|
||||
* Typically gets called every time a reply is built, before adding more
|
||||
* data to the clients output buffers. If the function returns REDIS_ERR no
|
||||
* data should be appended to the output buffers. */
|
||||
int prepareClientToWrite(redisClient *c) {
|
||||
/* If it's the Lua client we always return ok without installing any
|
||||
* handler since there is no socket at all. */
|
||||
if (c->flags & REDIS_LUA_CLIENT) return REDIS_OK;
|
||||
|
||||
/* Masters don't receive replies, unless REDIS_MASTER_FORCE_REPLY flag
|
||||
* is set. */
|
||||
if ((c->flags & REDIS_MASTER) &&
|
||||
!(c->flags & REDIS_MASTER_FORCE_REPLY)) return REDIS_ERR;
|
||||
if (c->fd <= 0) return REDIS_ERR; /* Fake client */
|
||||
|
||||
if (c->fd <= 0) return REDIS_ERR; /* Fake client for AOF loading. */
|
||||
|
||||
/* Only install the handler if not already installed and, in case of
|
||||
* slaves, if the client can actually receive writes. */
|
||||
if (c->bufpos == 0 && listLength(c->reply) == 0 &&
|
||||
(c->replstate == REDIS_REPL_NONE ||
|
||||
c->replstate == REDIS_REPL_ONLINE) &&
|
||||
aeCreateFileEvent(server.el, c->fd, AE_WRITABLE,
|
||||
sendReplyToClient, c) == AE_ERR) return REDIS_ERR;
|
||||
(c->replstate == REDIS_REPL_ONLINE && !c->repl_put_online_on_ack)))
|
||||
{
|
||||
/* Try to install the write handler. */
|
||||
if (aeCreateFileEvent(server.el, c->fd, AE_WRITABLE,
|
||||
sendReplyToClient, c) == AE_ERR)
|
||||
{
|
||||
freeClientAsync(c);
|
||||
return REDIS_ERR;
|
||||
}
|
||||
}
|
||||
|
||||
/* Authorize the caller to queue in the output buffer of this client. */
|
||||
return REDIS_OK;
|
||||
}
|
||||
|
||||
@@ -1689,7 +1715,9 @@ void pauseClients(mstime_t end) {
|
||||
/* Return non-zero if clients are currently paused. As a side effect the
|
||||
* function checks if the pause time was reached and clear it. */
|
||||
int clientsArePaused(void) {
|
||||
if (server.clients_paused && server.clients_pause_end_time < server.mstime) {
|
||||
if (server.clients_paused &&
|
||||
server.clients_pause_end_time < server.mstime)
|
||||
{
|
||||
listNode *ln;
|
||||
listIter li;
|
||||
redisClient *c;
|
||||
@@ -1702,7 +1730,10 @@ int clientsArePaused(void) {
|
||||
while ((ln = listNext(&li)) != NULL) {
|
||||
c = listNodeValue(ln);
|
||||
|
||||
if (c->flags & REDIS_SLAVE) continue;
|
||||
/* Don't touch slaves and blocked clients. The latter pending
|
||||
* requests be processed when unblocked. */
|
||||
if (c->flags & (REDIS_SLAVE|REDIS_BLOCKED)) continue;
|
||||
c->flags |= REDIS_UNBLOCKED;
|
||||
listAddNodeTail(server.unblocked_clients,c);
|
||||
}
|
||||
}
|
||||
|
||||
+12
-22
@@ -923,8 +923,14 @@ int clientsCronHandleTimeout(redisClient *c) {
|
||||
mstime_t now_ms = mstime();
|
||||
|
||||
if (c->bpop.timeout != 0 && c->bpop.timeout < now_ms) {
|
||||
/* Handle blocking operation specific timeout. */
|
||||
replyToBlockedClientTimedOut(c);
|
||||
unblockClient(c);
|
||||
} else if (server.cluster_enabled) {
|
||||
/* Cluster: handle unblock & redirect of clients blocked
|
||||
* into keys no longer served by this server. */
|
||||
if (clusterRedirectBlockedClientIfNeeded(c))
|
||||
unblockClient(c);
|
||||
}
|
||||
}
|
||||
return 0;
|
||||
@@ -1258,7 +1264,7 @@ void beforeSleep(struct aeEventLoop *eventLoop) {
|
||||
REDIS_NOTUSED(eventLoop);
|
||||
|
||||
/* Call the Redis Cluster before sleep function. Note that this function
|
||||
* may change the state of Redis Cluster (frok ok to fail or vice versa),
|
||||
* may change the state of Redis Cluster (from ok to fail or vice versa),
|
||||
* so it's a good idea to call it before serving the unblocked clients
|
||||
* later in this function. */
|
||||
if (server.cluster_enabled) clusterBeforeSleep();
|
||||
@@ -2158,38 +2164,22 @@ int processCommand(redisClient *c) {
|
||||
* 2) The command has no key arguments. */
|
||||
if (server.cluster_enabled &&
|
||||
!(c->flags & REDIS_MASTER) &&
|
||||
!(c->flags & REDIS_LUA_CLIENT &&
|
||||
server.lua_caller->flags & REDIS_MASTER) &&
|
||||
!(c->cmd->getkeys_proc == NULL && c->cmd->firstkey == 0))
|
||||
{
|
||||
int hashslot;
|
||||
|
||||
if (server.cluster->state != REDIS_CLUSTER_OK) {
|
||||
flagTransaction(c);
|
||||
addReplySds(c,sdsnew("-CLUSTERDOWN The cluster is down. Use CLUSTER INFO for more information\r\n"));
|
||||
clusterRedirectClient(c,NULL,0,REDIS_CLUSTER_REDIR_DOWN_STATE);
|
||||
return REDIS_OK;
|
||||
} else {
|
||||
int error_code;
|
||||
clusterNode *n = getNodeByQuery(c,c->cmd,c->argv,c->argc,&hashslot,&error_code);
|
||||
if (n == NULL) {
|
||||
if (n == NULL || n != server.cluster->myself) {
|
||||
flagTransaction(c);
|
||||
if (error_code == REDIS_CLUSTER_REDIR_CROSS_SLOT) {
|
||||
addReplySds(c,sdsnew("-CROSSSLOT Keys in request don't hash to the same slot\r\n"));
|
||||
} else if (error_code == REDIS_CLUSTER_REDIR_UNSTABLE) {
|
||||
/* The request spawns mutliple keys in the same slot,
|
||||
* but the slot is not "stable" currently as there is
|
||||
* a migration or import in progress. */
|
||||
addReplySds(c,sdsnew("-TRYAGAIN Multiple keys request during rehashing of slot\r\n"));
|
||||
} else if (error_code == REDIS_CLUSTER_REDIR_DOWN) {
|
||||
addReplySds(c,sdsnew("-CLUSTERDOWN The cluster is down. Hash slot is unbound\r\n"));
|
||||
} else {
|
||||
redisPanic("getNodeByQuery() unknown error.");
|
||||
}
|
||||
return REDIS_OK;
|
||||
} else if (n != server.cluster->myself) {
|
||||
flagTransaction(c);
|
||||
addReplySds(c,sdscatprintf(sdsempty(),
|
||||
"-%s %d %s:%d\r\n",
|
||||
(error_code == REDIS_CLUSTER_REDIR_ASK) ? "ASK" : "MOVED",
|
||||
hashslot,n->ip,n->port));
|
||||
clusterRedirectClient(c,n,hashslot,error_code);
|
||||
return REDIS_OK;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1366,6 +1366,7 @@ void blockClient(redisClient *c, int btype);
|
||||
void unblockClient(redisClient *c);
|
||||
void replyToBlockedClientTimedOut(redisClient *c);
|
||||
int getTimeoutFromObjectOrReply(redisClient *c, robj *object, mstime_t *timeout, int unit);
|
||||
void disconnectAllBlockedClients(void);
|
||||
|
||||
/* Git SHA1 */
|
||||
char *redisGitSHA1(void);
|
||||
|
||||
+5
-2
@@ -652,7 +652,8 @@ void replconfCommand(redisClient *c) {
|
||||
*
|
||||
* It does a few things:
|
||||
*
|
||||
* 1) Put the slave in ONLINE state.
|
||||
* 1) Put the slave in ONLINE state (useless when the function is called
|
||||
* because state is already ONLINE but repl_put_online_on_ack is true).
|
||||
* 2) Make sure the writable event is re-installed, since calling the SYNC
|
||||
* command disables it, so that we can accumulate output buffer without
|
||||
* sending it to the slave.
|
||||
@@ -660,7 +661,7 @@ void replconfCommand(redisClient *c) {
|
||||
void putSlaveOnline(redisClient *slave) {
|
||||
slave->replstate = REDIS_REPL_ONLINE;
|
||||
slave->repl_put_online_on_ack = 0;
|
||||
slave->repl_ack_time = server.unixtime;
|
||||
slave->repl_ack_time = server.unixtime; /* Prevent false timeout. */
|
||||
if (aeCreateFileEvent(server.el, slave->fd, AE_WRITABLE,
|
||||
sendReplyToClient, slave) == AE_ERR) {
|
||||
redisLog(REDIS_WARNING,"Unable to register writable event for slave bulk transfer: %s", strerror(errno));
|
||||
@@ -773,6 +774,7 @@ void updateSlavesWaitingBgsave(int bgsaveerr, int type) {
|
||||
* is technically online now. */
|
||||
slave->replstate = REDIS_REPL_ONLINE;
|
||||
slave->repl_put_online_on_ack = 1;
|
||||
slave->repl_ack_time = server.unixtime; /* Timeout otherwise. */
|
||||
} else {
|
||||
if (bgsaveerr != REDIS_OK) {
|
||||
freeClient(slave);
|
||||
@@ -1437,6 +1439,7 @@ void replicationSetMaster(char *ip, int port) {
|
||||
server.masterhost = sdsnew(ip);
|
||||
server.masterport = port;
|
||||
if (server.master) freeClient(server.master);
|
||||
disconnectAllBlockedClients(); /* Clients blocked in master, now slave. */
|
||||
disconnectSlaves(); /* Force our slaves to resync with us as well. */
|
||||
replicationDiscardCachedMaster(); /* Don't try a PSYNC. */
|
||||
freeReplicationBacklog(); /* Don't allow our chained slaves to PSYNC. */
|
||||
|
||||
+3
-2
@@ -357,8 +357,9 @@ int luaRedisGenericCommand(lua_State *lua, int raise_error) {
|
||||
if (cmd->flags & REDIS_CMD_WRITE) server.lua_write_dirty = 1;
|
||||
|
||||
/* If this is a Redis Cluster node, we need to make sure Lua is not
|
||||
* trying to access non-local keys. */
|
||||
if (server.cluster_enabled) {
|
||||
* trying to access non-local keys, with the exception of commands
|
||||
* received from our master. */
|
||||
if (server.cluster_enabled && !(server.lua_caller->flags & REDIS_MASTER)) {
|
||||
/* Duplicate relevant flags in the lua client. */
|
||||
c->flags &= ~(REDIS_READONLY|REDIS_ASKING);
|
||||
c->flags |= server.lua_caller->flags & (REDIS_READONLY|REDIS_ASKING);
|
||||
|
||||
+1
-1
@@ -1 +1 @@
|
||||
#define REDIS_VERSION "2.9.105"
|
||||
#define REDIS_VERSION "3.0.0"
|
||||
|
||||
@@ -17,6 +17,7 @@ proc main {} {
|
||||
}
|
||||
run_tests
|
||||
cleanup
|
||||
end_tests
|
||||
}
|
||||
|
||||
if {[catch main e]} {
|
||||
|
||||
@@ -0,0 +1,192 @@
|
||||
# Check the manual failover
|
||||
|
||||
source "../tests/includes/init-tests.tcl"
|
||||
|
||||
test "Create a 5 nodes cluster" {
|
||||
create_cluster 5 5
|
||||
}
|
||||
|
||||
test "Cluster is up" {
|
||||
assert_cluster_state ok
|
||||
}
|
||||
|
||||
test "Cluster is writable" {
|
||||
cluster_write_test 0
|
||||
}
|
||||
|
||||
test "Instance #5 is a slave" {
|
||||
assert {[RI 5 role] eq {slave}}
|
||||
}
|
||||
|
||||
test "Instance #5 synced with the master" {
|
||||
wait_for_condition 1000 50 {
|
||||
[RI 5 master_link_status] eq {up}
|
||||
} else {
|
||||
fail "Instance #5 master link status is not up"
|
||||
}
|
||||
}
|
||||
|
||||
set current_epoch [CI 1 cluster_current_epoch]
|
||||
|
||||
set numkeys 50000
|
||||
set numops 10000
|
||||
set cluster [redis_cluster 127.0.0.1:[get_instance_attrib redis 0 port]]
|
||||
catch {unset content}
|
||||
array set content {}
|
||||
|
||||
test "Send CLUSTER FAILOVER to #5, during load" {
|
||||
for {set j 0} {$j < $numops} {incr j} {
|
||||
# Write random data to random list.
|
||||
set listid [randomInt $numkeys]
|
||||
set key "key:$listid"
|
||||
set ele [randomValue]
|
||||
# We write both with Lua scripts and with plain commands.
|
||||
# This way we are able to stress Lua -> Redis command invocation
|
||||
# as well, that has tests to prevent Lua to write into wrong
|
||||
# hash slots.
|
||||
if {$listid % 2} {
|
||||
$cluster rpush $key $ele
|
||||
} else {
|
||||
$cluster eval {redis.call("rpush",KEYS[1],ARGV[1])} 1 $key $ele
|
||||
}
|
||||
lappend content($key) $ele
|
||||
|
||||
if {($j % 1000) == 0} {
|
||||
puts -nonewline W; flush stdout
|
||||
}
|
||||
|
||||
if {$j == $numops/2} {R 5 cluster failover}
|
||||
}
|
||||
}
|
||||
|
||||
test "Wait for failover" {
|
||||
wait_for_condition 1000 50 {
|
||||
[CI 1 cluster_current_epoch] > $current_epoch
|
||||
} else {
|
||||
fail "No failover detected"
|
||||
}
|
||||
}
|
||||
|
||||
test "Cluster should eventually be up again" {
|
||||
assert_cluster_state ok
|
||||
}
|
||||
|
||||
test "Cluster is writable" {
|
||||
cluster_write_test 1
|
||||
}
|
||||
|
||||
test "Instance #5 is now a master" {
|
||||
assert {[RI 5 role] eq {master}}
|
||||
}
|
||||
|
||||
test "Verify $numkeys keys for consistency with logical content" {
|
||||
# Check that the Redis Cluster content matches our logical content.
|
||||
foreach {key value} [array get content] {
|
||||
assert {[$cluster lrange $key 0 -1] eq $value}
|
||||
}
|
||||
}
|
||||
|
||||
test "Instance #0 gets converted into a slave" {
|
||||
wait_for_condition 1000 50 {
|
||||
[RI 0 role] eq {slave}
|
||||
} else {
|
||||
fail "Old master was not converted into slave"
|
||||
}
|
||||
}
|
||||
|
||||
## Check that manual failover does not happen if we can't talk with the master.
|
||||
|
||||
source "../tests/includes/init-tests.tcl"
|
||||
|
||||
test "Create a 5 nodes cluster" {
|
||||
create_cluster 5 5
|
||||
}
|
||||
|
||||
test "Cluster is up" {
|
||||
assert_cluster_state ok
|
||||
}
|
||||
|
||||
test "Cluster is writable" {
|
||||
cluster_write_test 0
|
||||
}
|
||||
|
||||
test "Instance #5 is a slave" {
|
||||
assert {[RI 5 role] eq {slave}}
|
||||
}
|
||||
|
||||
test "Instance #5 synced with the master" {
|
||||
wait_for_condition 1000 50 {
|
||||
[RI 5 master_link_status] eq {up}
|
||||
} else {
|
||||
fail "Instance #5 master link status is not up"
|
||||
}
|
||||
}
|
||||
|
||||
test "Make instance #0 unreachable without killing it" {
|
||||
R 0 deferred 1
|
||||
R 0 DEBUG SLEEP 10
|
||||
}
|
||||
|
||||
test "Send CLUSTER FAILOVER to instance #5" {
|
||||
R 5 cluster failover
|
||||
}
|
||||
|
||||
test "Instance #5 is still a slave after some time (no failover)" {
|
||||
after 5000
|
||||
assert {[RI 5 role] eq {master}}
|
||||
}
|
||||
|
||||
test "Wait for instance #0 to return back alive" {
|
||||
R 0 deferred 0
|
||||
assert {[R 0 read] eq {OK}}
|
||||
}
|
||||
|
||||
## Check with "force" failover happens anyway.
|
||||
|
||||
source "../tests/includes/init-tests.tcl"
|
||||
|
||||
test "Create a 5 nodes cluster" {
|
||||
create_cluster 5 5
|
||||
}
|
||||
|
||||
test "Cluster is up" {
|
||||
assert_cluster_state ok
|
||||
}
|
||||
|
||||
test "Cluster is writable" {
|
||||
cluster_write_test 0
|
||||
}
|
||||
|
||||
test "Instance #5 is a slave" {
|
||||
assert {[RI 5 role] eq {slave}}
|
||||
}
|
||||
|
||||
test "Instance #5 synced with the master" {
|
||||
wait_for_condition 1000 50 {
|
||||
[RI 5 master_link_status] eq {up}
|
||||
} else {
|
||||
fail "Instance #5 master link status is not up"
|
||||
}
|
||||
}
|
||||
|
||||
test "Make instance #0 unreachable without killing it" {
|
||||
R 0 deferred 1
|
||||
R 0 DEBUG SLEEP 10
|
||||
}
|
||||
|
||||
test "Send CLUSTER FAILOVER to instance #5" {
|
||||
R 5 cluster failover force
|
||||
}
|
||||
|
||||
test "Instance #5 is a master after some time" {
|
||||
wait_for_condition 1000 50 {
|
||||
[RI 5 role] eq {master}
|
||||
} else {
|
||||
fail "Instance #5 is not a master after some time regardless of FORCE"
|
||||
}
|
||||
}
|
||||
|
||||
test "Wait for instance #0 to return back alive" {
|
||||
R 0 deferred 0
|
||||
assert {[R 0 read] eq {OK}}
|
||||
}
|
||||
@@ -0,0 +1,59 @@
|
||||
# Manual takeover test
|
||||
|
||||
source "../tests/includes/init-tests.tcl"
|
||||
|
||||
test "Create a 5 nodes cluster" {
|
||||
create_cluster 5 5
|
||||
}
|
||||
|
||||
test "Cluster is up" {
|
||||
assert_cluster_state ok
|
||||
}
|
||||
|
||||
test "Cluster is writable" {
|
||||
cluster_write_test 0
|
||||
}
|
||||
|
||||
test "Killing majority of master nodes" {
|
||||
kill_instance redis 0
|
||||
kill_instance redis 1
|
||||
kill_instance redis 2
|
||||
}
|
||||
|
||||
test "Cluster should eventually be down" {
|
||||
assert_cluster_state fail
|
||||
}
|
||||
|
||||
test "Use takeover to bring slaves back" {
|
||||
R 5 cluster failover takeover
|
||||
R 6 cluster failover takeover
|
||||
R 7 cluster failover takeover
|
||||
}
|
||||
|
||||
test "Cluster should eventually be up again" {
|
||||
assert_cluster_state ok
|
||||
}
|
||||
|
||||
test "Cluster is writable" {
|
||||
cluster_write_test 4
|
||||
}
|
||||
|
||||
test "Instance #5, #6, #7 are now masters" {
|
||||
assert {[RI 5 role] eq {master}}
|
||||
assert {[RI 6 role] eq {master}}
|
||||
assert {[RI 7 role] eq {master}}
|
||||
}
|
||||
|
||||
test "Restarting the previously killed master nodes" {
|
||||
restart_instance redis 0
|
||||
restart_instance redis 1
|
||||
restart_instance redis 2
|
||||
}
|
||||
|
||||
test "Instance #0, #1, #2 gets converted into a slaves" {
|
||||
wait_for_condition 1000 50 {
|
||||
[RI 0 role] eq {slave} && [RI 1 role] eq {slave} && [RI 2 role] eq {slave}
|
||||
} else {
|
||||
fail "Old masters not converted into slaves"
|
||||
}
|
||||
}
|
||||
@@ -19,6 +19,7 @@ set ::verbose 0
|
||||
set ::valgrind 0
|
||||
set ::pause_on_error 0
|
||||
set ::simulate_error 0
|
||||
set ::failed 0
|
||||
set ::sentinel_instances {}
|
||||
set ::redis_instances {}
|
||||
set ::sentinel_base_port 20000
|
||||
@@ -231,6 +232,7 @@ proc test {descr code} {
|
||||
flush stdout
|
||||
|
||||
if {[catch {set retval [uplevel 1 $code]} error]} {
|
||||
incr ::failed
|
||||
if {[string match "assertion:*" $error]} {
|
||||
set msg [string range $error 10 end]
|
||||
puts [colorstr red $msg]
|
||||
@@ -246,6 +248,7 @@ proc test {descr code} {
|
||||
}
|
||||
}
|
||||
|
||||
# Execute all the units inside the 'tests' directory.
|
||||
proc run_tests {} {
|
||||
set tests [lsort [glob ../tests/*]]
|
||||
foreach test $tests {
|
||||
@@ -258,6 +261,17 @@ proc run_tests {} {
|
||||
}
|
||||
}
|
||||
|
||||
# Print a message and exists with 0 / 1 according to zero or more failures.
|
||||
proc end_tests {} {
|
||||
if {$::failed == 0} {
|
||||
puts "GOOD! No errors."
|
||||
exit 0
|
||||
} else {
|
||||
puts "WARNING $::failed tests faield."
|
||||
exit 1
|
||||
}
|
||||
}
|
||||
|
||||
# The "S" command is used to interact with the N-th Sentinel.
|
||||
# The general form is:
|
||||
#
|
||||
|
||||
@@ -1,10 +1,17 @@
|
||||
start_server {tags {"repl"}} {
|
||||
set A [srv 0 client]
|
||||
set A_host [srv 0 host]
|
||||
set A_port [srv 0 port]
|
||||
start_server {} {
|
||||
test {First server should have role slave after SLAVEOF} {
|
||||
r -1 slaveof [srv 0 host] [srv 0 port]
|
||||
set B [srv 0 client]
|
||||
set B_host [srv 0 host]
|
||||
set B_port [srv 0 port]
|
||||
|
||||
test {Set instance A as slave of B} {
|
||||
$A slaveof $B_host $B_port
|
||||
wait_for_condition 50 100 {
|
||||
[s -1 role] eq {slave} &&
|
||||
[string match {*master_link_status:up*} [r -1 info replication]]
|
||||
[lindex [$A role] 0] eq {slave} &&
|
||||
[string match {*master_link_status:up*} [$A info replication]]
|
||||
} else {
|
||||
fail "Can't turn the instance into a slave"
|
||||
}
|
||||
@@ -15,9 +22,9 @@ start_server {tags {"repl"}} {
|
||||
$rd brpoplpush a b 5
|
||||
r lpush a foo
|
||||
wait_for_condition 50 100 {
|
||||
[r debug digest] eq [r -1 debug digest]
|
||||
[$A debug digest] eq [$B debug digest]
|
||||
} else {
|
||||
fail "Master and slave have different digest: [r debug digest] VS [r -1 debug digest]"
|
||||
fail "Master and slave have different digest: [$A debug digest] VS [$B debug digest]"
|
||||
}
|
||||
}
|
||||
|
||||
@@ -28,7 +35,36 @@ start_server {tags {"repl"}} {
|
||||
r lpush c 3
|
||||
$rd brpoplpush c d 5
|
||||
after 1000
|
||||
assert_equal [r debug digest] [r -1 debug digest]
|
||||
assert_equal [$A debug digest] [$B debug digest]
|
||||
}
|
||||
|
||||
test {BLPOP followed by role change, issue #2473} {
|
||||
set rd [redis_deferring_client]
|
||||
$rd blpop foo 0 ; # Block while B is a master
|
||||
|
||||
# Turn B into master of A
|
||||
$A slaveof no one
|
||||
$B slaveof $A_host $A_port
|
||||
wait_for_condition 50 100 {
|
||||
[lindex [$B role] 0] eq {slave} &&
|
||||
[string match {*master_link_status:up*} [$B info replication]]
|
||||
} else {
|
||||
fail "Can't turn the instance into a slave"
|
||||
}
|
||||
|
||||
# Push elements into the "foo" list of the new slave.
|
||||
# If the client is still attached to the instance, we'll get
|
||||
# a desync between the two instances.
|
||||
$A rpush foo a b c
|
||||
after 100
|
||||
|
||||
wait_for_condition 50 100 {
|
||||
[$A debug digest] eq [$B debug digest] &&
|
||||
[$A lrange foo 0 -1] eq {a b c} &&
|
||||
[$B lrange foo 0 -1] eq {a b c}
|
||||
} else {
|
||||
fail "Master and slave have different digest: [$A debug digest] VS [$B debug digest]"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -113,7 +149,7 @@ foreach dl {no yes} {
|
||||
start_server {} {
|
||||
lappend slaves [srv 0 client]
|
||||
test "Connect multiple slaves at the same time (issue #141), diskless=$dl" {
|
||||
# Send SALVEOF commands to slaves
|
||||
# Send SLAVEOF commands to slaves
|
||||
[lindex $slaves 0] slaveof $master_host $master_port
|
||||
[lindex $slaves 1] slaveof $master_host $master_port
|
||||
[lindex $slaves 2] slaveof $master_host $master_port
|
||||
|
||||
@@ -13,6 +13,7 @@ proc main {} {
|
||||
spawn_instance redis $::redis_base_port $::instances_count
|
||||
run_tests
|
||||
cleanup
|
||||
end_tests
|
||||
}
|
||||
|
||||
if {[catch main e]} {
|
||||
|
||||
@@ -54,10 +54,15 @@ proc kill_server config {
|
||||
|
||||
# kill server and wait for the process to be totally exited
|
||||
catch {exec kill $pid}
|
||||
if {$::valgrind} {
|
||||
set max_wait 60000
|
||||
} else {
|
||||
set max_wait 10000
|
||||
}
|
||||
while {[is_alive $config]} {
|
||||
incr wait 10
|
||||
|
||||
if {$wait >= 5000} {
|
||||
if {$wait >= $max_wait} {
|
||||
puts "Forcing process $pid to exit..."
|
||||
catch {exec kill -KILL $pid}
|
||||
} elseif {$wait % 1000 == 0} {
|
||||
|
||||
@@ -43,7 +43,7 @@ then
|
||||
while [ $((PORT < ENDPORT)) != "0" ]; do
|
||||
PORT=$((PORT+1))
|
||||
echo "Stopping $PORT"
|
||||
redis-cli -p $PORT shutdown nosave
|
||||
../../src/redis-cli -p $PORT shutdown nosave
|
||||
done
|
||||
exit 0
|
||||
fi
|
||||
@@ -54,7 +54,7 @@ then
|
||||
while [ 1 ]; do
|
||||
clear
|
||||
date
|
||||
redis-cli -p $PORT cluster nodes | head -30
|
||||
../../src/redis-cli -p $PORT cluster nodes | head -30
|
||||
sleep 1
|
||||
done
|
||||
exit 0
|
||||
|
||||
Reference in New Issue
Block a user