Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion deps/xredis-gtid
1 change: 1 addition & 0 deletions src/config.c
Original file line number Diff line number Diff line change
Expand Up @@ -3184,6 +3184,7 @@ standardConfig configs[] = {
createBoolConfig("swap-rdb-bitmap-encode-enabled", NULL, MODIFIABLE_CONFIG, server.swap_rdb_bitmap_encode_enabled, 1, NULL, NULL),
createBoolConfig("swap-bitmap-subkeys-enabled", NULL, MODIFIABLE_CONFIG, server.swap_bitmap_subkeys_enabled, 1, NULL, NULL),
createBoolConfig("swap-ttl-compact-enabled", NULL, MODIFIABLE_CONFIG, server.swap_ttl_compact_enabled, 1, NULL, NULL),
createBoolConfig("swap-soft-block-cmd-enabled", NULL, MODIFIABLE_CONFIG, server.swap_soft_block_cmd_enabled, 0, NULL, NULL),
createBoolConfig("rocksdb.data.cache_index_and_filter_blocks", "rocksdb.cache_index_and_filter_blocks", IMMUTABLE_CONFIG, server.rocksdb_data_cache_index_and_filter_blocks, 0, NULL, NULL),
createBoolConfig("rocksdb.meta.cache_index_and_filter_blocks", NULL, IMMUTABLE_CONFIG, server.rocksdb_meta_cache_index_and_filter_blocks, 0, NULL, NULL),
createBoolConfig("rocksdb.enable_pipelined_write", NULL, IMMUTABLE_CONFIG, server.rocksdb_enable_pipelined_write, 0, NULL, NULL),
Expand Down
4 changes: 4 additions & 0 deletions src/ctrip_swap.h
Original file line number Diff line number Diff line change
Expand Up @@ -2089,6 +2089,10 @@ void swapApplySwapInfo(int swap_info_argc, sds *swap_info_argv);
int submitReplClientRequests(client *c);
sds genSwapReplInfoString(sds info);

/* Command */
sds genSwapCommandInfoString(sds info);
int swapCmdShouldBlock(client *c, struct redisCommand *cmd);

/* Swap */
void swapInit(void);
int dbSwap(client *c);
Expand Down
69 changes: 48 additions & 21 deletions src/ctrip_swap_cmd.c
Original file line number Diff line number Diff line change
Expand Up @@ -870,64 +870,64 @@ struct redisCommand redisCommandTable[SWAP_CMD_COUNT] = {
0,NULL,NULL,SWAP_IN,0,2,2,1,0,0,0},

{"xadd",xaddCommand,-5,
"write use-memory fast random @stream",
0,NULL,NULL,SWAP_IN,0,1,1,1,0,0,0},
"write use-memory fast random swap-soft-blocked @stream",
0,NULL,NULL,SWAP_NOP,0,1,1,1,0,0,0},

{"xrange",xrangeCommand,-4,
"read-only @stream",
0,NULL,NULL,SWAP_IN,0,1,1,1,0,0,0},
0,NULL,NULL,SWAP_NOP,0,1,1,1,0,0,0},

{"xrevrange",xrevrangeCommand,-4,
"read-only @stream",
0,NULL,NULL,SWAP_IN,0,1,1,1,0,0,0},
0,NULL,NULL,SWAP_NOP,0,1,1,1,0,0,0},

{"xlen",xlenCommand,2,
"read-only fast @stream",
0,NULL,NULL,SWAP_IN,0,1,1,1,0,0,0},
0,NULL,NULL,SWAP_NOP,0,1,1,1,0,0,0},

{"xread",xreadCommand,-4,
"read-only @stream @blocking",
0,xreadGetKeys,NULL,SWAP_IN,0,0,0,0,0,0,0},
0,xreadGetKeys,NULL,SWAP_NOP,0,0,0,0,0,0,0},

{"xreadgroup",xreadCommand,-7,
"write @stream @blocking",
0,xreadGetKeys,NULL,SWAP_IN,0,0,0,0,0,0,0},
"write swap-soft-blocked @stream @blocking",
0,xreadGetKeys,NULL,SWAP_NOP,0,0,0,0,0,0,0},

{"xgroup",xgroupCommand,-2,
"write use-memory @stream",
"write use-memory swap-soft-blocked @stream",
0,NULL,NULL,SWAP_NOP,0,2,2,1,0,0,0},

{"xsetid",xsetidCommand,3,
"write use-memory fast @stream",
"write use-memory fast swap-soft-blocked @stream",
0,NULL,NULL,SWAP_NOP,0,1,1,1,0,0,0},

{"xack",xackCommand,-4,
"write fast random @stream",
"write fast random swap-soft-blocked @stream",
0,NULL,NULL,SWAP_NOP,0,1,1,1,0,0,0},

