Compare commits
10
Commits
conduct
...
no-mo-second
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
352408a8ad | ||
|
|
5418c81312 | ||
|
|
33cbdac124 | ||
|
|
62441004ae | ||
|
|
f135aef0db | ||
|
|
32378b7fad | ||
|
|
ee22a8bdd3 | ||
|
|
0471dcaf91 | ||
|
|
2eb2f3635f | ||
|
|
8dd017368b |
+1
-1
@@ -110,7 +110,7 @@ void processUnblockedClients(void) {
|
||||
* the code is conceptually more correct this way. */
|
||||
if (!(c->flags & CLIENT_BLOCKED)) {
|
||||
if (c->querybuf && sdslen(c->querybuf) > 0) {
|
||||
processInputBuffer(c);
|
||||
processInputBufferAndReplicate(c);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+48
-67
@@ -1037,25 +1037,11 @@ static void freeClientArgv(client *c) {
|
||||
|
||||
/* Close all the slaves connections. This is useful in chained replication
|
||||
* when we resync with our own master and want to force all our slaves to
|
||||
* resync with us as well.
|
||||
*
|
||||
* If 'async' is non-zero we free the clients asynchronously. This is needed
|
||||
* when we call this function from a context where in the chain of the
|
||||
* callers somebody is iterating the list of clients. For instance when
|
||||
* CLIENT KILL TYPE master is called, caching the master client may
|
||||
* adjust the meaningful offset of replication, and in turn call
|
||||
* discionectSlaves(). Since CLIENT KILL iterates the clients this is
|
||||
* not safe. */
|
||||
void disconnectSlaves(int async) {
|
||||
listIter li;
|
||||
listNode *ln;
|
||||
listRewind(server.slaves,&li);
|
||||
while((ln = listNext(&li))) {
|
||||
* resync with us as well. */
|
||||
void disconnectSlaves(void) {
|
||||
while (listLength(server.slaves)) {
|
||||
listNode *ln = listFirst(server.slaves);
|
||||
if (async)
|
||||
freeClientAsync((client*)ln->value);
|
||||
else
|
||||
freeClient((client*)ln->value);
|
||||
freeClient((client*)ln->value);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1763,55 +1749,12 @@ int processMultibulkBuffer(client *c) {
|
||||
return C_ERR;
|
||||
}
|
||||
|
||||
/* Perform necessary tasks after a command was executed:
|
||||
*
|
||||
* 1. The client is reset unless there are reasons to avoid doing it.
|
||||
* 2. In the case of master clients, the replication offset is updated.
|
||||
* 3. Propagate commands we got from our master to replicas down the line. */
|
||||
void commandProcessed(client *c) {
|
||||
int cmd_is_ping = c->cmd && c->cmd->proc == pingCommand;
|
||||
long long prev_offset = c->reploff;
|
||||
if (c->flags & CLIENT_MASTER && !(c->flags & CLIENT_MULTI)) {
|
||||
/* Update the applied replication offset of our master. */
|
||||
c->reploff = c->read_reploff - sdslen(c->querybuf) + c->qb_pos;
|
||||
}
|
||||
|
||||
/* Don't reset the client structure for clients blocked in a
|
||||
* module blocking command, so that the reply callback will
|
||||
* still be able to access the client argv and argc field.
|
||||
* The client will be reset in unblockClientFromModule(). */
|
||||
if (!(c->flags & CLIENT_BLOCKED) ||
|
||||
c->btype != BLOCKED_MODULE)
|
||||
{
|
||||
resetClient(c);
|
||||
}
|
||||
|
||||
/* If the client is a master we need to compute the difference
|
||||
* between the applied offset before and after processing the buffer,
|
||||
* to understand how much of the replication stream was actually
|
||||
* applied to the master state: this quantity, and its corresponding
|
||||
* part of the replication stream, will be propagated to the
|
||||
* sub-replicas and to the replication backlog. */
|
||||
if (c->flags & CLIENT_MASTER) {
|
||||
long long applied = c->reploff - prev_offset;
|
||||
long long prev_master_repl_meaningful_offset = server.master_repl_meaningful_offset;
|
||||
if (applied) {
|
||||
replicationFeedSlavesFromMasterStream(server.slaves,
|
||||
c->pending_querybuf, applied);
|
||||
sdsrange(c->pending_querybuf,applied,-1);
|
||||
}
|
||||
/* The server.master_repl_meaningful_offset variable represents
|
||||
* the offset of the replication stream without the pending PINGs. */
|
||||
if (cmd_is_ping)
|
||||
server.master_repl_meaningful_offset = prev_master_repl_meaningful_offset;
|
||||
}
|
||||
}
|
||||
|
||||
/* This function calls processCommand(), but also performs a few sub tasks
|
||||
* for the client that are useful in that context:
|
||||
* that are useful in that context:
|
||||
*
|
||||
* 1. It sets the current client to the client 'c'.
|
||||
* 2. calls commandProcessed() if the command was handled.
|
||||
* 2. In the case of master clients, the replication offset is updated.
|
||||
* 3. The client is reset unless there are reasons to avoid doing it.
|
||||
*
|
||||
* The function returns C_ERR in case the client was freed as a side effect
|
||||
* of processing the command, otherwise C_OK is returned. */
|
||||
@@ -1819,7 +1762,20 @@ int processCommandAndResetClient(client *c) {
|
||||
int deadclient = 0;
|
||||
server.current_client = c;
|
||||
if (processCommand(c) == C_OK) {
|
||||
commandProcessed(c);
|
||||
if (c->flags & CLIENT_MASTER && !(c->flags & CLIENT_MULTI)) {
|
||||
/* Update the applied replication offset of our master. */
|
||||
c->reploff = c->read_reploff - sdslen(c->querybuf) + c->qb_pos;
|
||||
}
|
||||
|
||||
/* Don't reset the client structure for clients blocked in a
|
||||
* module blocking command, so that the reply callback will
|
||||
* still be able to access the client argv and argc field.
|
||||
* The client will be reset in unblockClientFromModule(). */
|
||||
if (!(c->flags & CLIENT_BLOCKED) ||
|
||||
c->btype != BLOCKED_MODULE)
|
||||
{
|
||||
resetClient(c);
|
||||
}
|
||||
}
|
||||
if (server.current_client == NULL) deadclient = 1;
|
||||
server.current_client = NULL;
|
||||
@@ -1916,6 +1872,31 @@ void processInputBuffer(client *c) {
|
||||
}
|
||||
}
|
||||
|
||||
/* This is a wrapper for processInputBuffer that also cares about handling
|
||||
* the replication forwarding to the sub-replicas, in case the client 'c'
|
||||
* is flagged as master. Usually you want to call this instead of the
|
||||
* raw processInputBuffer(). */
|
||||
void processInputBufferAndReplicate(client *c) {
|
||||
if (!(c->flags & CLIENT_MASTER)) {
|
||||
processInputBuffer(c);
|
||||
} else {
|
||||
/* If the client is a master we need to compute the difference
|
||||
* between the applied offset before and after processing the buffer,
|
||||
* to understand how much of the replication stream was actually
|
||||
* applied to the master state: this quantity, and its corresponding
|
||||
* part of the replication stream, will be propagated to the
|
||||
* sub-replicas and to the replication backlog. */
|
||||
size_t prev_offset = c->reploff;
|
||||
processInputBuffer(c);
|
||||
size_t applied = c->reploff - prev_offset;
|
||||
if (applied) {
|
||||
replicationFeedSlavesFromMasterStream(server.slaves,
|
||||
c->pending_querybuf, applied);
|
||||
sdsrange(c->pending_querybuf,applied,-1);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void readQueryFromClient(connection *conn) {
|
||||
client *c = connGetPrivateData(conn);
|
||||
int nread, readlen;
|
||||
@@ -1983,7 +1964,7 @@ void readQueryFromClient(connection *conn) {
|
||||
|
||||
/* There is more data in the client input buffer, continue parsing it
|
||||
* in case to check if there is a full command to execute. */
|
||||
processInputBuffer(c);
|
||||
processInputBufferAndReplicate(c);
|
||||
}
|
||||
|
||||
void getClientsMaxBuffers(unsigned long *longest_output_list,
|
||||
@@ -3212,7 +3193,7 @@ int handleClientsWithPendingReadsUsingThreads(void) {
|
||||
continue;
|
||||
}
|
||||
}
|
||||
processInputBuffer(c);
|
||||
processInputBufferAndReplicate(c);
|
||||
}
|
||||
return processed;
|
||||
}
|
||||
|
||||
+10
-81
@@ -39,7 +39,6 @@
|
||||
#include <sys/socket.h>
|
||||
#include <sys/stat.h>
|
||||
|
||||
long long adjustMeaningfulReplOffset(int *adjusted);
|
||||
void replicationDiscardCachedMaster(void);
|
||||
void replicationResurrectCachedMaster(connection *conn);
|
||||
void replicationSendAck(void);
|
||||
@@ -163,7 +162,6 @@ void feedReplicationBacklog(void *ptr, size_t len) {
|
||||
unsigned char *p = ptr;
|
||||
|
||||
server.master_repl_offset += len;
|
||||
server.master_repl_meaningful_offset = server.master_repl_offset;
|
||||
|
||||
/* This is a circular buffer, so write as much data we can at every
|
||||
* iteration and rewind the "idx" index if we reach the limit. */
|
||||
@@ -1831,7 +1829,6 @@ void readSyncBulkPayload(connection *conn) {
|
||||
* we are starting a new history. */
|
||||
memcpy(server.replid,server.master->replid,sizeof(server.replid));
|
||||
server.master_repl_offset = server.master->reploff;
|
||||
server.master_repl_meaningful_offset = server.master->reploff;
|
||||
clearReplicationId2();
|
||||
|
||||
/* Let's create the replication backlog if needed. Slaves need to
|
||||
@@ -2086,7 +2083,7 @@ int slaveTryPartialResynchronization(connection *conn, int read_reply) {
|
||||
memcpy(server.cached_master->replid,new,sizeof(server.replid));
|
||||
|
||||
/* Disconnect all the sub-slaves: they need to be notified. */
|
||||
disconnectSlaves(0);
|
||||
disconnectSlaves();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2359,7 +2356,7 @@ void syncWithMaster(connection *conn) {
|
||||
* as well, if we have any sub-slaves. The master may transfer us an
|
||||
* entirely different data set and we have no way to incrementally feed
|
||||
* our slaves after that. */
|
||||
disconnectSlaves(0); /* Force our slaves to resync with us as well. */
|
||||
disconnectSlaves(); /* Force our slaves to resync with us as well. */
|
||||
freeReplicationBacklog(); /* Don't allow our chained slaves to PSYNC. */
|
||||
|
||||
/* Fall back to SYNC if needed. Otherwise psync_result == PSYNC_FULLRESYNC
|
||||
@@ -2506,7 +2503,7 @@ void replicationSetMaster(char *ip, int port) {
|
||||
|
||||
/* Force our slaves to resync with us as well. They may hopefully be able
|
||||
* to partially resync with us, but we can notify the replid change. */
|
||||
disconnectSlaves(0);
|
||||
disconnectSlaves();
|
||||
cancelReplicationHandshake();
|
||||
/* Before destroying our master state, create a cached master using
|
||||
* our own parameters, to later PSYNC with the new master. */
|
||||
@@ -2541,23 +2538,19 @@ void replicationUnsetMaster(void) {
|
||||
|
||||
sdsfree(server.masterhost);
|
||||
server.masterhost = NULL;
|
||||
if (server.master) freeClient(server.master);
|
||||
replicationDiscardCachedMaster();
|
||||
cancelReplicationHandshake();
|
||||
/* When a slave is turned into a master, the current replication ID
|
||||
* (that was inherited from the master at synchronization time) is
|
||||
* used as secondary ID up to the current offset, and a new replication
|
||||
* ID is created to continue with a new replication history.
|
||||
*
|
||||
* NOTE: this function MUST be called after we call
|
||||
* freeClient(server.master), since there we adjust the replication
|
||||
* offset trimming the final PINGs. See Github issue #7320. */
|
||||
* ID is created to continue with a new replication history. */
|
||||
shiftReplicationId();
|
||||
if (server.master) freeClient(server.master);
|
||||
replicationDiscardCachedMaster();
|
||||
cancelReplicationHandshake();
|
||||
/* Disconnecting all the slaves is required: we need to inform slaves
|
||||
* of the replication ID change (see shiftReplicationId() call). However
|
||||
* the slaves will be able to partially resync with us, so it will be
|
||||
* a very fast reconnection. */
|
||||
disconnectSlaves(0);
|
||||
disconnectSlaves();
|
||||
server.repl_state = REPL_STATE_NONE;
|
||||
|
||||
/* We need to make sure the new master will start the replication stream
|
||||
@@ -2759,11 +2752,6 @@ void replicationCacheMaster(client *c) {
|
||||
* pending outputs to the master. */
|
||||
sdsclear(server.master->querybuf);
|
||||
sdsclear(server.master->pending_querybuf);
|
||||
|
||||
/* Adjust reploff and read_reploff to the last meaningful offset we
|
||||
* executed. This is the offset the replica will use for future PSYNC. */
|
||||
int offset_adjusted;
|
||||
server.master->reploff = adjustMeaningfulReplOffset(&offset_adjusted);
|
||||
server.master->read_reploff = server.master->reploff;
|
||||
if (c->flags & CLIENT_MULTI) discardTransaction(c);
|
||||
listEmpty(c->reply);
|
||||
@@ -2786,53 +2774,6 @@ void replicationCacheMaster(client *c) {
|
||||
* so make sure to adjust the replication state. This function will
|
||||
* also set server.master to NULL. */
|
||||
replicationHandleMasterDisconnection();
|
||||
|
||||
/* If we trimmed this replica backlog, we need to disconnect our chained
|
||||
* replicas (if any), otherwise they may have the PINGs we removed
|
||||
* from the stream and their offset would no longer match: upon
|
||||
* disconnection they will also trim the final PINGs and will be able
|
||||
* to incrementally sync without issues. */
|
||||
if (offset_adjusted) disconnectSlaves(1);
|
||||
}
|
||||
|
||||
/* If the "meaningful" offset, that is the offset without the final PINGs
|
||||
* in the stream, is different than the last offset, use it instead:
|
||||
* often when the master is no longer reachable, replicas will never
|
||||
* receive the PINGs, however the master will end with an incremented
|
||||
* offset because of the PINGs and will not be able to incrementally
|
||||
* PSYNC with the new master.
|
||||
* This function trims the replication backlog when needed, and returns
|
||||
* the offset to be used for future partial sync.
|
||||
*
|
||||
* If the integer 'adjusted' was passed by reference, it is set to 1
|
||||
* if the function call actually modified the offset and the replication
|
||||
* backlog, otherwise it is set to 0. It can be NULL if the caller is
|
||||
* not interested in getting this info. */
|
||||
long long adjustMeaningfulReplOffset(int *adjusted) {
|
||||
if (server.master_repl_offset > server.master_repl_meaningful_offset) {
|
||||
long long delta = server.master_repl_offset -
|
||||
server.master_repl_meaningful_offset;
|
||||
serverLog(LL_NOTICE,
|
||||
"Using the meaningful offset %lld instead of %lld to exclude "
|
||||
"the final PINGs (%lld bytes difference)",
|
||||
server.master_repl_meaningful_offset,
|
||||
server.master_repl_offset,
|
||||
delta);
|
||||
server.master_repl_offset = server.master_repl_meaningful_offset;
|
||||
if (server.repl_backlog_histlen <= delta) {
|
||||
server.repl_backlog_histlen = 0;
|
||||
server.repl_backlog_idx = 0;
|
||||
} else {
|
||||
server.repl_backlog_histlen -= delta;
|
||||
server.repl_backlog_idx =
|
||||
(server.repl_backlog_idx + (server.repl_backlog_size - delta)) %
|
||||
server.repl_backlog_size;
|
||||
}
|
||||
if (adjusted) *adjusted = 1;
|
||||
} else {
|
||||
if (adjusted) *adjusted = 0;
|
||||
}
|
||||
return server.master_repl_offset;
|
||||
}
|
||||
|
||||
/* This function is called when a master is turend into a slave, in order to
|
||||
@@ -2845,16 +2786,11 @@ long long adjustMeaningfulReplOffset(int *adjusted) {
|
||||
* current offset if no data was lost during the failover. So we use our
|
||||
* current replication ID and offset in order to synthesize a cached master. */
|
||||
void replicationCacheMasterUsingMyself(void) {
|
||||
serverLog(LL_NOTICE,
|
||||
"Before turning into a replica, using my own master parameters "
|
||||
"to synthesize a cached master: I may be able to synchronize with "
|
||||
"the new master with just a partial transfer.");
|
||||
|
||||
/* This will be used to populate the field server.master->reploff
|
||||
* by replicationCreateMasterClient(). We'll later set the created
|
||||
* master as server.cached_master, so the replica will use such
|
||||
* offset for PSYNC. */
|
||||
server.master_initial_offset = adjustMeaningfulReplOffset(NULL);
|
||||
server.master_initial_offset = server.master_repl_offset;
|
||||
|
||||
/* The master client we create can be set to any DBID, because
|
||||
* the new master will start its replication stream with SELECT. */
|
||||
@@ -2867,6 +2803,7 @@ void replicationCacheMasterUsingMyself(void) {
|
||||
unlinkClient(server.master);
|
||||
server.cached_master = server.master;
|
||||
server.master = NULL;
|
||||
serverLog(LL_NOTICE,"Before turning into a replica, using my master parameters to synthesize a cached master: I may be able to synchronize with the new master with just a partial transfer.");
|
||||
}
|
||||
|
||||
/* Free a cached master, called when there are no longer the conditions for
|
||||
@@ -3246,18 +3183,10 @@ void replicationCron(void) {
|
||||
clientsArePaused();
|
||||
|
||||
if (!manual_failover_in_progress) {
|
||||
long long before_ping = server.master_repl_meaningful_offset;
|
||||
ping_argv[0] = createStringObject("PING",4);
|
||||
replicationFeedSlaves(server.slaves, server.slaveseldb,
|
||||
ping_argv, 1);
|
||||
decrRefCount(ping_argv[0]);
|
||||
/* The server.master_repl_meaningful_offset variable represents
|
||||
* the offset of the replication stream without the pending PINGs.
|
||||
* This is useful to set the right replication offset for PSYNC
|
||||
* when the master is turned into a replica. Otherwise pending
|
||||
* PINGs may not allow it to perform an incremental sync with the
|
||||
* new master. */
|
||||
server.master_repl_meaningful_offset = before_ping;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -2394,7 +2394,6 @@ void initServerConfig(void) {
|
||||
server.repl_syncio_timeout = CONFIG_REPL_SYNCIO_TIMEOUT;
|
||||
server.repl_down_since = 0; /* Never connected, repl is down since EVER. */
|
||||
server.master_repl_offset = 0;
|
||||
server.master_repl_meaningful_offset = 0;
|
||||
|
||||
/* Replication partial resync backlog */
|
||||
server.repl_backlog = NULL;
|
||||
@@ -4471,7 +4470,6 @@ sds genRedisInfoString(const char *section) {
|
||||
"master_replid:%s\r\n"
|
||||
"master_replid2:%s\r\n"
|
||||
"master_repl_offset:%lld\r\n"
|
||||
"master_repl_meaningful_offset:%lld\r\n"
|
||||
"second_repl_offset:%lld\r\n"
|
||||
"repl_backlog_active:%d\r\n"
|
||||
"repl_backlog_size:%lld\r\n"
|
||||
@@ -4480,7 +4478,6 @@ sds genRedisInfoString(const char *section) {
|
||||
server.replid,
|
||||
server.replid2,
|
||||
server.master_repl_offset,
|
||||
server.master_repl_meaningful_offset,
|
||||
server.second_replid_offset,
|
||||
server.repl_backlog != NULL,
|
||||
server.repl_backlog_size,
|
||||
@@ -4858,7 +4855,6 @@ void loadDataFromDisk(void) {
|
||||
{
|
||||
memcpy(server.replid,rsi.repl_id,sizeof(server.replid));
|
||||
server.master_repl_offset = rsi.repl_offset;
|
||||
server.master_repl_meaningful_offset = rsi.repl_offset;
|
||||
/* If we are a slave, create a cached master from this
|
||||
* information, in order to allow partial resynchronizations
|
||||
* with masters. */
|
||||
|
||||
+2
-2
@@ -1261,7 +1261,6 @@ struct redisServer {
|
||||
char replid[CONFIG_RUN_ID_SIZE+1]; /* My current replication ID. */
|
||||
char replid2[CONFIG_RUN_ID_SIZE+1]; /* replid inherited from master*/
|
||||
long long master_repl_offset; /* My current replication offset */
|
||||
long long master_repl_meaningful_offset; /* Offset minus latest PINGs. */
|
||||
long long second_replid_offset; /* Accept offsets up to this for replid2. */
|
||||
int slaveseldb; /* Last SELECTed DB in replication output */
|
||||
int repl_ping_slave_period; /* Master pings the slave every N seconds */
|
||||
@@ -1608,6 +1607,7 @@ void setDeferredSetLen(client *c, void *node, long length);
|
||||
void setDeferredAttributeLen(client *c, void *node, long length);
|
||||
void setDeferredPushLen(client *c, void *node, long length);
|
||||
void processInputBuffer(client *c);
|
||||
void processInputBufferAndReplicate(client *c);
|
||||
void processGopherRequest(client *c);
|
||||
void acceptHandler(aeEventLoop *el, int fd, void *privdata, int mask);
|
||||
void acceptTcpHandler(aeEventLoop *el, int fd, void *privdata, int mask);
|
||||
@@ -1660,7 +1660,7 @@ int getClientType(client *c);
|
||||
int getClientTypeByName(char *name);
|
||||
char *getClientTypeName(int class);
|
||||
void flushSlavesOutputBuffers(void);
|
||||
void disconnectSlaves(int async);
|
||||
void disconnectSlaves(void);
|
||||
int listenToPort(int port, int *fds, int *count);
|
||||
void pauseClients(mstime_t duration);
|
||||
int clientsArePaused(void);
|
||||
|
||||
@@ -1,112 +0,0 @@
|
||||
# Test the meaningful offset implementation to make sure masters
|
||||
# are able to PSYNC with replicas even if the replication stream
|
||||
# has pending PINGs at the end.
|
||||
|
||||
start_server {tags {"psync2"}} {
|
||||
start_server {} {
|
||||
# Config
|
||||
set debug_msg 0 ; # Enable additional debug messages
|
||||
|
||||
for {set j 0} {$j < 2} {incr j} {
|
||||
set R($j) [srv [expr 0-$j] client]
|
||||
set R_host($j) [srv [expr 0-$j] host]
|
||||
set R_port($j) [srv [expr 0-$j] port]
|
||||
$R($j) CONFIG SET repl-ping-replica-period 1
|
||||
if {$debug_msg} {puts "Log file: [srv [expr 0-$j] stdout]"}
|
||||
}
|
||||
|
||||
# Setup replication
|
||||
test "PSYNC2 meaningful offset: setup" {
|
||||
$R(1) replicaof $R_host(0) $R_port(0)
|
||||
$R(0) set foo bar
|
||||
wait_for_condition 50 1000 {
|
||||
[status $R(1) master_link_status] == "up" &&
|
||||
[$R(0) dbsize] == 1 && [$R(1) dbsize] == 1
|
||||
} else {
|
||||
fail "Replicas not replicating from master"
|
||||
}
|
||||
}
|
||||
|
||||
test "PSYNC2 meaningful offset: write and wait replication" {
|
||||
$R(0) INCR counter
|
||||
$R(0) INCR counter
|
||||
$R(0) INCR counter
|
||||
wait_for_condition 50 1000 {
|
||||
[$R(0) GET counter] eq [$R(1) GET counter]
|
||||
} else {
|
||||
fail "Master and replica don't agree about counter"
|
||||
}
|
||||
}
|
||||
|
||||
# In this test we'll make sure the replica will get stuck, but with
|
||||
# an active connection: this way the master will continue to send PINGs
|
||||
# every second (we modified the PING period earlier)
|
||||
test "PSYNC2 meaningful offset: pause replica and promote it" {
|
||||
$R(1) MULTI
|
||||
$R(1) DEBUG SLEEP 5
|
||||
$R(1) SLAVEOF NO ONE
|
||||
$R(1) EXEC
|
||||
$R(1) ping ; # Wait for it to return back available
|
||||
}
|
||||
|
||||
test "Make the old master a replica of the new one and check conditions" {
|
||||
set sync_partial [status $R(1) sync_partial_ok]
|
||||
assert {$sync_partial == 0}
|
||||
$R(0) REPLICAOF $R_host(1) $R_port(1)
|
||||
wait_for_condition 50 1000 {
|
||||
[status $R(1) sync_partial_ok] == 1
|
||||
} else {
|
||||
fail "The new master was not able to partial sync"
|
||||
}
|
||||
}
|
||||
}}
|
||||
|
||||
|
||||
start_server {tags {"psync2"}} {
|
||||
start_server {} {
|
||||
start_server {} {
|
||||
|
||||
for {set j 0} {$j < 3} {incr j} {
|
||||
set R($j) [srv [expr 0-$j] client]
|
||||
set R_host($j) [srv [expr 0-$j] host]
|
||||
set R_port($j) [srv [expr 0-$j] port]
|
||||
$R($j) CONFIG SET repl-ping-replica-period 1
|
||||
}
|
||||
|
||||
test "Chained replicas disconnect when replica re-connect with the same master" {
|
||||
# Add a second replica as a chained replica of the current replica
|
||||
$R(1) replicaof $R_host(0) $R_port(0)
|
||||
$R(2) replicaof $R_host(1) $R_port(1)
|
||||
wait_for_condition 50 1000 {
|
||||
[status $R(2) master_link_status] == "up"
|
||||
} else {
|
||||
fail "Chained replica not replicating from its master"
|
||||
}
|
||||
|
||||
# Do a write on the master, and wait for 3 seconds for the master to
|
||||
# send some PINGs to its replica
|
||||
$R(0) INCR counter2
|
||||
after 2000
|
||||
set sync_partial_master [status $R(0) sync_partial_ok]
|
||||
set sync_partial_replica [status $R(1) sync_partial_ok]
|
||||
$R(0) CONFIG SET repl-ping-replica-period 100
|
||||
|
||||
# Disconnect the master's direct replica
|
||||
$R(0) client kill type replica
|
||||
wait_for_condition 50 1000 {
|
||||
[status $R(1) master_link_status] == "up" &&
|
||||
[status $R(2) master_link_status] == "up" &&
|
||||
[status $R(0) sync_partial_ok] == $sync_partial_master + 1 &&
|
||||
[status $R(1) sync_partial_ok] == $sync_partial_replica + 1
|
||||
} else {
|
||||
fail "Disconnected replica failed to PSYNC with master"
|
||||
}
|
||||
|
||||
# Verify that the replica and its replica's meaningful and real
|
||||
# offsets match with the master
|
||||
assert_equal [status $R(0) master_repl_offset] [status $R(1) master_repl_offset]
|
||||
assert_equal [status $R(0) master_repl_offset] [status $R(2) master_repl_offset]
|
||||
assert_equal [status $R(0) master_repl_meaningful_offset] [status $R(1) master_repl_meaningful_offset]
|
||||
assert_equal [status $R(0) master_repl_meaningful_offset] [status $R(2) master_repl_meaningful_offset]
|
||||
}
|
||||
}}}
|
||||
@@ -117,7 +117,6 @@ start_server {} {
|
||||
set used [list $master_id]
|
||||
test "PSYNC2: \[NEW LAYOUT\] Set #$master_id as master" {
|
||||
$R($master_id) slaveof no one
|
||||
$R($master_id) config set repl-ping-replica-period 1 ;# increse the chance that random ping will cause issues
|
||||
if {$counter_value == 0} {
|
||||
$R($master_id) set x $counter_value
|
||||
}
|
||||
@@ -259,14 +258,9 @@ start_server {} {
|
||||
$R($j) slaveof $master_host $master_port
|
||||
}
|
||||
|
||||
# Wait for replicas to sync. it is not enough to just wait for connected_slaves==4
|
||||
# since we might do the check before the master realized that they're disconnected
|
||||
# Wait for slaves to sync
|
||||
wait_for_condition 50 1000 {
|
||||
[status $R($master_id) connected_slaves] == 4 &&
|
||||
[status $R([expr {($master_id+1)%5}]) master_link_status] == "up" &&
|
||||
[status $R([expr {($master_id+2)%5}]) master_link_status] == "up" &&
|
||||
[status $R([expr {($master_id+3)%5}]) master_link_status] == "up" &&
|
||||
[status $R([expr {($master_id+4)%5}]) master_link_status] == "up"
|
||||
[status $R($master_id) connected_slaves] == 4
|
||||
} else {
|
||||
show_cluster_status
|
||||
fail "Replica not reconnecting"
|
||||
@@ -278,7 +272,6 @@ start_server {} {
|
||||
set slave_id [expr {($master_id+1)%5}]
|
||||
set sync_count [status $R($master_id) sync_full]
|
||||
set sync_partial [status $R($master_id) sync_partial_ok]
|
||||
set sync_partial_err [status $R($master_id) sync_partial_err]
|
||||
catch {
|
||||
$R($slave_id) config rewrite
|
||||
$R($slave_id) debug restart
|
||||
@@ -370,103 +363,3 @@ start_server {} {
|
||||
}
|
||||
|
||||
}}}}}
|
||||
|
||||
start_server {tags {"psync2"}} {
|
||||
start_server {} {
|
||||
start_server {} {
|
||||
start_server {} {
|
||||
start_server {} {
|
||||
test {pings at the end of replication stream are ignored for psync} {
|
||||
set master [srv -4 client]
|
||||
set master_host [srv -4 host]
|
||||
set master_port [srv -4 port]
|
||||
set replica1 [srv -3 client]
|
||||
set replica2 [srv -2 client]
|
||||
set replica3 [srv -1 client]
|
||||
set replica4 [srv -0 client]
|
||||
|
||||
$replica1 replicaof $master_host $master_port
|
||||
$replica2 replicaof $master_host $master_port
|
||||
$replica3 replicaof $master_host $master_port
|
||||
$replica4 replicaof $master_host $master_port
|
||||
wait_for_condition 50 1000 {
|
||||
[status $master connected_slaves] == 4
|
||||
} else {
|
||||
fail "replicas didn't connect"
|
||||
}
|
||||
|
||||
$master incr x
|
||||
wait_for_condition 50 1000 {
|
||||
[$replica1 get x] == 1 && [$replica2 get x] == 1 &&
|
||||
[$replica3 get x] == 1 && [$replica4 get x] == 1
|
||||
} else {
|
||||
fail "replicas didn't get incr"
|
||||
}
|
||||
|
||||
# disconnect replica1 and replica2
|
||||
# and wait for the master to send a ping to replica3 and replica4
|
||||
$replica1 replicaof no one
|
||||
$replica2 replicaof 127.0.0.1 1 ;# we can't promote it to master since that will cycle the replication id
|
||||
$master config set repl-ping-replica-period 1
|
||||
after 1500
|
||||
|
||||
# make everyone sync from the replica1 that didn't get the last ping from the old master
|
||||
# replica4 will keep syncing from the old master which now syncs from replica1
|
||||
# and replica2 will re-connect to the old master (which went back in time)
|
||||
set new_master_host [srv -3 host]
|
||||
set new_master_port [srv -3 port]
|
||||
$replica3 replicaof $new_master_host $new_master_port
|
||||
$master replicaof $new_master_host $new_master_port
|
||||
$replica2 replicaof $master_host $master_port
|
||||
wait_for_condition 50 1000 {
|
||||
[status $replica2 master_link_status] == "up" &&
|
||||
[status $replica3 master_link_status] == "up" &&
|
||||
[status $replica4 master_link_status] == "up" &&
|
||||
[status $master master_link_status] == "up"
|
||||
} else {
|
||||
fail "replicas didn't connect"
|
||||
}
|
||||
|
||||
# make sure replication is still alive and kicking
|
||||
$replica1 incr x
|
||||
wait_for_condition 50 1000 {
|
||||
[$replica2 get x] == 2 &&
|
||||
[$replica3 get x] == 2 &&
|
||||
[$replica4 get x] == 2 &&
|
||||
[$master get x] == 2
|
||||
} else {
|
||||
fail "replicas didn't get incr"
|
||||
}
|
||||
|
||||
# make sure there are full syncs other than the initial ones
|
||||
assert_equal [status $master sync_full] 4
|
||||
assert_equal [status $replica1 sync_full] 0
|
||||
assert_equal [status $replica2 sync_full] 0
|
||||
assert_equal [status $replica3 sync_full] 0
|
||||
assert_equal [status $replica4 sync_full] 0
|
||||
|
||||
# force psync
|
||||
$master client kill type master
|
||||
$replica2 client kill type master
|
||||
$replica3 client kill type master
|
||||
$replica4 client kill type master
|
||||
|
||||
# make sure replication is still alive and kicking
|
||||
$replica1 incr x
|
||||
wait_for_condition 50 1000 {
|
||||
[$replica2 get x] == 3 &&
|
||||
[$replica3 get x] == 3 &&
|
||||
[$replica4 get x] == 3 &&
|
||||
[$master get x] == 3
|
||||
} else {
|
||||
fail "replicas didn't get incr"
|
||||
}
|
||||
|
||||
# make sure there are full syncs other than the initial ones
|
||||
assert_equal [status $master sync_full] 4
|
||||
assert_equal [status $replica1 sync_full] 0
|
||||
assert_equal [status $replica2 sync_full] 0
|
||||
assert_equal [status $replica3 sync_full] 0
|
||||
assert_equal [status $replica4 sync_full] 0
|
||||
}
|
||||
}}}}}
|
||||
|
||||
@@ -47,7 +47,6 @@ set ::all_tests {
|
||||
integration/logging
|
||||
integration/psync2
|
||||
integration/psync2-reg
|
||||
integration/psync2-pingoff
|
||||
unit/pubsub
|
||||
unit/slowlog
|
||||
unit/scripting
|
||||
|
||||
Reference in New Issue
Block a user