Compare commits
5
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
81428a24a8 | ||
|
|
eb706b4202 | ||
|
|
9e67df2a39 | ||
|
|
f472bb1069 | ||
|
|
fe63046e2a |
+30
@@ -318,6 +318,36 @@ repl-diskless-sync no
|
||||
# it entirely just set it to 0 seconds and the transfer will start ASAP.
|
||||
repl-diskless-sync-delay 5
|
||||
|
||||
# Enable diskless replication on slave side.
|
||||
#
|
||||
# When this option is on, the slave loads the RDB directly from the socket
|
||||
# rather than saving it to disk first. However there are data loss risks
|
||||
# associated with this feature, so make sure to read the following WARNING
|
||||
# section.
|
||||
#
|
||||
# WARNING: Note that this means that the dataset in the slave gets flushed
|
||||
# before the slave is actually sure the RDB transfer is complete, so if the
|
||||
# replication link is disconnected after the slave already flushed away its
|
||||
# dataset, but before successfully loading the new one, the slave will
|
||||
# remain empty (for all the time needed to attempt a new synchornization with
|
||||
# the master).
|
||||
#
|
||||
# This means that you should carefully consider the effects of this feature
|
||||
# on slaves that may be promoted to masters:
|
||||
#
|
||||
# 1) Sentinel checks the disconnection time and the offset of slaves before
|
||||
# promotion. However it is possible that after the check, the slave
|
||||
# attempts to connect with the master again and flushes its dataset.
|
||||
# In order to run Sentinel safely in this setup, make sure to enable
|
||||
# the "slave-protected-restart" option.
|
||||
#
|
||||
# 2) Redis Cluster slaves will refuse to try to be promoted to masters if
|
||||
# if the dataset was flushed, so this is safe in the context of Redis Cluster.
|
||||
#
|
||||
# 3) If you are using your own HA setup, make sure to enable slave
|
||||
# "slave-protected-restart".
|
||||
repl-diskless-load no
|
||||
|
||||
# Slaves send PINGs to server in a predefined interval. It's possible to change
|
||||
# this interval with the repl_ping_slave_period option. The default value is 10
|
||||
# seconds.
|
||||
|
||||
+14
@@ -193,6 +193,20 @@ int anetSendTimeout(char *err, int fd, long long ms) {
|
||||
return ANET_OK;
|
||||
}
|
||||
|
||||
/* Set the socket receive timeout (SO_RCVTIMEO socket option) to the specified
|
||||
* number of milliseconds, or disable it if the 'ms' argument is zero. */
|
||||
int anetRecvTimeout(char *err, int fd, long long ms) {
|
||||
struct timeval tv;
|
||||
|
||||
tv.tv_sec = ms/1000;
|
||||
tv.tv_usec = (ms%1000)*1000;
|
||||
if (setsockopt(fd, SOL_SOCKET, SO_RCVTIMEO, &tv, sizeof(tv)) == -1) {
|
||||
anetSetError(err, "setsockopt SO_RCVTIMEO: %s", strerror(errno));
|
||||
return ANET_ERR;
|
||||
}
|
||||
return ANET_OK;
|
||||
}
|
||||
|
||||
/* anetGenericResolve() is called by anetResolve() and anetResolveIP() to
|
||||
* do the actual work. It resolves the hostname "host" and set the string
|
||||
* representation of the IP address into the buffer pointed by "ipbuf".
|
||||
|
||||
@@ -68,6 +68,7 @@ int anetEnableTcpNoDelay(char *err, int fd);
|
||||
int anetDisableTcpNoDelay(char *err, int fd);
|
||||
int anetTcpKeepAlive(char *err, int fd);
|
||||
int anetSendTimeout(char *err, int fd, long long ms);
|
||||
int anetRecvTimeout(char *err, int fd, long long ms);
|
||||
int anetPeerToString(int fd, char *ip, size_t ip_len, int *port);
|
||||
int anetKeepAlive(char *err, int fd, int interval);
|
||||
int anetSockName(int fd, char *ip, size_t ip_len, int *port);
|
||||
|
||||
@@ -620,7 +620,7 @@ int loadAppendOnlyFile(char *filename) {
|
||||
server.aof_state = REDIS_AOF_OFF;
|
||||
|
||||
fakeClient = createFakeClient();
|
||||
startLoading(fp);
|
||||
startLoadingFile(fp);
|
||||
|
||||
while(1) {
|
||||
int argc, j;
|
||||
|
||||
@@ -329,6 +329,10 @@ void loadServerConfigFromString(char *config) {
|
||||
if ((server.repl_diskless_sync = yesnotoi(argv[1])) == -1) {
|
||||
err = "argument must be 'yes' or 'no'"; goto loaderr;
|
||||
}
|
||||
} else if (!strcasecmp(argv[0],"repl-diskless-load") && argc==2) {
|
||||
if ((server.repl_diskless_load = yesnotoi(argv[1])) == -1) {
|
||||
err = "argument must be 'yes' or 'no'"; goto loaderr;
|
||||
}
|
||||
} else if (!strcasecmp(argv[0],"repl-diskless-sync-delay") && argc==2) {
|
||||
server.repl_diskless_sync_delay = atoi(argv[1]);
|
||||
if (server.repl_diskless_sync_delay < 0) {
|
||||
@@ -861,6 +865,8 @@ void configSetCommand(redisClient *c) {
|
||||
"repl-disable-tcp-nodelay",server.repl_disable_tcp_nodelay) {
|
||||
} config_set_bool_field(
|
||||
"repl-diskless-sync",server.repl_diskless_sync) {
|
||||
} config_set_bool_field(
|
||||
"repl-diskless-load",server.repl_diskless_load) {
|
||||
} config_set_bool_field(
|
||||
"cluster-require-full-coverage",server.cluster_require_full_coverage) {
|
||||
} config_set_bool_field(
|
||||
@@ -1109,6 +1115,8 @@ void configGetCommand(redisClient *c) {
|
||||
server.repl_disable_tcp_nodelay);
|
||||
config_get_bool_field("repl-diskless-sync",
|
||||
server.repl_diskless_sync);
|
||||
config_get_bool_field("repl-diskless-load",
|
||||
server.repl_diskless_load);
|
||||
config_get_bool_field("aof-rewrite-incremental-fsync",
|
||||
server.aof_rewrite_incremental_fsync);
|
||||
config_get_bool_field("aof-load-truncated",
|
||||
@@ -1780,6 +1788,7 @@ int rewriteConfig(char *path) {
|
||||
rewriteConfigBytesOption(state,"repl-backlog-ttl",server.repl_backlog_time_limit,REDIS_DEFAULT_REPL_BACKLOG_TIME_LIMIT);
|
||||
rewriteConfigYesNoOption(state,"repl-disable-tcp-nodelay",server.repl_disable_tcp_nodelay,REDIS_DEFAULT_REPL_DISABLE_TCP_NODELAY);
|
||||
rewriteConfigYesNoOption(state,"repl-diskless-sync",server.repl_diskless_sync,REDIS_DEFAULT_REPL_DISKLESS_SYNC);
|
||||
rewriteConfigYesNoOption(state,"repl-diskless-load",server.repl_diskless_load,REDIS_DEFAULT_REPL_DISKLESS_LOAD);
|
||||
rewriteConfigNumericalOption(state,"repl-diskless-sync-delay",server.repl_diskless_sync_delay,REDIS_DEFAULT_REPL_DISKLESS_SYNC_DELAY);
|
||||
rewriteConfigNumericalOption(state,"slave-priority",server.slave_priority,REDIS_DEFAULT_SLAVE_PRIORITY);
|
||||
rewriteConfigNumericalOption(state,"min-slaves-to-write",server.repl_min_slaves_to_write,REDIS_DEFAULT_MIN_SLAVES_TO_WRITE);
|
||||
|
||||
@@ -121,6 +121,7 @@ redisClient *createClient(int fd) {
|
||||
c->pubsub_channels = dictCreate(&setDictType,NULL);
|
||||
c->pubsub_patterns = listCreate();
|
||||
c->peerid = NULL;
|
||||
c->repl_eof_supported = 0;
|
||||
listSetFreeMethod(c->pubsub_patterns,decrRefCountVoid);
|
||||
listSetMatchMethod(c->pubsub_patterns,listMatchObjects);
|
||||
if (fd != -1) listAddNodeTail(server.clients,c);
|
||||
|
||||
@@ -1187,18 +1187,19 @@ robj *rdbLoadObject(int rdbtype, rio *rdb) {
|
||||
|
||||
/* Mark that we are loading in the global state and setup the fields
|
||||
* needed to provide loading stats. */
|
||||
void startLoading(FILE *fp) {
|
||||
struct stat sb;
|
||||
|
||||
void startLoading(size_t size) {
|
||||
/* Load the DB */
|
||||
server.loading = 1;
|
||||
server.loading_start_time = time(NULL);
|
||||
server.loading_loaded_bytes = 0;
|
||||
if (fstat(fileno(fp), &sb) == -1) {
|
||||
server.loading_total_bytes = 0;
|
||||
} else {
|
||||
server.loading_total_bytes = sb.st_size;
|
||||
}
|
||||
server.loading_total_bytes = size;
|
||||
}
|
||||
|
||||
void startLoadingFile(FILE *fp) {
|
||||
struct stat sb;
|
||||
if (fstat(fileno(fp), &sb) == -1)
|
||||
sb.st_size = 0;
|
||||
startLoading(sb.st_size);
|
||||
}
|
||||
|
||||
/* Refresh the loading progress info */
|
||||
@@ -1233,66 +1234,74 @@ void rdbLoadProgressCallback(rio *r, const void *buf, size_t len) {
|
||||
}
|
||||
|
||||
int rdbLoad(char *filename) {
|
||||
FILE *fp;
|
||||
int retval;
|
||||
rio rdb;
|
||||
struct stat sb;
|
||||
if ((fp = fopen(filename,"r")) == NULL) return REDIS_ERR;
|
||||
if (fstat(fileno(fp), &sb) == -1)
|
||||
sb.st_size = 0;
|
||||
rioInitWithFile(&rdb,fp);
|
||||
retval = rdbLoadRio(&rdb, sb.st_size);
|
||||
fclose(fp);
|
||||
return retval;
|
||||
}
|
||||
|
||||
int rdbLoadRio(rio *rdb, size_t size) {
|
||||
uint32_t dbid;
|
||||
int type, rdbver;
|
||||
redisDb *db = server.db+0;
|
||||
char buf[1024];
|
||||
long long expiretime, now = mstime();
|
||||
FILE *fp;
|
||||
rio rdb;
|
||||
|
||||
if ((fp = fopen(filename,"r")) == NULL) return REDIS_ERR;
|
||||
|
||||
rioInitWithFile(&rdb,fp);
|
||||
rdb.update_cksum = rdbLoadProgressCallback;
|
||||
rdb.max_processing_chunk = server.loading_process_events_interval_bytes;
|
||||
if (rioRead(&rdb,buf,9) == 0) goto eoferr;
|
||||
rdb->update_cksum = rdbLoadProgressCallback;
|
||||
rdb->max_processing_chunk = server.loading_process_events_interval_bytes;
|
||||
if (rioRead(rdb,buf,9) == 0) goto eoferr;
|
||||
buf[9] = '\0';
|
||||
if (memcmp(buf,"REDIS",5) != 0) {
|
||||
fclose(fp);
|
||||
redisLog(REDIS_WARNING,"Wrong signature trying to load DB from file");
|
||||
errno = EINVAL;
|
||||
return REDIS_ERR;
|
||||
}
|
||||
rdbver = atoi(buf+5);
|
||||
if (rdbver < 1 || rdbver > REDIS_RDB_VERSION) {
|
||||
fclose(fp);
|
||||
redisLog(REDIS_WARNING,"Can't handle RDB format version %d",rdbver);
|
||||
errno = EINVAL;
|
||||
return REDIS_ERR;
|
||||
}
|
||||
|
||||
startLoading(fp);
|
||||
startLoading(size);
|
||||
while(1) {
|
||||
robj *key, *val;
|
||||
expiretime = -1;
|
||||
|
||||
/* Read type. */
|
||||
if ((type = rdbLoadType(&rdb)) == -1) goto eoferr;
|
||||
if ((type = rdbLoadType(rdb)) == -1) goto eoferr;
|
||||
|
||||
/* Handle special types. */
|
||||
if (type == REDIS_RDB_OPCODE_EXPIRETIME) {
|
||||
/* EXPIRETIME: load an expire associated with the next key
|
||||
* to load. Note that after loading an expire we need to
|
||||
* load the actual type, and continue. */
|
||||
if ((expiretime = rdbLoadTime(&rdb)) == -1) goto eoferr;
|
||||
if ((expiretime = rdbLoadTime(rdb)) == -1) goto eoferr;
|
||||
/* We read the time so we need to read the object type again. */
|
||||
if ((type = rdbLoadType(&rdb)) == -1) goto eoferr;
|
||||
if ((type = rdbLoadType(rdb)) == -1) goto eoferr;
|
||||
/* the EXPIRETIME opcode specifies time in seconds, so convert
|
||||
* into milliseconds. */
|
||||
expiretime *= 1000;
|
||||
} else if (type == REDIS_RDB_OPCODE_EXPIRETIME_MS) {
|
||||
/* EXPIRETIME_MS: milliseconds precision expire times introduced
|
||||
* with RDB v3. Like EXPIRETIME but no with more precision. */
|
||||
if ((expiretime = rdbLoadMillisecondTime(&rdb)) == -1) goto eoferr;
|
||||
if ((expiretime = rdbLoadMillisecondTime(rdb)) == -1) goto eoferr;
|
||||
/* We read the time so we need to read the object type again. */
|
||||
if ((type = rdbLoadType(&rdb)) == -1) goto eoferr;
|
||||
if ((type = rdbLoadType(rdb)) == -1) goto eoferr;
|
||||
} else if (type == REDIS_RDB_OPCODE_EOF) {
|
||||
/* EOF: End of file, exit the main loop. */
|
||||
break;
|
||||
} else if (type == REDIS_RDB_OPCODE_SELECTDB) {
|
||||
/* SELECTDB: Select the specified database. */
|
||||
if ((dbid = rdbLoadLen(&rdb,NULL)) == REDIS_RDB_LENERR)
|
||||
if ((dbid = rdbLoadLen(rdb,NULL)) == REDIS_RDB_LENERR)
|
||||
goto eoferr;
|
||||
if (dbid >= (unsigned)server.dbnum) {
|
||||
redisLog(REDIS_WARNING,
|
||||
@@ -1307,9 +1316,9 @@ int rdbLoad(char *filename) {
|
||||
/* RESIZEDB: Hint about the size of the keys in the currently
|
||||
* selected data base, in order to avoid useless rehashing. */
|
||||
uint32_t db_size, expires_size;
|
||||
if ((db_size = rdbLoadLen(&rdb,NULL)) == REDIS_RDB_LENERR)
|
||||
if ((db_size = rdbLoadLen(rdb,NULL)) == REDIS_RDB_LENERR)
|
||||
goto eoferr;
|
||||
if ((expires_size = rdbLoadLen(&rdb,NULL)) == REDIS_RDB_LENERR)
|
||||
if ((expires_size = rdbLoadLen(rdb,NULL)) == REDIS_RDB_LENERR)
|
||||
goto eoferr;
|
||||
dictExpand(db->dict,db_size);
|
||||
dictExpand(db->expires,expires_size);
|
||||
@@ -1321,8 +1330,8 @@ int rdbLoad(char *filename) {
|
||||
*
|
||||
* An AUX field is composed of two strings: key and value. */
|
||||
robj *auxkey, *auxval;
|
||||
if ((auxkey = rdbLoadStringObject(&rdb)) == NULL) goto eoferr;
|
||||
if ((auxval = rdbLoadStringObject(&rdb)) == NULL) goto eoferr;
|
||||
if ((auxkey = rdbLoadStringObject(rdb)) == NULL) goto eoferr;
|
||||
if ((auxval = rdbLoadStringObject(rdb)) == NULL) goto eoferr;
|
||||
|
||||
if (((char*)auxkey->ptr)[0] == '%') {
|
||||
/* All the fields with a name staring with '%' are considered
|
||||
@@ -1344,9 +1353,9 @@ int rdbLoad(char *filename) {
|
||||
}
|
||||
|
||||
/* Read key */
|
||||
if ((key = rdbLoadStringObject(&rdb)) == NULL) goto eoferr;
|
||||
if ((key = rdbLoadStringObject(rdb)) == NULL) goto eoferr;
|
||||
/* Read value */
|
||||
if ((val = rdbLoadObject(type,&rdb)) == NULL) goto eoferr;
|
||||
if ((val = rdbLoadObject(type,rdb)) == NULL) goto eoferr;
|
||||
/* Check if the key already expired. This function is used when loading
|
||||
* an RDB file from disk, either at startup, or when an RDB was
|
||||
* received from the master. In the latter case, the master is
|
||||
@@ -1367,9 +1376,9 @@ int rdbLoad(char *filename) {
|
||||
}
|
||||
/* Verify the checksum if RDB version is >= 5 */
|
||||
if (rdbver >= 5 && server.rdb_checksum) {
|
||||
uint64_t cksum, expected = rdb.cksum;
|
||||
uint64_t cksum, expected = rdb->cksum;
|
||||
|
||||
if (rioRead(&rdb,&cksum,8) == 0) goto eoferr;
|
||||
if (rioRead(rdb,&cksum,8) == 0) goto eoferr;
|
||||
memrev64ifbe(&cksum);
|
||||
if (cksum == 0) {
|
||||
redisLog(REDIS_WARNING,"RDB file was saved with checksum disabled: no check performed.");
|
||||
@@ -1379,7 +1388,6 @@ int rdbLoad(char *filename) {
|
||||
}
|
||||
}
|
||||
|
||||
fclose(fp);
|
||||
stopLoading();
|
||||
return REDIS_OK;
|
||||
|
||||
@@ -1572,7 +1580,7 @@ int rdbSaveToSlavesSockets(void) {
|
||||
clientids[numfds] = slave->id;
|
||||
fds[numfds++] = slave->fd;
|
||||
slave->replstate = REDIS_REPL_WAIT_BGSAVE_END;
|
||||
/* Put the socket in non-blocking mode to simplify RDB transfer.
|
||||
/* Put the socket in blocking mode to simplify RDB transfer.
|
||||
* We'll restore it when the children returns (since duped socket
|
||||
* will share the O_NONBLOCK attribute with the parent). */
|
||||
anetBlock(NULL,slave->fd);
|
||||
@@ -1646,6 +1654,7 @@ int rdbSaveToSlavesSockets(void) {
|
||||
zfree(msg);
|
||||
}
|
||||
zfree(clientids);
|
||||
rioFreeFdset(&slave_sockets);
|
||||
exitFromChild((retval == REDIS_OK) ? 0 : 1);
|
||||
} else {
|
||||
/* Parent */
|
||||
|
||||
@@ -105,6 +105,7 @@ uint32_t rdbLoadLen(rio *rdb, int *isencoded);
|
||||
int rdbSaveObjectType(rio *rdb, robj *o);
|
||||
int rdbLoadObjectType(rio *rdb);
|
||||
int rdbLoad(char *filename);
|
||||
int rdbLoadRio(rio *rdb, size_t size);
|
||||
int rdbSaveBackground(char *filename);
|
||||
int rdbSaveToSlavesSockets(void);
|
||||
void rdbRemoveTempFile(pid_t childpid);
|
||||
|
||||
@@ -1513,12 +1513,16 @@ void initServerConfig(void) {
|
||||
server.cached_master = NULL;
|
||||
server.repl_master_initial_offset = -1;
|
||||
server.repl_state = REDIS_REPL_NONE;
|
||||
server.repl_transfer_tmpfile = NULL;
|
||||
server.repl_transfer_fd = -1;
|
||||
server.repl_transfer_s = -1;
|
||||
server.repl_syncio_timeout = REDIS_REPL_SYNCIO_TIMEOUT;
|
||||
server.repl_serve_stale_data = REDIS_DEFAULT_SLAVE_SERVE_STALE_DATA;
|
||||
server.repl_slave_ro = REDIS_DEFAULT_SLAVE_READ_ONLY;
|
||||
server.repl_down_since = 0; /* Never connected, repl is down since EVER. */
|
||||
server.repl_disable_tcp_nodelay = REDIS_DEFAULT_REPL_DISABLE_TCP_NODELAY;
|
||||
server.repl_diskless_sync = REDIS_DEFAULT_REPL_DISKLESS_SYNC;
|
||||
server.repl_diskless_load = REDIS_DEFAULT_REPL_DISKLESS_LOAD;
|
||||
server.repl_diskless_sync_delay = REDIS_DEFAULT_REPL_DISKLESS_SYNC_DELAY;
|
||||
server.slave_priority = REDIS_DEFAULT_SLAVE_PRIORITY;
|
||||
server.master_repl_offset = 0;
|
||||
|
||||
+6
-2
@@ -118,6 +118,7 @@ typedef long long mstime_t; /* millisecond time type. */
|
||||
#define REDIS_DEFAULT_RDB_CHECKSUM 1
|
||||
#define REDIS_DEFAULT_RDB_FILENAME "dump.rdb"
|
||||
#define REDIS_DEFAULT_REPL_DISKLESS_SYNC 0
|
||||
#define REDIS_DEFAULT_REPL_DISKLESS_LOAD 0
|
||||
#define REDIS_DEFAULT_REPL_DISKLESS_SYNC_DELAY 5
|
||||
#define REDIS_DEFAULT_SLAVE_SERVE_STALE_DATA 1
|
||||
#define REDIS_DEFAULT_SLAVE_READ_ONLY 1
|
||||
@@ -570,6 +571,7 @@ typedef struct redisClient {
|
||||
dict *pubsub_channels; /* channels a client is interested in (SUBSCRIBE) */
|
||||
list *pubsub_patterns; /* patterns a client is interested in (SUBSCRIBE) */
|
||||
sds peerid; /* Cached peer ID. */
|
||||
int repl_eof_supported; /* slave supports EOF based, diskless replication */
|
||||
|
||||
/* Response buffer */
|
||||
int bufpos;
|
||||
@@ -828,7 +830,8 @@ struct redisServer {
|
||||
int repl_min_slaves_to_write; /* Min number of slaves to write. */
|
||||
int repl_min_slaves_max_lag; /* Max lag of <count> slaves to write. */
|
||||
int repl_good_slaves_count; /* Number of slaves with lag <= max_lag. */
|
||||
int repl_diskless_sync; /* Send RDB to slaves sockets directly. */
|
||||
int repl_diskless_load; /* Slave parse RDB directly from the socket. */
|
||||
int repl_diskless_sync; /* Master send RDB to slaves sockets directly. */
|
||||
int repl_diskless_sync_delay; /* Delay to start a diskless repl BGSAVE. */
|
||||
/* Replication (slave) */
|
||||
char *masterauth; /* AUTH with this password with master */
|
||||
@@ -1196,7 +1199,8 @@ long long replicationGetSlaveOffset(void);
|
||||
char *replicationGetSlaveName(redisClient *c);
|
||||
|
||||
/* Generic persistence functions */
|
||||
void startLoading(FILE *fp);
|
||||
void startLoading(size_t size);
|
||||
void startLoadingFile(FILE *fp);
|
||||
void loadingProgress(off_t pos);
|
||||
void stopLoading(void);
|
||||
|
||||
|
||||
+226
-122
@@ -441,14 +441,19 @@ need_full_resync:
|
||||
* socket target depending on the configuration, and making sure that
|
||||
* the script cache is flushed before to start.
|
||||
*
|
||||
* Returns REDIS_OK on success or REDIS_ERR otherwise. */
|
||||
int startBgsaveForReplication(void) {
|
||||
* Returns REDIS_OK on success or REDIS_ERR otherwise.
|
||||
*
|
||||
* The caller should pass '1' as the function argument if all the slaves
|
||||
* currently waiting for a BGSAVE all claimed to support the EOF-style
|
||||
* streaming format for RDB transfer. Otherwise it should be '0'. */
|
||||
int startBgsaveForReplication(int all_slaves_supprot_eof) {
|
||||
int retval;
|
||||
int use_eof = all_slaves_support_eof && server.repl_diskless_sync;
|
||||
|
||||
redisLog(REDIS_NOTICE,"Starting BGSAVE for SYNC with target: %s",
|
||||
server.repl_diskless_sync ? "slaves sockets" : "disk");
|
||||
use_eof ? "slaves sockets" : "disk");
|
||||
|
||||
if (server.repl_diskless_sync)
|
||||
if (use_eof)
|
||||
retval = rdbSaveToSlavesSockets();
|
||||
else
|
||||
retval = rdbSaveBackground(server.rdb_filename);
|
||||
@@ -553,7 +558,7 @@ void syncCommand(redisClient *c) {
|
||||
c->replstate = REDIS_REPL_WAIT_BGSAVE_START;
|
||||
redisLog(REDIS_NOTICE,"Waiting for next BGSAVE for SYNC");
|
||||
} else {
|
||||
if (server.repl_diskless_sync) {
|
||||
if (server.repl_diskless_sync && c->repl_eof_supported) {
|
||||
/* Diskless replication RDB child is created inside
|
||||
* replicationCron() since we want to delay its start a
|
||||
* few seconds to wait for more slaves to arrive. */
|
||||
@@ -562,7 +567,7 @@ void syncCommand(redisClient *c) {
|
||||
redisLog(REDIS_NOTICE,"Delay next BGSAVE for SYNC");
|
||||
} else {
|
||||
/* Ok we don't have a BGSAVE in progress, let's start one. */
|
||||
if (startBgsaveForReplication() != REDIS_OK) {
|
||||
if (startBgsaveForReplication(0) != REDIS_OK) {
|
||||
redisLog(REDIS_NOTICE,"Replication failed, can't BGSAVE");
|
||||
addReplyError(c,"Unable to perform background save");
|
||||
return;
|
||||
@@ -613,6 +618,8 @@ void replconfCommand(redisClient *c) {
|
||||
&port,NULL) != REDIS_OK))
|
||||
return;
|
||||
c->slave_listening_port = port;
|
||||
} else if (!strcasecmp(c->argv[j]->ptr,"eof-supported")) {
|
||||
c->repl_eof_supported = 1;
|
||||
} else if (!strcasecmp(c->argv[j]->ptr,"ack")) {
|
||||
/* REPLCONF ACK is used by slave to inform the master the amount
|
||||
* of replication stream that it processed so far. It is an
|
||||
@@ -745,7 +752,8 @@ void sendBulkToSlave(aeEventLoop *el, int fd, void *privdata, int mask) {
|
||||
* (if it had a disk or socket target). */
|
||||
void updateSlavesWaitingBgsave(int bgsaveerr, int type) {
|
||||
listNode *ln;
|
||||
int startbgsave = 0;
|
||||
int slaves_waiting_eof = 0;
|
||||
int slaves_waiting_noneof = 0;
|
||||
listIter li;
|
||||
|
||||
listRewind(server.slaves,&li);
|
||||
@@ -753,7 +761,10 @@ void updateSlavesWaitingBgsave(int bgsaveerr, int type) {
|
||||
redisClient *slave = ln->value;
|
||||
|
||||
if (slave->replstate == REDIS_REPL_WAIT_BGSAVE_START) {
|
||||
startbgsave = 1;
|
||||
if (slave->repl_eof_supported)
|
||||
slaves_waiting_eof++;
|
||||
else
|
||||
slaves_waiting_noneof++;
|
||||
slave->replstate = REDIS_REPL_WAIT_BGSAVE_END;
|
||||
} else if (slave->replstate == REDIS_REPL_WAIT_BGSAVE_END) {
|
||||
struct redis_stat buf;
|
||||
@@ -801,8 +812,10 @@ void updateSlavesWaitingBgsave(int bgsaveerr, int type) {
|
||||
}
|
||||
}
|
||||
}
|
||||
if (startbgsave) {
|
||||
if (startBgsaveForReplication() != REDIS_OK) {
|
||||
if (slaves_waiting_eof || slaves_waiting_noneof) {
|
||||
/* if there is at least one slave that doesn't support EOF, we'll
|
||||
* start an non-eof replication */
|
||||
if (startBgsaveForReplication(slaves_waiting_noneof==0) != REDIS_OK) {
|
||||
listIter li;
|
||||
|
||||
listRewind(server.slaves,&li);
|
||||
@@ -825,9 +838,13 @@ void replicationAbortSyncTransfer(void) {
|
||||
|
||||
aeDeleteFileEvent(server.el,server.repl_transfer_s,AE_READABLE);
|
||||
close(server.repl_transfer_s);
|
||||
close(server.repl_transfer_fd);
|
||||
unlink(server.repl_transfer_tmpfile);
|
||||
zfree(server.repl_transfer_tmpfile);
|
||||
if (server.repl_transfer_fd!=-1) {
|
||||
close(server.repl_transfer_fd);
|
||||
unlink(server.repl_transfer_tmpfile);
|
||||
zfree(server.repl_transfer_tmpfile);
|
||||
server.repl_transfer_tmpfile = NULL;
|
||||
server.repl_transfer_fd = -1;
|
||||
}
|
||||
server.repl_state = REDIS_REPL_CONNECT;
|
||||
}
|
||||
|
||||
@@ -879,6 +896,7 @@ void readSyncBulkPayload(aeEventLoop *el, int fd, void *privdata, int mask) {
|
||||
char buf[4096];
|
||||
ssize_t nread, readlen;
|
||||
off_t left;
|
||||
sds initialCommandStream = NULL;
|
||||
REDIS_NOTUSED(el);
|
||||
REDIS_NOTUSED(privdata);
|
||||
REDIS_NOTUSED(mask);
|
||||
@@ -933,129 +951,189 @@ void readSyncBulkPayload(aeEventLoop *el, int fd, void *privdata, int mask) {
|
||||
* at the next call. */
|
||||
server.repl_transfer_size = 0;
|
||||
redisLog(REDIS_NOTICE,
|
||||
"MASTER <-> SLAVE sync: receiving streamed RDB from master");
|
||||
"MASTER <-> SLAVE sync: receiving streamed RDB from master with EOF %s",
|
||||
server.repl_diskless_load? "to parser":"to disk");
|
||||
} else {
|
||||
usemark = 0;
|
||||
server.repl_transfer_size = strtol(buf+1,NULL,10);
|
||||
redisLog(REDIS_NOTICE,
|
||||
"MASTER <-> SLAVE sync: receiving %lld bytes from master",
|
||||
(long long) server.repl_transfer_size);
|
||||
"MASTER <-> SLAVE sync: receiving %lld bytes from master %s",
|
||||
(long long) server.repl_transfer_size,
|
||||
server.repl_diskless_load? "to parser":"to disk");
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
/* Read bulk data */
|
||||
if (usemark) {
|
||||
readlen = sizeof(buf);
|
||||
} else {
|
||||
left = server.repl_transfer_size - server.repl_transfer_read;
|
||||
readlen = (left < (signed)sizeof(buf)) ? left : (signed)sizeof(buf);
|
||||
}
|
||||
|
||||
nread = read(fd,buf,readlen);
|
||||
if (nread <= 0) {
|
||||
redisLog(REDIS_WARNING,"I/O error trying to sync with MASTER: %s",
|
||||
(nread == -1) ? strerror(errno) : "connection lost");
|
||||
replicationAbortSyncTransfer();
|
||||
return;
|
||||
}
|
||||
server.stat_net_input_bytes += nread;
|
||||
|
||||
/* When a mark is used, we want to detect EOF asap in order to avoid
|
||||
* writing the EOF mark into the file... */
|
||||
int eof_reached = 0;
|
||||
|
||||
if (usemark) {
|
||||
/* Update the last bytes array, and check if it matches our delimiter.*/
|
||||
if (nread >= REDIS_RUN_ID_SIZE) {
|
||||
memcpy(lastbytes,buf+nread-REDIS_RUN_ID_SIZE,REDIS_RUN_ID_SIZE);
|
||||
if (!server.repl_diskless_load) {
|
||||
/* read the data from the socket, store it to a file and search for the EOF */
|
||||
if (usemark) {
|
||||
readlen = sizeof(buf);
|
||||
} else {
|
||||
int rem = REDIS_RUN_ID_SIZE-nread;
|
||||
memmove(lastbytes,lastbytes+nread,rem);
|
||||
memcpy(lastbytes+rem,buf,nread);
|
||||
left = server.repl_transfer_size - server.repl_transfer_read;
|
||||
readlen = (left < (signed)sizeof(buf)) ? left : (signed)sizeof(buf);
|
||||
}
|
||||
if (memcmp(lastbytes,eofmark,REDIS_RUN_ID_SIZE) == 0) eof_reached = 1;
|
||||
}
|
||||
|
||||
server.repl_transfer_lastio = server.unixtime;
|
||||
if (write(server.repl_transfer_fd,buf,nread) != nread) {
|
||||
redisLog(REDIS_WARNING,"Write error or short write writing to the DB dump file needed for MASTER <-> SLAVE synchronization: %s", strerror(errno));
|
||||
goto error;
|
||||
}
|
||||
server.repl_transfer_read += nread;
|
||||
nread = read(fd,buf,readlen);
|
||||
if (nread <= 0) {
|
||||
redisLog(REDIS_WARNING,"I/O error trying to sync with MASTER: %s",
|
||||
(nread == -1) ? strerror(errno) : "connection lost");
|
||||
replicationAbortSyncTransfer();
|
||||
return;
|
||||
}
|
||||
server.stat_net_input_bytes += nread;
|
||||
|
||||
/* Delete the last 40 bytes from the file if we reached EOF. */
|
||||
if (usemark && eof_reached) {
|
||||
if (ftruncate(server.repl_transfer_fd,
|
||||
server.repl_transfer_read - REDIS_RUN_ID_SIZE) == -1)
|
||||
{
|
||||
redisLog(REDIS_WARNING,"Error truncating the RDB file received from the master for SYNC: %s", strerror(errno));
|
||||
/* When a mark is used, we want to detect EOF asap in order to avoid
|
||||
* writing the EOF mark into the file... */
|
||||
int eof_reached = 0;
|
||||
|
||||
if (usemark) {
|
||||
/* Update the last bytes array, and check if it matches our delimiter.*/
|
||||
if (nread >= REDIS_RUN_ID_SIZE) {
|
||||
memcpy(lastbytes,buf+nread-REDIS_RUN_ID_SIZE,REDIS_RUN_ID_SIZE);
|
||||
} else {
|
||||
int rem = REDIS_RUN_ID_SIZE-nread;
|
||||
memmove(lastbytes,lastbytes+nread,rem);
|
||||
memcpy(lastbytes+rem,buf,nread);
|
||||
}
|
||||
if (memcmp(lastbytes,eofmark,REDIS_RUN_ID_SIZE) == 0) eof_reached = 1;
|
||||
}
|
||||
|
||||
server.repl_transfer_lastio = server.unixtime;
|
||||
if (write(server.repl_transfer_fd,buf,nread) != nread) {
|
||||
redisLog(REDIS_WARNING,"Write error or short write writing to the DB dump file needed for MASTER <-> SLAVE synchronization: %s", strerror(errno));
|
||||
goto error;
|
||||
}
|
||||
server.repl_transfer_read += nread;
|
||||
|
||||
/* Delete the last 40 bytes from the file if we reached EOF. */
|
||||
if (usemark && eof_reached) {
|
||||
if (ftruncate(server.repl_transfer_fd,
|
||||
server.repl_transfer_read - REDIS_RUN_ID_SIZE) == -1)
|
||||
{
|
||||
redisLog(REDIS_WARNING,"Error truncating the RDB file received from the master for SYNC: %s", strerror(errno));
|
||||
goto error;
|
||||
}
|
||||
}
|
||||
|
||||
/* Sync data on disk from time to time, otherwise at the end of the transfer
|
||||
* we may suffer a big delay as the memory buffers are copied into the
|
||||
* actual disk. */
|
||||
if (server.repl_transfer_read >=
|
||||
server.repl_transfer_last_fsync_off + REPL_MAX_WRITTEN_BEFORE_FSYNC)
|
||||
{
|
||||
off_t sync_size = server.repl_transfer_read -
|
||||
server.repl_transfer_last_fsync_off;
|
||||
rdb_fsync_range(server.repl_transfer_fd,
|
||||
server.repl_transfer_last_fsync_off, sync_size);
|
||||
server.repl_transfer_last_fsync_off += sync_size;
|
||||
}
|
||||
|
||||
/* Check if the transfer is now complete */
|
||||
if (!usemark) {
|
||||
if (server.repl_transfer_read == server.repl_transfer_size)
|
||||
eof_reached = 1;
|
||||
}
|
||||
if (!eof_reached)
|
||||
return;
|
||||
}
|
||||
|
||||
/* Sync data on disk from time to time, otherwise at the end of the transfer
|
||||
* we may suffer a big delay as the memory buffers are copied into the
|
||||
* actual disk. */
|
||||
if (server.repl_transfer_read >=
|
||||
server.repl_transfer_last_fsync_off + REPL_MAX_WRITTEN_BEFORE_FSYNC)
|
||||
{
|
||||
off_t sync_size = server.repl_transfer_read -
|
||||
server.repl_transfer_last_fsync_off;
|
||||
rdb_fsync_range(server.repl_transfer_fd,
|
||||
server.repl_transfer_last_fsync_off, sync_size);
|
||||
server.repl_transfer_last_fsync_off += sync_size;
|
||||
}
|
||||
/* We reach here when the slave is using diskless replication,
|
||||
* or when we are done reading from the socket to the rdb file. */
|
||||
redisLog(REDIS_NOTICE, "MASTER <-> SLAVE sync: Flushing old data");
|
||||
signalFlushedDb(-1);
|
||||
emptyDb(replicationEmptyDbCallback);
|
||||
/* Before loading the DB into memory we need to delete the readable
|
||||
* handler, otherwise it will get called recursively since
|
||||
* rdbLoad() will call the event loop to process events from time to
|
||||
* time for non blocking loading. */
|
||||
aeDeleteFileEvent(server.el,server.repl_transfer_s,AE_READABLE);
|
||||
redisLog(REDIS_NOTICE, "MASTER <-> SLAVE sync: Loading DB in memory");
|
||||
if (server.repl_diskless_load) {
|
||||
rio rdb;
|
||||
rioInitWithFd(&rdb,fd);
|
||||
/* Put the socket in blocking mode to simplify RDB transfer.
|
||||
* We'll restore it when the RDB is received. */
|
||||
anetBlock(NULL,fd);
|
||||
anetRecvTimeout(NULL,fd,server.repl_timeout*1000);
|
||||
|
||||
/* Check if the transfer is now complete */
|
||||
if (!usemark) {
|
||||
if (server.repl_transfer_read == server.repl_transfer_size)
|
||||
eof_reached = 1;
|
||||
}
|
||||
|
||||
if (eof_reached) {
|
||||
if (rdbLoadRio(&rdb, server.repl_transfer_size) != REDIS_OK) {
|
||||
redisLog(REDIS_WARNING,"Failed trying to load the MASTER synchronization DB from disk");
|
||||
replicationAbortSyncTransfer();
|
||||
rioFreeFd(&rdb, NULL);
|
||||
/* Remove the half-loaded data, and load back the old dataset
|
||||
* if we have persistence turned on.
|
||||
*
|
||||
* TODO:
|
||||
* 1) Actually allow rdbLoadRio() to don't fail with exit().
|
||||
* 2) Load RDB / AOF.
|
||||
*
|
||||
* Right now this code path is not entered when the connection
|
||||
* breaks between master and slave AFAIK.
|
||||
*/
|
||||
emptyDb(NULL);
|
||||
return;
|
||||
}
|
||||
if (usemark) {
|
||||
if (!rioRead(&rdb,buf,REDIS_RUN_ID_SIZE) || memcmp(buf,eofmark,REDIS_RUN_ID_SIZE) != 0) {
|
||||
redisLog(REDIS_WARNING,"replication stream EOF marker is broken");
|
||||
replicationAbortSyncTransfer();
|
||||
rioFreeFd(&rdb, NULL);
|
||||
return;
|
||||
}
|
||||
}
|
||||
/* get the unread command stream from the rio buffer */
|
||||
rioFreeFd(&rdb, &initialCommandStream);
|
||||
/* Restore the socket as non-blocking. */
|
||||
anetNonBlock(NULL,fd);
|
||||
anetRecvTimeout(NULL,fd,0);
|
||||
} else {
|
||||
if (rename(server.repl_transfer_tmpfile,server.rdb_filename) == -1) {
|
||||
redisLog(REDIS_WARNING,"Failed trying to rename the temp DB into dump.rdb in MASTER <-> SLAVE synchronization: %s", strerror(errno));
|
||||
replicationAbortSyncTransfer();
|
||||
return;
|
||||
}
|
||||
redisLog(REDIS_NOTICE, "MASTER <-> SLAVE sync: Flushing old data");
|
||||
signalFlushedDb(-1);
|
||||
emptyDb(replicationEmptyDbCallback);
|
||||
/* Before loading the DB into memory we need to delete the readable
|
||||
* handler, otherwise it will get called recursively since
|
||||
* rdbLoad() will call the event loop to process events from time to
|
||||
* time for non blocking loading. */
|
||||
aeDeleteFileEvent(server.el,server.repl_transfer_s,AE_READABLE);
|
||||
redisLog(REDIS_NOTICE, "MASTER <-> SLAVE sync: Loading DB in memory");
|
||||
if (rdbLoad(server.rdb_filename) != REDIS_OK) {
|
||||
redisLog(REDIS_WARNING,"Failed trying to load the MASTER synchronization DB from disk");
|
||||
replicationAbortSyncTransfer();
|
||||
return;
|
||||
}
|
||||
/* Final setup of the connected slave <- master link */
|
||||
zfree(server.repl_transfer_tmpfile);
|
||||
close(server.repl_transfer_fd);
|
||||
replicationCreateMasterClient(server.repl_transfer_s);
|
||||
redisLog(REDIS_NOTICE, "MASTER <-> SLAVE sync: Finished with success");
|
||||
/* Restart the AOF subsystem now that we finished the sync. This
|
||||
* will trigger an AOF rewrite, and when done will start appending
|
||||
* to the new file. */
|
||||
if (server.aof_state != REDIS_AOF_OFF) {
|
||||
int retry = 10;
|
||||
|
||||
stopAppendOnly();
|
||||
while (retry-- && startAppendOnly() == REDIS_ERR) {
|
||||
redisLog(REDIS_WARNING,"Failed enabling the AOF after successful master synchronization! Trying it again in one second.");
|
||||
sleep(1);
|
||||
}
|
||||
if (!retry) {
|
||||
redisLog(REDIS_WARNING,"FATAL: this slave instance finished the synchronization with its master, but the AOF can't be turned on. Exiting now.");
|
||||
exit(1);
|
||||
}
|
||||
}
|
||||
server.repl_transfer_fd = -1;
|
||||
server.repl_transfer_tmpfile = NULL;
|
||||
}
|
||||
/* Final setup of the connected slave <- master link */
|
||||
replicationCreateMasterClient(server.repl_transfer_s);
|
||||
/* feed the initial command stream into the client input buffer */
|
||||
if (initialCommandStream) {
|
||||
int nread = sdslen(initialCommandStream);
|
||||
server.current_client = server.master;
|
||||
server.master->querybuf = sdscatsds(server.master->querybuf, initialCommandStream);
|
||||
server.master->reploff += nread;
|
||||
server.master->lastinteraction = server.unixtime;
|
||||
server.stat_net_input_bytes += nread;
|
||||
processInputBuffer(server.master);
|
||||
server.current_client = NULL;
|
||||
sdsfree(initialCommandStream);
|
||||
}
|
||||
|
||||
redisLog(REDIS_NOTICE, "MASTER <-> SLAVE sync: Finished with success");
|
||||
/* Restart the AOF subsystem now that we finished the sync. This
|
||||
* will trigger an AOF rewrite, and when done will start appending
|
||||
* to the new file. */
|
||||
if (server.aof_state != REDIS_AOF_OFF) {
|
||||
int retry = 10;
|
||||
|
||||
stopAppendOnly();
|
||||
while (retry-- && startAppendOnly() == REDIS_ERR) {
|
||||
redisLog(REDIS_WARNING,"Failed enabling the AOF after successful master synchronization! Trying it again in one second.");
|
||||
sleep(1);
|
||||
}
|
||||
if (!retry) {
|
||||
redisLog(REDIS_WARNING,"FATAL: this slave instance finished the synchronization with its master, but the AOF can't be turned on. Exiting now.");
|
||||
exit(1);
|
||||
}
|
||||
}
|
||||
return;
|
||||
|
||||
error:
|
||||
@@ -1319,6 +1397,18 @@ void syncWithMaster(aeEventLoop *el, int fd, void *privdata, int mask) {
|
||||
sdsfree(err);
|
||||
}
|
||||
|
||||
/* Inform the master that this slave supports EOF marker of diskless-sync */
|
||||
{
|
||||
err = sendSynchronousCommand(fd,"REPLCONF","eof-supported","yes",
|
||||
NULL);
|
||||
/* Ignore the error if any, not all the Redis versions support
|
||||
* REPLCONF eof-supported. */
|
||||
if (err[0] == '-') {
|
||||
redisLog(REDIS_NOTICE,"(Non critical) Master does not understand REPLCONF eof-supported: %s", err);
|
||||
}
|
||||
sdsfree(err);
|
||||
}
|
||||
|
||||
/* Try a partial resynchonization. If we don't have a cached master
|
||||
* slaveTryPartialResynchronization() will at least try to use PSYNC
|
||||
* to start a full resynchronization so that we get the master run id
|
||||
@@ -1343,16 +1433,20 @@ void syncWithMaster(aeEventLoop *el, int fd, void *privdata, int mask) {
|
||||
}
|
||||
|
||||
/* Prepare a suitable temp file for bulk transfer */
|
||||
while(maxtries--) {
|
||||
snprintf(tmpfile,256,
|
||||
"temp-%d.%ld.rdb",(int)server.unixtime,(long int)getpid());
|
||||
dfd = open(tmpfile,O_CREAT|O_WRONLY|O_EXCL,0644);
|
||||
if (dfd != -1) break;
|
||||
sleep(1);
|
||||
}
|
||||
if (dfd == -1) {
|
||||
redisLog(REDIS_WARNING,"Opening the temp file needed for MASTER <-> SLAVE synchronization: %s",strerror(errno));
|
||||
goto error;
|
||||
if (!server.repl_diskless_load) {
|
||||
while(maxtries--) {
|
||||
snprintf(tmpfile,256,
|
||||
"temp-%d.%ld.rdb",(int)server.unixtime,(long int)getpid());
|
||||
dfd = open(tmpfile,O_CREAT|O_WRONLY|O_EXCL,0644);
|
||||
if (dfd != -1) break;
|
||||
sleep(1);
|
||||
}
|
||||
if (dfd == -1) {
|
||||
redisLog(REDIS_WARNING,"Opening the temp file needed for MASTER <-> SLAVE synchronization: %s",strerror(errno));
|
||||
goto error;
|
||||
}
|
||||
server.repl_transfer_tmpfile = zstrdup(tmpfile);
|
||||
server.repl_transfer_fd = dfd;
|
||||
}
|
||||
|
||||
/* Setup the non blocking download of the bulk file. */
|
||||
@@ -1369,13 +1463,17 @@ void syncWithMaster(aeEventLoop *el, int fd, void *privdata, int mask) {
|
||||
server.repl_transfer_size = -1;
|
||||
server.repl_transfer_read = 0;
|
||||
server.repl_transfer_last_fsync_off = 0;
|
||||
server.repl_transfer_fd = dfd;
|
||||
server.repl_transfer_lastio = server.unixtime;
|
||||
server.repl_transfer_tmpfile = zstrdup(tmpfile);
|
||||
return;
|
||||
|
||||
error:
|
||||
close(fd);
|
||||
if (server.repl_transfer_fd != -1)
|
||||
close(server.repl_transfer_fd);
|
||||
if (server.repl_transfer_tmpfile)
|
||||
zfree(server.repl_transfer_tmpfile);
|
||||
server.repl_transfer_tmpfile = NULL;
|
||||
server.repl_transfer_fd = -1;
|
||||
server.repl_transfer_s = -1;
|
||||
server.repl_state = REDIS_REPL_CONNECT;
|
||||
return;
|
||||
@@ -2073,7 +2171,8 @@ void replicationCron(void) {
|
||||
* slaves in WAIT_BGSAVE_START state. */
|
||||
if (server.rdb_child_pid == -1 && server.aof_child_pid == -1) {
|
||||
time_t idle, max_idle = 0;
|
||||
int slaves_waiting = 0;
|
||||
int slaves_waiting_eof = 0;
|
||||
int slaves_waiting_noneof = 0;
|
||||
listNode *ln;
|
||||
listIter li;
|
||||
|
||||
@@ -2083,14 +2182,19 @@ void replicationCron(void) {
|
||||
if (slave->replstate == REDIS_REPL_WAIT_BGSAVE_START) {
|
||||
idle = server.unixtime - slave->lastinteraction;
|
||||
if (idle > max_idle) max_idle = idle;
|
||||
slaves_waiting++;
|
||||
if (slave->repl_eof_supported)
|
||||
slaves_waiting_eof++;
|
||||
else
|
||||
slaves_waiting_noneof++;
|
||||
}
|
||||
}
|
||||
|
||||
if (slaves_waiting && max_idle > server.repl_diskless_sync_delay) {
|
||||
if ((slaves_waiting_eof || slaves_waiting_noneof) && max_idle > server.repl_diskless_sync_delay) {
|
||||
/* Start a BGSAVE. Usually with socket target, or with disk target
|
||||
* if there was a recent socket -> disk config change. */
|
||||
if (startBgsaveForReplication() == REDIS_OK) {
|
||||
* if there was a recent socket -> disk config change.
|
||||
* if there is at least one slave that doesn't support EOF, we'll
|
||||
* start an non-eof replication */
|
||||
if (startBgsaveForReplication(slaves_waiting_noneof==0) == REDIS_OK){
|
||||
/* It started! We need to change the state of slaves
|
||||
* from WAIT_BGSAVE_START to WAIT_BGSAVE_END in case
|
||||
* the current target is disk. Otherwise it was already done
|
||||
|
||||
@@ -157,13 +157,106 @@ void rioInitWithFile(rio *r, FILE *fp) {
|
||||
r->io.file.autosync = 0;
|
||||
}
|
||||
|
||||
/* ------------------- File descriptor implementation ------------------- */
|
||||
|
||||
static size_t rioFdWrite(rio *r, const void *buf, size_t len) {
|
||||
REDIS_NOTUSED(r);
|
||||
REDIS_NOTUSED(buf);
|
||||
REDIS_NOTUSED(len);
|
||||
return 0; /* Error, this target does not yet support writing. */
|
||||
}
|
||||
|
||||
/* Returns 1 or 0 for success/failure. */
|
||||
static size_t rioFdRead(rio *r, void *buf, size_t len) {
|
||||
size_t avail = sdslen(r->io.fd.buf)-r->io.fd.pos;
|
||||
|
||||
/* if the buffer is too small for the entire request: realloc */
|
||||
if (sdslen(r->io.fd.buf) + sdsavail(r->io.fd.buf) < len)
|
||||
r->io.fd.buf = sdsMakeRoomFor(r->io.fd.buf, len - sdslen(r->io.fd.buf));
|
||||
|
||||
/* if the remaining unused buffer is not large enough: memmove so that we can read the rest */
|
||||
if (len > avail && sdsavail(r->io.fd.buf) < len - avail) {
|
||||
sdsrange(r->io.fd.buf, r->io.fd.pos, -1);
|
||||
r->io.fd.pos = 0;
|
||||
}
|
||||
|
||||
/* if we don't already have all the data in the sds, read more */
|
||||
while (len > sdslen(r->io.fd.buf) - r->io.fd.pos) {
|
||||
size_t toread = len - (sdslen(r->io.fd.buf) - r->io.fd.pos);
|
||||
/* read either what's missing, or REDIS_IOBUF_LEN, the bigger of the two */
|
||||
if (toread < REDIS_IOBUF_LEN)
|
||||
toread = REDIS_IOBUF_LEN;
|
||||
if (toread > sdsavail(r->io.fd.buf))
|
||||
toread = sdsavail(r->io.fd.buf);
|
||||
int retval = read(r->io.fd.fd, (char*)r->io.fd.buf + sdslen(r->io.fd.buf), toread);
|
||||
if (retval <= 0) {
|
||||
if (errno == EWOULDBLOCK) errno = ETIMEDOUT;
|
||||
return 0;
|
||||
}
|
||||
sdsIncrLen(r->io.fd.buf, retval);
|
||||
}
|
||||
|
||||
memcpy(buf, (char*)r->io.fd.buf + r->io.fd.pos, len);
|
||||
r->io.fd.pos += len;
|
||||
return len;
|
||||
}
|
||||
|
||||
/* Returns read/write position in file. */
|
||||
static off_t rioFdTell(rio *r) {
|
||||
off_t pos = lseek(r->io.fd.fd, 0, SEEK_CUR);
|
||||
return pos - sdslen(r->io.fd.buf) + r->io.fd.pos;
|
||||
}
|
||||
|
||||
/* Flushes any buffer to target device if applicable. Returns 1 on success
|
||||
* and 0 on failures. */
|
||||
static int rioFdFlush(rio *r) {
|
||||
/* Our flush is implemented by the write method, that recognizes a
|
||||
* buffer set to NULL with a count of zero as a flush request. */
|
||||
return rioFdWrite(r,NULL,0);
|
||||
}
|
||||
|
||||
static const rio rioFdIO = {
|
||||
rioFdRead,
|
||||
rioFdWrite,
|
||||
rioFdTell,
|
||||
rioFdFlush,
|
||||
NULL, /* update_checksum */
|
||||
0, /* current checksum */
|
||||
0, /* bytes read or written */
|
||||
0, /* read/write chunk size */
|
||||
{ { NULL, 0 } } /* union for io-specific vars */
|
||||
};
|
||||
|
||||
void rioInitWithFd(rio *r, int fd) {
|
||||
*r = rioFdIO;
|
||||
r->io.fd.fd = fd;
|
||||
r->io.fd.pos = 0;
|
||||
r->io.fd.buf = sdsnewlen(NULL, REDIS_IOBUF_LEN);
|
||||
sdsclear(r->io.fd.buf);
|
||||
}
|
||||
|
||||
/* release the rio stream.
|
||||
* optionally returns the unread buffered data. */
|
||||
void rioFreeFd(rio *r, sds* out_remainingBufferedData) {
|
||||
if(out_remainingBufferedData && (size_t)r->io.fd.pos < sdslen(r->io.fd.buf)) {
|
||||
if (r->io.fd.pos > 0)
|
||||
sdsrange(r->io.fd.buf, r->io.fd.pos, -1);
|
||||
*out_remainingBufferedData = r->io.fd.buf;
|
||||
} else {
|
||||
sdsfree(r->io.fd.buf);
|
||||
if (out_remainingBufferedData)
|
||||
*out_remainingBufferedData = NULL;
|
||||
}
|
||||
r->io.fd.buf = NULL;
|
||||
}
|
||||
|
||||
/* ------------------- File descriptors set implementation ------------------- */
|
||||
|
||||
/* Returns 1 or 0 for success/failure.
|
||||
* The function returns success as long as we are able to correctly write
|
||||
* to at least one file descriptor.
|
||||
*
|
||||
* When buf is NULL adn len is 0, the function performs a flush operation
|
||||
* When buf is NULL and len is 0, the function performs a flush operation
|
||||
* if there is some pending buffer, so this function is also used in order
|
||||
* to implement rioFdsetFlush(). */
|
||||
static size_t rioFdsetWrite(rio *r, const void *buf, size_t len) {
|
||||
@@ -176,7 +269,7 @@ static size_t rioFdsetWrite(rio *r, const void *buf, size_t len) {
|
||||
* a given size, we actually write to the sockets. */
|
||||
if (len) {
|
||||
r->io.fdset.buf = sdscatlen(r->io.fdset.buf,buf,len);
|
||||
len = 0; /* Prevent entering the while belove if we don't flush. */
|
||||
len = 0; /* Prevent entering the while below if we don't flush. */
|
||||
if (sdslen(r->io.fdset.buf) > REDIS_IOBUF_LEN) doflush = 1;
|
||||
}
|
||||
|
||||
@@ -276,6 +369,7 @@ void rioInitWithFdset(rio *r, int *fds, int numfds) {
|
||||
r->io.fdset.buf = sdsempty();
|
||||
}
|
||||
|
||||
/* release the rio stream. */
|
||||
void rioFreeFdset(rio *r) {
|
||||
zfree(r->io.fdset.fds);
|
||||
zfree(r->io.fdset.state);
|
||||
|
||||
@@ -73,6 +73,12 @@ struct _rio {
|
||||
off_t buffered; /* Bytes written since last fsync. */
|
||||
off_t autosync; /* fsync after 'autosync' bytes written. */
|
||||
} file;
|
||||
/* file descriptor */
|
||||
struct {
|
||||
int fd; /* File descriptor. */
|
||||
off_t pos;
|
||||
sds buf;
|
||||
} fd;
|
||||
/* Multiple FDs target (used to write to N sockets). */
|
||||
struct {
|
||||
int *fds; /* File descriptors. */
|
||||
@@ -126,8 +132,12 @@ static inline int rioFlush(rio *r) {
|
||||
|
||||
void rioInitWithFile(rio *r, FILE *fp);
|
||||
void rioInitWithBuffer(rio *r, sds s);
|
||||
void rioInitWithFd(rio *r, int fd);
|
||||
void rioInitWithFdset(rio *r, int *fds, int numfds);
|
||||
|
||||
void rioFreeFdset(rio *r);
|
||||
void rioFreeFd(rio *r, sds* out_remainingBufferedData);
|
||||
|
||||
size_t rioWriteBulkCount(rio *r, char prefix, int count);
|
||||
size_t rioWriteBulkString(rio *r, const char *buf, size_t len);
|
||||
size_t rioWriteBulkLongLong(rio *r, long long l);
|
||||
|
||||
@@ -1,12 +1,3 @@
|
||||
proc start_bg_complex_data {host port db ops} {
|
||||
set tclsh [info nameofexecutable]
|
||||
exec $tclsh tests/helpers/bg_complex_data.tcl $host $port $db $ops &
|
||||
}
|
||||
|
||||
proc stop_bg_complex_data {handle} {
|
||||
catch {exec /bin/kill -9 $handle}
|
||||
}
|
||||
|
||||
start_server {tags {"repl"}} {
|
||||
start_server {} {
|
||||
|
||||
|
||||
@@ -1,12 +1,3 @@
|
||||
proc start_bg_complex_data {host port db ops} {
|
||||
set tclsh [info nameofexecutable]
|
||||
exec $tclsh tests/helpers/bg_complex_data.tcl $host $port $db $ops &
|
||||
}
|
||||
|
||||
proc stop_bg_complex_data {handle} {
|
||||
catch {exec /bin/kill -9 $handle}
|
||||
}
|
||||
|
||||
# Creates a master-slave pair and breaks the link continuously to force
|
||||
# partial resyncs attempts, all this while flooding the master with
|
||||
# write queries.
|
||||
|
||||
@@ -130,85 +130,91 @@ start_server {tags {"repl"}} {
|
||||
}
|
||||
}
|
||||
|
||||
foreach dl {no yes} {
|
||||
start_server {tags {"repl"}} {
|
||||
set master [srv 0 client]
|
||||
$master config set repl-diskless-sync $dl
|
||||
set master_host [srv 0 host]
|
||||
set master_port [srv 0 port]
|
||||
set slaves {}
|
||||
set load_handle0 [start_write_load $master_host $master_port 3]
|
||||
set load_handle1 [start_write_load $master_host $master_port 5]
|
||||
set load_handle2 [start_write_load $master_host $master_port 20]
|
||||
set load_handle3 [start_write_load $master_host $master_port 8]
|
||||
set load_handle4 [start_write_load $master_host $master_port 4]
|
||||
start_server {} {
|
||||
lappend slaves [srv 0 client]
|
||||
foreach mdl {no yes} {
|
||||
foreach sdl {no yes} {
|
||||
start_server {tags {"repl"}} {
|
||||
set master [srv 0 client]
|
||||
$master config set repl-diskless-sync $mdl
|
||||
$master config set repl-diskless-sync-delay 1
|
||||
set master_host [srv 0 host]
|
||||
set master_port [srv 0 port]
|
||||
set slaves {}
|
||||
set load_handle0 [start_bg_complex_data $master_host $master_port 9 100000000]
|
||||
set load_handle1 [start_bg_complex_data $master_host $master_port 11 100000000]
|
||||
set load_handle2 [start_bg_complex_data $master_host $master_port 12 100000000]
|
||||
set load_handle3 [start_write_load $master_host $master_port 8]
|
||||
set load_handle4 [start_write_load $master_host $master_port 4]
|
||||
start_server {} {
|
||||
lappend slaves [srv 0 client]
|
||||
start_server {} {
|
||||
lappend slaves [srv 0 client]
|
||||
test "Connect multiple slaves at the same time (issue #141), diskless=$dl" {
|
||||
# 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
|
||||
start_server {} {
|
||||
lappend slaves [srv 0 client]
|
||||
test "Connect multiple slaves at the same time (issue #141), master diskless=$mdl, slave diskless=$sdl" {
|
||||
# Send SLAVEOF commands to slaves
|
||||
[lindex $slaves 0] config set repl-diskless-load $sdl
|
||||
[lindex $slaves 1] config set repl-diskless-load $sdl
|
||||
[lindex $slaves 2] config set repl-diskless-load $sdl
|
||||
[lindex $slaves 0] slaveof $master_host $master_port
|
||||
[lindex $slaves 1] slaveof $master_host $master_port
|
||||
[lindex $slaves 2] slaveof $master_host $master_port
|
||||
|
||||
# Wait for all the three slaves to reach the "online"
|
||||
# state from the POV of the master.
|
||||
set retry 500
|
||||
while {$retry} {
|
||||
set info [r -3 info]
|
||||
if {[string match {*slave0:*state=online*slave1:*state=online*slave2:*state=online*} $info]} {
|
||||
break
|
||||
} else {
|
||||
incr retry -1
|
||||
after 100
|
||||
# Wait for all the three slaves to reach the "online"
|
||||
# state from the POV of the master.
|
||||
set retry 500
|
||||
while {$retry} {
|
||||
set info [r -3 info]
|
||||
if {[string match {*slave0:*state=online*slave1:*state=online*slave2:*state=online*} $info]} {
|
||||
break
|
||||
} else {
|
||||
incr retry -1
|
||||
after 100
|
||||
}
|
||||
}
|
||||
if {$retry == 0} {
|
||||
error "assertion:Slaves not correctly synchronized"
|
||||
}
|
||||
}
|
||||
if {$retry == 0} {
|
||||
error "assertion:Slaves not correctly synchronized"
|
||||
}
|
||||
|
||||
# Wait that slaves acknowledge they are online so
|
||||
# we are sure that DBSIZE and DEBUG DIGEST will not
|
||||
# fail because of timing issues.
|
||||
wait_for_condition 500 100 {
|
||||
[lindex [[lindex $slaves 0] role] 3] eq {connected} &&
|
||||
[lindex [[lindex $slaves 1] role] 3] eq {connected} &&
|
||||
[lindex [[lindex $slaves 2] role] 3] eq {connected}
|
||||
} else {
|
||||
fail "Slaves still not connected after some time"
|
||||
# Wait that slaves acknowledge they are online so
|
||||
# we are sure that DBSIZE and DEBUG DIGEST will not
|
||||
# fail because of timing issues.
|
||||
wait_for_condition 500 100 {
|
||||
[lindex [[lindex $slaves 0] role] 3] eq {connected} &&
|
||||
[lindex [[lindex $slaves 1] role] 3] eq {connected} &&
|
||||
[lindex [[lindex $slaves 2] role] 3] eq {connected}
|
||||
} else {
|
||||
fail "Slaves still not connected after some time"
|
||||
}
|
||||
|
||||
# Stop the write load
|
||||
stop_bg_complex_data $load_handle0
|
||||
stop_bg_complex_data $load_handle1
|
||||
stop_bg_complex_data $load_handle2
|
||||
stop_write_load $load_handle3
|
||||
stop_write_load $load_handle4
|
||||
|
||||
# Make sure that slaves and master have same
|
||||
# number of keys
|
||||
wait_for_condition 500 100 {
|
||||
[$master dbsize] == [[lindex $slaves 0] dbsize] &&
|
||||
[$master dbsize] == [[lindex $slaves 1] dbsize] &&
|
||||
[$master dbsize] == [[lindex $slaves 2] dbsize]
|
||||
} else {
|
||||
fail "Different number of keys between masted and slave after too long time."
|
||||
}
|
||||
|
||||
# Check digests
|
||||
set digest [$master debug digest]
|
||||
set digest0 [[lindex $slaves 0] debug digest]
|
||||
set digest1 [[lindex $slaves 1] debug digest]
|
||||
set digest2 [[lindex $slaves 2] debug digest]
|
||||
assert {$digest ne 0000000000000000000000000000000000000000}
|
||||
assert {$digest eq $digest0}
|
||||
assert {$digest eq $digest1}
|
||||
assert {$digest eq $digest2}
|
||||
}
|
||||
|
||||
# Stop the write load
|
||||
stop_write_load $load_handle0
|
||||
stop_write_load $load_handle1
|
||||
stop_write_load $load_handle2
|
||||
stop_write_load $load_handle3
|
||||
stop_write_load $load_handle4
|
||||
|
||||
# Make sure that slaves and master have same
|
||||
# number of keys
|
||||
wait_for_condition 500 100 {
|
||||
[$master dbsize] == [[lindex $slaves 0] dbsize] &&
|
||||
[$master dbsize] == [[lindex $slaves 1] dbsize] &&
|
||||
[$master dbsize] == [[lindex $slaves 2] dbsize]
|
||||
} else {
|
||||
fail "Different number of keys between masted and slave after too long time."
|
||||
}
|
||||
|
||||
# Check digests
|
||||
set digest [$master debug digest]
|
||||
set digest0 [[lindex $slaves 0] debug digest]
|
||||
set digest1 [[lindex $slaves 1] debug digest]
|
||||
set digest2 [[lindex $slaves 2] debug digest]
|
||||
assert {$digest ne 0000000000000000000000000000000000000000}
|
||||
assert {$digest eq $digest0}
|
||||
assert {$digest eq $digest1}
|
||||
assert {$digest eq $digest2}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -371,3 +371,15 @@ proc start_write_load {host port seconds} {
|
||||
proc stop_write_load {handle} {
|
||||
catch {exec /bin/kill -9 $handle}
|
||||
}
|
||||
|
||||
# Execute a background process writing complex data for the specified number
|
||||
# of ops to the specified Redis instance.
|
||||
proc start_bg_complex_data {host port db ops} {
|
||||
set tclsh [info nameofexecutable]
|
||||
exec $tclsh tests/helpers/bg_complex_data.tcl $host $port $db $ops &
|
||||
}
|
||||
|
||||
# Stop a process generating write load executed with start_bg_complex_data.
|
||||
proc stop_bg_complex_data {handle} {
|
||||
catch {exec /bin/kill -9 $handle}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user