{"xpending",xpendingCommand,-3,
"read-only random @stream",
0,NULL,NULL,SWAP_IN,0,1,1,1,0,0,0},
0,NULL,NULL,SWAP_NOP,0,1,1,1,0,0,0},

{"xclaim",xclaimCommand,-6,
"write random fast @stream",
0,NULL,NULL,SWAP_IN,0,1,1,1,0,0,0},
"write random fast swap-soft-blocked @stream",
0,NULL,NULL,SWAP_NOP,0,1,1,1,0,0,0},

{"xautoclaim",xautoclaimCommand,-6,
"write random fast @stream",
0,NULL,NULL,SWAP_IN,0,1,1,1,0,0,0},
"write random fast swap-soft-blocked @stream",
0,NULL,NULL,SWAP_NOP,0,1,1,1,0,0,0},

{"xinfo",xinfoCommand,-2,
"read-only random @stream",
0,NULL,NULL,SWAP_IN,0,2,2,1,0,0,0},
0,NULL,NULL,SWAP_NOP,0,2,2,1,0,0,0},

{"xdel",xdelCommand,-3,
"write fast @stream",
0,NULL,NULL,SWAP_IN,0,1,1,1,0,0,0},
"write fast swap-soft-blocked @stream",
0,NULL,NULL,SWAP_NOP,0,1,1,1,0,0,0},

{"xtrim",xtrimCommand,-4,
"write random @stream",
0,NULL,NULL,SWAP_IN,0,1,1,1,0,0,0},
"write random swap-soft-blocked @stream",
0,NULL,NULL,SWAP_NOP,0,1,1,1,0,0,0},

{"post",securityWarningCommand,-1,
"ok-loading ok-stale read-only",
Expand Down Expand Up @@ -2312,6 +2312,33 @@ void swapInfoCommand(client *c) {
addReply(c,shared.ok);
}

/* Return 1 when cmd must be rejected. The swap check lives inside so
* callers do not wrap the call in ENABLE_SWAP. */
int swapCmdShouldBlock(client *c, struct redisCommand *cmd) {

if (!cmd ||
!(cmd->flags & (CMD_SWAP_HARD_BLOCKED|CMD_SWAP_SOFT_BLOCKED)))
return 0;

if (c && (c->flags & CLIENT_MASTER)) {
server.swap_master_allowed_cmd_count++;
return 0;
}
/* AOF replay applies commands that were already accepted. */
if (c && c->id == CLIENT_ID_AOF)
return 0;

if (cmd->flags & CMD_SWAP_HARD_BLOCKED) {
server.swap_hard_blocked_cmd_count++;
return 1;
}
if (server.swap_soft_block_cmd_enabled) {
server.swap_soft_blocked_cmd_count++;
return 1;
}
return 0;
}

#ifdef REDIS_TEST

void rewriteResetClientCommandCString(client *c, int argc, ...) {
Expand Down
3 changes: 3 additions & 0 deletions src/ctrip_swap_server.c
Original file line number Diff line number Diff line change
Expand Up @@ -244,6 +244,9 @@ void swapInitServer(void) {
server.rocksdb_rdb_checkpoint_dir = NULL;
server.rocksdb_internal_stats = NULL;
server.swap_util_task_manager = createRocksdbUtilTaskManager();
server.swap_hard_blocked_cmd_count = 0;
server.swap_soft_blocked_cmd_count = 0;
server.swap_master_allowed_cmd_count = 0;

asyncCompleteQueueInit();
parallelSyncInit(server.swap_ps_parallism_rdb);
Expand Down
9 changes: 9 additions & 0 deletions src/ctrip_swap_server.h
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,11 @@
#define CMD_SWAP_DATATYPE_LIST (1ULL<<46)
#define CMD_SWAP_DATATYPE_BITMAP (1ULL<<47)

/*cmd flag*/
#define CMD_SWAP_HARD_BLOCKED (1ULL<<48)
#define CMD_SWAP_SOFT_BLOCKED (1ULL<<49)


/* CHECK: CLIENT_REPL_RDBONLY is the last CLIENT_xx flag */
#define CLIENT_SWAPPING (1ULL<<43) /* The client is waiting swap. */
#define CLIENT_SWAP_UNLOCKING (1ULL<<44) /* Client is releasing swap lock. */
Expand Down Expand Up @@ -346,6 +351,10 @@ typedef struct swapBatchLimitsConfig {
/* for swap.info command, which propagate system info to replica */ \
int swap_swap_info_supported; \
int swap_swap_info_propagate_mode; \
int swap_soft_block_cmd_enabled; \
long long swap_hard_blocked_cmd_count; \
long long swap_soft_blocked_cmd_count; \
long long swap_master_allowed_cmd_count; \
unsigned long long swap_swap_info_slave_period; /* Master send cmd swap.info to the slave every N seconds */

#ifdef __APPLE__
Expand Down
15 changes: 15 additions & 0 deletions src/ctrip_swap_stat.c
Original file line number Diff line number Diff line change
Expand Up @@ -348,6 +348,20 @@ void trackSwapInstantaneousMetrics() {
trackSwapRateLimitInstantaneousMetrics();
}


sds genSwapCommandInfoString(sds info) {

info = sdscatprintf(info,
"swap_cmd_stats:"
"hard_blocked=%lld,"
"soft_blocked=%lld,"
"master_allowed=%lld\r\n",
server.swap_hard_blocked_cmd_count,
server.swap_soft_blocked_cmd_count,
server.swap_master_allowed_cmd_count);
return info;
}

sds genSwapInfoString(sds info) {
info = genSwapStorageInfoString(info);
info = genSwapHitInfoString(info);
Expand All @@ -365,6 +379,7 @@ sds genSwapInfoString(sds info) {
info = genSwapBitmapStringSwitchedInfoString(info);
info = genSwapTtlCompactInfoString(info);
info = genSwapFullCompactInfoString(info);
info = genSwapCommandInfoString(info);
return info;
}

Expand Down
4 changes: 4 additions & 0 deletions src/module.c
Original file line number Diff line number Diff line change
Expand Up @@ -811,6 +811,10 @@ int64_t commandFlagsFromString(char *s) {
else if (!strcasecmp(t,"getkeys-api")) flags |= CMD_MODULE_GETKEYS;
else if (!strcasecmp(t,"no-cluster")) flags |= CMD_MODULE_NO_CLUSTER;
else if (!strcasecmp(t,"gtid-non-determinism")) flags |= CMD_GTID_NON_DETERMINISM;
#ifdef ENABLE_SWAP
else if (!strcasecmp(t, "swap-hard-blocked")) flags |= CMD_SWAP_HARD_BLOCKED;
else if (!strcasecmp(t, "swap-soft-blocked")) flags |= CMD_SWAP_SOFT_BLOCKED;
#endif
else break;
}
sdsfreesplitres(tokens,count);
Expand Down
8 changes: 8 additions & 0 deletions src/scripting.c
Original file line number Diff line number Diff line change
Expand Up @@ -633,6 +633,14 @@ int luaRedisGenericCommand(lua_State *lua, int raise_error) {
goto cleanup;
}

#ifdef ENABLE_SWAP
if (swapCmdShouldBlock(server.lua_caller, cmd)) {
luaPushError(lua,
"Can't execute this command: data type is not supported by swap");
goto cleanup;
}
#endif

/* Check the ACLs. */
int acl_errpos;
int acl_retval = ACLCheckAllPerm(c,&acl_errpos);
Expand Down
14 changes: 14 additions & 0 deletions src/server.c
Original file line number Diff line number Diff line change
Expand Up @@ -3782,6 +3782,12 @@ int populateCommandTableParseFlags(struct redisCommand *c, char *strflags) {
c->flags |= CMD_MAY_REPLICATE;
} else if (!strcasecmp(flag, "gtid-non-determinism")) {
c->flags |= CMD_GTID_NON_DETERMINISM;
#ifdef ENABLE_SWAP
} else if (!strcasecmp(flag, "swap-hard-blocked")) {
c->flags |= CMD_SWAP_HARD_BLOCKED;
} else if (!strcasecmp(flag, "swap-soft-blocked")) {
c->flags |= CMD_SWAP_SOFT_BLOCKED;
#endif
} else {
/* Parse ACL categories here if the flag name starts with @. */
uint64_t catflag;
Expand Down Expand Up @@ -4571,6 +4577,14 @@ int processCommand(client *c) {
rejectCommandFormat(c, "Previous master draining.");
return C_OK;
}


if (swapCmdShouldBlock(c, c->cmd)) {
rejectCommandFormat(c,
"-ERR Can't execute '%s' command: data type is not supported by swap",
c->cmd->name);
return C_OK;
}
#endif

/* Only allow a subset of commands in the context of Pub/Sub if the
Expand Down
4 changes: 4 additions & 0 deletions src/server.h
Original file line number Diff line number Diff line change
Expand Up @@ -2369,6 +2369,10 @@ int getMaxmemoryState(size_t *total, size_t *logical, size_t *tofree, float *lev
size_t freeMemoryGetNotCountedMemory();
int overMaxmemoryAfterAlloc(size_t moremem);
int processCommand(client *c);
/* Return 1 when cmd must be rejected. Updates swap_cmd_stats when swap is
* enabled. caller is the originating client. Lua passes server.lua_caller
* because the fake Lua client does not carry CLIENT_MASTER. */
int swapCmdShouldBlock(client *caller, struct redisCommand *cmd);
int processPendingCommandsAndResetClient(client *c);
void setupSignalHandlers(void);
void removeSignalHandlers(void);
Expand Down
Loading
Loading