Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
01888d1e58 | ||
|
|
30278061cc | ||
|
|
49efe300af | ||
|
|
702f914771 | ||
|
|
e0b2e24830 | ||
|
|
0da453160b | ||
|
|
1d8973c47d | ||
|
|
df7add9e70 | ||
|
|
9003483d43 | ||
|
|
a13d6378c1 | ||
|
|
fa17e2daf0 | ||
|
|
ff7c1faa12 | ||
|
|
e252e9f231 | ||
|
|
9c0a68861e | ||
|
|
5844f5d0d1 | ||
|
|
a97658f293 | ||
|
|
138e7af57b | ||
|
|
3a9f41ad86 | ||
|
|
71fba427c2 | ||
|
|
08d4df8d31 | ||
|
|
e7422ef166 | ||
|
|
362032e43a | ||
|
|
10323dc5fe | ||
|
|
e213c408fa | ||
|
|
5674656db7 | ||
|
|
35d71b1ffc | ||
|
|
6637862838 | ||
|
|
b065f4441b | ||
|
|
bd99b26bc5 | ||
|
|
0560738f6b | ||
|
|
88d58661db | ||
|
|
315e3b14ef | ||
|
|
7ff051f6c1 | ||
|
|
f387a5acf8 | ||
|
|
1fab07e078 | ||
|
|
8ebae5d630 | ||
|
|
21c3d77118 | ||
|
|
60a28fad8a | ||
|
|
e42baed4c3 | ||
|
|
aa67aec84e | ||
|
|
93959bc09f | ||
|
|
2e92d0f04a | ||
|
|
adcb470130 | ||
|
|
1b71fea998 | ||
|
|
2b5cf6bf78 | ||
|
|
7e78ab4b6f | ||
|
|
3468cd3664 | ||
|
|
d1b5c5defd | ||
|
|
66899a42fc | ||
|
|
76b18c7a0e | ||
|
|
ca804a1022 | ||
|
|
d15d9fecd2 | ||
|
|
c2717911db | ||
|
|
1641f41cfc | ||
|
|
b37b2b5c14 | ||
|
|
c43c970344 | ||
|
|
b64c861171 | ||
|
|
47bbaa17b0 | ||
|
|
0595420b1e | ||
|
|
2d34ec60bf | ||
|
|
2d7d75adb3 |
+107
-15
@@ -1,8 +1,6 @@
|
||||
Redis 3.0 release notes
|
||||
=======================
|
||||
|
||||
WARNING: Redis 3.0 is currently a BETA not suitable for production environments.
|
||||
|
||||
--------------------------------------------------------------------------------
|
||||
Upgrade urgency levels:
|
||||
|
||||
@@ -12,6 +10,111 @@ HIGH: There is a critical bug that may affect a subset of users. Upgrade!
|
||||
CRITICAL: There is a critical bug affecting MOST USERS. Upgrade ASAP.
|
||||
--------------------------------------------------------------------------------
|
||||
|
||||
--[ Redis 3.0.2 ] Release date: 4 Jun 2015
|
||||
|
||||
Upgrade urgency: HIGH for Redis because of a security issue.
|
||||
LOW for Sentinel.
|
||||
|
||||
* [FIX] Critical security issue fix by Ben Murphy: http://t.co/LpGTyZmfS7
|
||||
* [FIX] SMOVE reply fixed when src and dst keys are the same. (Glenn Nethercutt)
|
||||
* [FIX] Lua cmsgpack lib updated to support str8 type. (Sebastian Waisbrot)
|
||||
|
||||
* [NEW] ZADD support for options: NX, XX, CH. See new doc at redis.io.
|
||||
(Salvatore Sanfilippo)
|
||||
* [NEW] Senitnel: CKQUORUM and FLUSHCONFIG commands back ported.
|
||||
(Salvatore Sanfilippo)
|
||||
|
||||
--[ Redis 3.0.1 ] Release date: 5 May 2015
|
||||
|
||||
Upgrade urgency: LOW for Redis and Cluster, MODERATE for Sentinel.
|
||||
|
||||
* [FIX] Sentinel memory leak due to hiredis fixed. (Salvatore Sanfilippo)
|
||||
* [FIX] Sentinel memory leak on duplicated instance. (Charsyam)
|
||||
* [FIX] Redis crash on Lua reaching output buffer limits. (Yossi Gottlieb)
|
||||
* [FIX] Sentinel flushes config on +slave events. (Bill Anderson)
|
||||
|
||||
--[ Redis 3.0.0 ] Release date: 1 Apr 2015
|
||||
|
||||
>> What's new in Redis 3.0 compared to Redis 2.8?
|
||||
|
||||
* Redis Cluster: a distributed implementation of a subset of Redis.
|
||||
* New "embedded string" object encoding resulting in less cache
|
||||
misses. Big speed gain under certain work loads.
|
||||
* AOF child -> parent final data transmission to minimize latency due
|
||||
to "last write" during AOF rewrites.
|
||||
* Much improved LRU approximation algorithm for keys eviction.
|
||||
* WAIT command to block waiting for a write to be transmitted to
|
||||
the specified number of slaves.
|
||||
* MIGRATE connection caching. Much faster keys migraitons.
|
||||
* MIGRATE new options COPY and REPLACE.
|
||||
* CLIENT PAUSE command: stop processing client requests for a
|
||||
specified amount of time.
|
||||
* BITCOUNT performance improvements.
|
||||
* CONFIG SET accepts memory values in different units (for example
|
||||
you can use "CONFIG SET maxmemory 1gb").
|
||||
* Redis log format slightly changed reporting in each line the role of the
|
||||
instance (master/slave) or if it's a saving child log.
|
||||
* INCR performance improvements.
|
||||
|
||||
>> Refactoring changes (no new features nor bug fixes)
|
||||
|
||||
* Blocking operations full refactoring (blocked.c)
|
||||
* Client output buffer memory tracking refactored.
|
||||
|
||||
Changes between RC6 and 3.0.0 stable:
|
||||
|
||||
>> General changes
|
||||
|
||||
* Fixes to diskless replication. (Oran Agra)
|
||||
* Test for BLPOP replication on role change. (Salvatore Sanfilippo)
|
||||
* prepareClientToWrite() error handling improvements. (Salvatore Sanfilippo)
|
||||
* Remove dict.c no longer used function. (Salvatore Sanfilippo)
|
||||
|
||||
>> Cluster changes
|
||||
|
||||
None
|
||||
|
||||
>> Sentinel changes
|
||||
|
||||
None
|
||||
|
||||
--[ Redis 3.0.0 RC6 (version 2.9.106) ] Release date: 24 mar 2015
|
||||
|
||||
Upgrade urgency: HIGH because of bugs related to Redis Custer and replication.
|
||||
|
||||
This is the 6th release candidate of Redis 3.0.0. This release fixes important
|
||||
issues discovered during stress testing, and implements safest behavior
|
||||
for blocking operations during clients reshardings, and a new much needed
|
||||
functionality of Redis Cluster manual failovers.
|
||||
|
||||
In order to fix certain bugs quite a bit of refactoring was needed which
|
||||
is usually non advisabble in a Release Candidate, but needed in order to
|
||||
end with a clean fix.
|
||||
|
||||
>> General changes
|
||||
|
||||
* [FIX] Redis (non clustered & clustered) replication bug involving blocking
|
||||
operations: see issue #2473. (Salvatore Sanfilippo)
|
||||
|
||||
>> Cluster changes
|
||||
|
||||
* [FIX] clientsArePaused() fix crashing the old master during manual failover.
|
||||
(Salvatore Sanfilippo)
|
||||
* [FIX] Lua scripts replication in Redis Cluster was totally broken.
|
||||
(Salvatore Sanfilippo)
|
||||
* [FIX] Redirect clients blocked into list operations when the hash slot
|
||||
they are blocked into is migrated to another instance or the cluster
|
||||
state turns into "fail". (Salvatore Sanfilippo)
|
||||
|
||||
* [NEW] TAKEOVER option for CLUSTER FAILOVER implemented. It is now possible
|
||||
to fix a cluster manually in the minority side of the partition, for
|
||||
example in order to allow for multi DC setups & recovery.
|
||||
(Salvatore Sanfilippo)
|
||||
|
||||
>> Sentinel changes
|
||||
|
||||
No changes in Sentinel.
|
||||
|
||||
--[ Redis 3.0.0 RC5 (version 2.9.105) ] Release date: 20 mar 2015
|
||||
|
||||
Upgrade urgency: Moderate for Redis Cluster users, low otherwise.
|
||||
@@ -29,6 +132,7 @@ process of finishing the documentation for Redis Cluster).
|
||||
* [FIX] Fix for backtrace generation issue. (Mariano Pérez Rodríguez, Matt Stancliff, Salvatore Sanfilippo)
|
||||
|
||||
* [NEW] Redis-cli --latency-dist backported from unstable.
|
||||
(Salvatore Sanfilippo)
|
||||
|
||||
>> Cluster changes
|
||||
|
||||
@@ -492,22 +596,10 @@ This is the second beta of Redis 3.0.0.
|
||||
|
||||
This is the first beta of Redis 3.0.0.
|
||||
|
||||
The following is a list of improvements in Redis 3.0, compared to Redis 2.8.
|
||||
|
||||
* [NEW] Redis Cluster: a distributed implementation of a subset of Redis.
|
||||
* [NEW] New "embedded string" object encoding resulting in less cache
|
||||
misses. Big speed gain under certain work loads.
|
||||
* [NEW] WAIT command to block waiting for a write to be transmitted to
|
||||
the specified number of slaves.
|
||||
* [NEW] MIGRATE connection caching. Much faster keys migraitons.
|
||||
* [NEW] MIGARTE new options COPY and REPLACE.
|
||||
* [NEW] CLIENT PAUSE command: stop processing client requests for a
|
||||
specified amount of time.
|
||||
|
||||
Migrating from 2.8 to 3.0
|
||||
=========================
|
||||
|
||||
Redis 3.0 is mostly a strict subset of 2.8, you should not have any problem
|
||||
Redis 2.8 is mostly a strict subset of 3.0, you should not have any problem
|
||||
upgrading your application from 2.8 to 3.0. However this is a list of small
|
||||
non-backward compatible changes introduced in the 3.0 release:
|
||||
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
Copyright (c) 2006-2014, Salvatore Sanfilippo
|
||||
Copyright (c) 2006-2015, Salvatore Sanfilippo
|
||||
All rights reserved.
|
||||
|
||||
Redistribution and use in source and binary forms, with or without modification, are permitted provided that the following conditions are met:
|
||||
|
||||
Vendored
+1
@@ -443,6 +443,7 @@ void redisProcessCallbacks(redisAsyncContext *ac) {
|
||||
if (((redisReply*)reply)->type == REDIS_REPLY_ERROR) {
|
||||
c->err = REDIS_ERR_OTHER;
|
||||
snprintf(c->errstr,sizeof(c->errstr),"%s",((redisReply*)reply)->str);
|
||||
c->reader->fn->freeObject(reply);
|
||||
__redisAsyncDisconnect(ac);
|
||||
return;
|
||||
}
|
||||
|
||||
Vendored
+1
-1
@@ -495,7 +495,7 @@ static void f_parser (lua_State *L, void *ud) {
|
||||
struct SParser *p = cast(struct SParser *, ud);
|
||||
int c = luaZ_lookahead(p->z);
|
||||
luaC_checkGC(L);
|
||||
tf = ((c == LUA_SIGNATURE[0]) ? luaU_undump : luaY_parser)(L, p->z,
|
||||
tf = (luaY_parser)(L, p->z,
|
||||
&p->buff, p->name);
|
||||
cl = luaF_newLclosure(L, tf->nups, hvalue(gt(L)));
|
||||
cl->l.p = tf;
|
||||
|
||||
Vendored
+42
-29
@@ -66,7 +66,7 @@
|
||||
/* Reverse memory bytes if arch is little endian. Given the conceptual
|
||||
* simplicity of the Lua build system we prefer check for endianess at runtime.
|
||||
* The performance difference should be acceptable. */
|
||||
static void memrevifle(void *ptr, size_t len) {
|
||||
void memrevifle(void *ptr, size_t len) {
|
||||
unsigned char *p = (unsigned char *)ptr,
|
||||
*e = (unsigned char *)p+len-1,
|
||||
aux;
|
||||
@@ -96,7 +96,7 @@ typedef struct mp_buf {
|
||||
size_t len, free;
|
||||
} mp_buf;
|
||||
|
||||
static void *mp_realloc(lua_State *L, void *target, size_t osize,size_t nsize) {
|
||||
void *mp_realloc(lua_State *L, void *target, size_t osize,size_t nsize) {
|
||||
void *(*local_realloc) (void *, void *, size_t osize, size_t nsize) = NULL;
|
||||
void *ud;
|
||||
|
||||
@@ -105,7 +105,7 @@ static void *mp_realloc(lua_State *L, void *target, size_t osize,size_t nsize) {
|
||||
return local_realloc(ud, target, osize, nsize);
|
||||
}
|
||||
|
||||
static mp_buf *mp_buf_new(lua_State *L) {
|
||||
mp_buf *mp_buf_new(lua_State *L) {
|
||||
mp_buf *buf = NULL;
|
||||
|
||||
/* Old size = 0; new size = sizeof(*buf) */
|
||||
@@ -117,7 +117,7 @@ static mp_buf *mp_buf_new(lua_State *L) {
|
||||
return buf;
|
||||
}
|
||||
|
||||
static void mp_buf_append(mp_buf *buf, const unsigned char *s, size_t len) {
|
||||
void mp_buf_append(mp_buf *buf, const unsigned char *s, size_t len) {
|
||||
if (buf->free < len) {
|
||||
size_t newlen = buf->len+len;
|
||||
|
||||
@@ -153,7 +153,7 @@ typedef struct mp_cur {
|
||||
int err;
|
||||
} mp_cur;
|
||||
|
||||
static void mp_cur_init(mp_cur *cursor, const unsigned char *s, size_t len) {
|
||||
void mp_cur_init(mp_cur *cursor, const unsigned char *s, size_t len) {
|
||||
cursor->p = s;
|
||||
cursor->left = len;
|
||||
cursor->err = MP_CUR_ERROR_NONE;
|
||||
@@ -173,13 +173,17 @@ static void mp_cur_init(mp_cur *cursor, const unsigned char *s, size_t len) {
|
||||
|
||||
/* ------------------------- Low level MP encoding -------------------------- */
|
||||
|
||||
static void mp_encode_bytes(mp_buf *buf, const unsigned char *s, size_t len) {
|
||||
void mp_encode_bytes(mp_buf *buf, const unsigned char *s, size_t len) {
|
||||
unsigned char hdr[5];
|
||||
int hdrlen;
|
||||
|
||||
if (len < 32) {
|
||||
hdr[0] = 0xa0 | (len&0xff); /* fix raw */
|
||||
hdrlen = 1;
|
||||
} else if (len <= 0xff) {
|
||||
hdr[0] = 0xd9;
|
||||
hdr[1] = len;
|
||||
hdrlen = 2;
|
||||
} else if (len <= 0xffff) {
|
||||
hdr[0] = 0xda;
|
||||
hdr[1] = (len&0xff00)>>8;
|
||||
@@ -198,7 +202,7 @@ static void mp_encode_bytes(mp_buf *buf, const unsigned char *s, size_t len) {
|
||||
}
|
||||
|
||||
/* we assume IEEE 754 internal format for single and double precision floats. */
|
||||
static void mp_encode_double(mp_buf *buf, double d) {
|
||||
void mp_encode_double(mp_buf *buf, double d) {
|
||||
unsigned char b[9];
|
||||
float f = d;
|
||||
|
||||
@@ -216,7 +220,7 @@ static void mp_encode_double(mp_buf *buf, double d) {
|
||||
}
|
||||
}
|
||||
|
||||
static void mp_encode_int(mp_buf *buf, int64_t n) {
|
||||
void mp_encode_int(mp_buf *buf, int64_t n) {
|
||||
unsigned char b[9];
|
||||
int enclen;
|
||||
|
||||
@@ -288,7 +292,7 @@ static void mp_encode_int(mp_buf *buf, int64_t n) {
|
||||
mp_buf_append(buf,b,enclen);
|
||||
}
|
||||
|
||||
static void mp_encode_array(mp_buf *buf, int64_t n) {
|
||||
void mp_encode_array(mp_buf *buf, int64_t n) {
|
||||
unsigned char b[5];
|
||||
int enclen;
|
||||
|
||||
@@ -311,7 +315,7 @@ static void mp_encode_array(mp_buf *buf, int64_t n) {
|
||||
mp_buf_append(buf,b,enclen);
|
||||
}
|
||||
|
||||
static void mp_encode_map(mp_buf *buf, int64_t n) {
|
||||
void mp_encode_map(mp_buf *buf, int64_t n) {
|
||||
unsigned char b[5];
|
||||
int enclen;
|
||||
|
||||
@@ -336,7 +340,7 @@ static void mp_encode_map(mp_buf *buf, int64_t n) {
|
||||
|
||||
/* --------------------------- Lua types encoding --------------------------- */
|
||||
|
||||
static void mp_encode_lua_string(lua_State *L, mp_buf *buf) {
|
||||
void mp_encode_lua_string(lua_State *L, mp_buf *buf) {
|
||||
size_t len;
|
||||
const char *s;
|
||||
|
||||
@@ -344,13 +348,13 @@ static void mp_encode_lua_string(lua_State *L, mp_buf *buf) {
|
||||
mp_encode_bytes(buf,(const unsigned char*)s,len);
|
||||
}
|
||||
|
||||
static void mp_encode_lua_bool(lua_State *L, mp_buf *buf) {
|
||||
void mp_encode_lua_bool(lua_State *L, mp_buf *buf) {
|
||||
unsigned char b = lua_toboolean(L,-1) ? 0xc3 : 0xc2;
|
||||
mp_buf_append(buf,&b,1);
|
||||
}
|
||||
|
||||
/* Lua 5.3 has a built in 64-bit integer type */
|
||||
static void mp_encode_lua_integer(lua_State *L, mp_buf *buf) {
|
||||
void mp_encode_lua_integer(lua_State *L, mp_buf *buf) {
|
||||
#if (LUA_VERSION_NUM < 503) && BITS_32
|
||||
lua_Number i = lua_tonumber(L,-1);
|
||||
#else
|
||||
@@ -362,7 +366,7 @@ static void mp_encode_lua_integer(lua_State *L, mp_buf *buf) {
|
||||
/* Lua 5.2 and lower only has 64-bit doubles, so we need to
|
||||
* detect if the double may be representable as an int
|
||||
* for Lua < 5.3 */
|
||||
static void mp_encode_lua_number(lua_State *L, mp_buf *buf) {
|
||||
void mp_encode_lua_number(lua_State *L, mp_buf *buf) {
|
||||
lua_Number n = lua_tonumber(L,-1);
|
||||
|
||||
if (IS_INT64_EQUIVALENT(n)) {
|
||||
@@ -372,10 +376,10 @@ static void mp_encode_lua_number(lua_State *L, mp_buf *buf) {
|
||||
}
|
||||
}
|
||||
|
||||
static void mp_encode_lua_type(lua_State *L, mp_buf *buf, int level);
|
||||
void mp_encode_lua_type(lua_State *L, mp_buf *buf, int level);
|
||||
|
||||
/* Convert a lua table into a message pack list. */
|
||||
static void mp_encode_lua_table_as_array(lua_State *L, mp_buf *buf, int level) {
|
||||
void mp_encode_lua_table_as_array(lua_State *L, mp_buf *buf, int level) {
|
||||
#if LUA_VERSION_NUM < 502
|
||||
size_t len = lua_objlen(L,-1), j;
|
||||
#else
|
||||
@@ -391,7 +395,7 @@ static void mp_encode_lua_table_as_array(lua_State *L, mp_buf *buf, int level) {
|
||||
}
|
||||
|
||||
/* Convert a lua table into a message pack key-value map. */
|
||||
static void mp_encode_lua_table_as_map(lua_State *L, mp_buf *buf, int level) {
|
||||
void mp_encode_lua_table_as_map(lua_State *L, mp_buf *buf, int level) {
|
||||
size_t len = 0;
|
||||
|
||||
/* First step: count keys into table. No other way to do it with the
|
||||
@@ -418,7 +422,7 @@ static void mp_encode_lua_table_as_map(lua_State *L, mp_buf *buf, int level) {
|
||||
/* Returns true if the Lua table on top of the stack is exclusively composed
|
||||
* of keys from numerical keys from 1 up to N, with N being the total number
|
||||
* of elements, without any hole in the middle. */
|
||||
static int table_is_an_array(lua_State *L) {
|
||||
int table_is_an_array(lua_State *L) {
|
||||
int count = 0, max = 0;
|
||||
#if LUA_VERSION_NUM < 503
|
||||
lua_Number n;
|
||||
@@ -461,14 +465,14 @@ static int table_is_an_array(lua_State *L) {
|
||||
/* If the length operator returns non-zero, that is, there is at least
|
||||
* an object at key '1', we serialize to message pack list. Otherwise
|
||||
* we use a map. */
|
||||
static void mp_encode_lua_table(lua_State *L, mp_buf *buf, int level) {
|
||||
void mp_encode_lua_table(lua_State *L, mp_buf *buf, int level) {
|
||||
if (table_is_an_array(L))
|
||||
mp_encode_lua_table_as_array(L,buf,level);
|
||||
else
|
||||
mp_encode_lua_table_as_map(L,buf,level);
|
||||
}
|
||||
|
||||
static void mp_encode_lua_null(lua_State *L, mp_buf *buf) {
|
||||
void mp_encode_lua_null(lua_State *L, mp_buf *buf) {
|
||||
unsigned char b[1];
|
||||
(void)L;
|
||||
|
||||
@@ -476,7 +480,7 @@ static void mp_encode_lua_null(lua_State *L, mp_buf *buf) {
|
||||
mp_buf_append(buf,b,1);
|
||||
}
|
||||
|
||||
static void mp_encode_lua_type(lua_State *L, mp_buf *buf, int level) {
|
||||
void mp_encode_lua_type(lua_State *L, mp_buf *buf, int level) {
|
||||
int t = lua_type(L,-1);
|
||||
|
||||
/* Limit the encoding of nested tables to a specified maximum depth, so that
|
||||
@@ -506,7 +510,7 @@ static void mp_encode_lua_type(lua_State *L, mp_buf *buf, int level) {
|
||||
* Packs all arguments as a stream for multiple upacking later.
|
||||
* Returns error if no arguments provided.
|
||||
*/
|
||||
static int mp_pack(lua_State *L) {
|
||||
int mp_pack(lua_State *L) {
|
||||
int nargs = lua_gettop(L);
|
||||
int i;
|
||||
mp_buf *buf;
|
||||
@@ -687,6 +691,15 @@ void mp_decode_to_lua_type(lua_State *L, mp_cur *c) {
|
||||
mp_cur_consume(c,9);
|
||||
}
|
||||
break;
|
||||
case 0xd9: /* raw 8 */
|
||||
mp_cur_need(c,2);
|
||||
{
|
||||
size_t l = c->p[1];
|
||||
mp_cur_need(c,2+l);
|
||||
lua_pushlstring(L,(char*)c->p+2,l);
|
||||
mp_cur_consume(c,2+l);
|
||||
}
|
||||
break;
|
||||
case 0xda: /* raw 16 */
|
||||
mp_cur_need(c,3);
|
||||
{
|
||||
@@ -773,7 +786,7 @@ void mp_decode_to_lua_type(lua_State *L, mp_cur *c) {
|
||||
}
|
||||
}
|
||||
|
||||
static int mp_unpack_full(lua_State *L, int limit, int offset) {
|
||||
int mp_unpack_full(lua_State *L, int limit, int offset) {
|
||||
size_t len;
|
||||
const char *s;
|
||||
mp_cur c;
|
||||
@@ -826,18 +839,18 @@ static int mp_unpack_full(lua_State *L, int limit, int offset) {
|
||||
return cnt;
|
||||
}
|
||||
|
||||
static int mp_unpack(lua_State *L) {
|
||||
int mp_unpack(lua_State *L) {
|
||||
return mp_unpack_full(L, 0, 0);
|
||||
}
|
||||
|
||||
static int mp_unpack_one(lua_State *L) {
|
||||
int mp_unpack_one(lua_State *L) {
|
||||
int offset = luaL_optinteger(L, 2, 0);
|
||||
/* Variable pop because offset may not exist */
|
||||
lua_pop(L, lua_gettop(L)-1);
|
||||
return mp_unpack_full(L, 1, offset);
|
||||
}
|
||||
|
||||
static int mp_unpack_limit(lua_State *L) {
|
||||
int mp_unpack_limit(lua_State *L) {
|
||||
int limit = luaL_checkinteger(L, 2);
|
||||
int offset = luaL_optinteger(L, 3, 0);
|
||||
/* Variable pop because offset may not exist */
|
||||
@@ -846,7 +859,7 @@ static int mp_unpack_limit(lua_State *L) {
|
||||
return mp_unpack_full(L, limit, offset);
|
||||
}
|
||||
|
||||
static int mp_safe(lua_State *L) {
|
||||
int mp_safe(lua_State *L) {
|
||||
int argc, err, total_results;
|
||||
|
||||
argc = lua_gettop(L);
|
||||
@@ -869,7 +882,7 @@ static int mp_safe(lua_State *L) {
|
||||
}
|
||||
|
||||
/* -------------------------------------------------------------------------- */
|
||||
static const struct luaL_Reg cmds[] = {
|
||||
const struct luaL_Reg cmds[] = {
|
||||
{"pack", mp_pack},
|
||||
{"unpack", mp_unpack},
|
||||
{"unpack_one", mp_unpack_one},
|
||||
@@ -877,7 +890,7 @@ static const struct luaL_Reg cmds[] = {
|
||||
{0}
|
||||
};
|
||||
|
||||
static int luaopen_create(lua_State *L) {
|
||||
int luaopen_create(lua_State *L) {
|
||||
int i;
|
||||
/* Manually construct our module table instead of
|
||||
* relying on _register or _newlib */
|
||||
|
||||
+26
-1
@@ -59,6 +59,8 @@
|
||||
* When implementing a new type of blocking opeation, the implementation
|
||||
* should modify unblockClient() and replyToBlockedClientTimedOut() in order
|
||||
* to handle the btype-specific behavior of this two functions.
|
||||
* If the blocking operation waits for certain keys to change state, the
|
||||
* clusterRedirectBlockedClientIfNeeded() function should also be updated.
|
||||
*/
|
||||
|
||||
#include "redis.h"
|
||||
@@ -114,7 +116,6 @@ void processUnblockedClients(void) {
|
||||
c = ln->value;
|
||||
listDelNode(server.unblocked_clients,ln);
|
||||
c->flags &= ~REDIS_UNBLOCKED;
|
||||
c->btype = REDIS_BLOCKED_NONE;
|
||||
|
||||
/* Process remaining data in the input buffer. */
|
||||
if (c->querybuf && sdslen(c->querybuf) > 0) {
|
||||
@@ -156,3 +157,27 @@ void replyToBlockedClientTimedOut(redisClient *c) {
|
||||
}
|
||||
}
|
||||
|
||||
/* Mass-unblock clients because something changed in the instance that makes
|
||||
* blocking no longer safe. For example clients blocked in list operations
|
||||
* in an instance which turns from master to slave is unsafe, so this function
|
||||
* is called when a master turns into a slave.
|
||||
*
|
||||
* The semantics is to send an -UNBLOCKED error to the client, disconnecting
|
||||
* it at the same time. */
|
||||
void disconnectAllBlockedClients(void) {
|
||||
listNode *ln;
|
||||
listIter li;
|
||||
|
||||
listRewind(server.clients,&li);
|
||||
while((ln = listNext(&li))) {
|
||||
redisClient *c = listNodeValue(ln);
|
||||
|
||||
if (c->flags & REDIS_BLOCKED) {
|
||||
addReplySds(c,sdsnew(
|
||||
"-UNBLOCKED force unblock from blocking operation, "
|
||||
"instance state changed (master -> slave?)\r\n"));
|
||||
unblockClient(c);
|
||||
c->flags |= REDIS_CLOSE_AFTER_REPLY;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+287
-131
@@ -74,27 +74,13 @@ void clusterCloseAllSlots(void);
|
||||
void clusterSetNodeAsMaster(clusterNode *n);
|
||||
void clusterDelNode(clusterNode *delnode);
|
||||
sds representRedisNodeFlags(sds ci, uint16_t flags);
|
||||
uint64_t clusterGetMaxEpoch(void);
|
||||
int clusterBumpConfigEpochWithoutConsensus(void);
|
||||
|
||||
/* -----------------------------------------------------------------------------
|
||||
* Initialization
|
||||
* -------------------------------------------------------------------------- */
|
||||
|
||||
/* Return the greatest configEpoch found in the cluster. */
|
||||
uint64_t clusterGetMaxEpoch(void) {
|
||||
uint64_t max = 0;
|
||||
dictIterator *di;
|
||||
dictEntry *de;
|
||||
|
||||
di = dictGetSafeIterator(server.cluster->nodes);
|
||||
while((de = dictNext(di)) != NULL) {
|
||||
clusterNode *node = dictGetVal(de);
|
||||
if (node->configEpoch > max) max = node->configEpoch;
|
||||
}
|
||||
dictReleaseIterator(di);
|
||||
if (max < server.cluster->currentEpoch) max = server.cluster->currentEpoch;
|
||||
return max;
|
||||
}
|
||||
|
||||
/* Load the cluster config from 'filename'.
|
||||
*
|
||||
* If the file does not exist or is zero-length (this may happen because
|
||||
@@ -927,6 +913,138 @@ void clusterRenameNode(clusterNode *node, char *newname) {
|
||||
clusterAddNode(node);
|
||||
}
|
||||
|
||||
/* -----------------------------------------------------------------------------
|
||||
* CLUSTER config epoch handling
|
||||
* -------------------------------------------------------------------------- */
|
||||
|
||||
/* Return the greatest configEpoch found in the cluster, or the current
|
||||
* epoch if greater than any node configEpoch. */
|
||||
uint64_t clusterGetMaxEpoch(void) {
|
||||
uint64_t max = 0;
|
||||
dictIterator *di;
|
||||
dictEntry *de;
|
||||
|
||||
di = dictGetSafeIterator(server.cluster->nodes);
|
||||
while((de = dictNext(di)) != NULL) {
|
||||
clusterNode *node = dictGetVal(de);
|
||||
if (node->configEpoch > max) max = node->configEpoch;
|
||||
}
|
||||
dictReleaseIterator(di);
|
||||
if (max < server.cluster->currentEpoch) max = server.cluster->currentEpoch;
|
||||
return max;
|
||||
}
|
||||
|
||||
/* If this node epoch is zero or is not already the greatest across the
|
||||
* cluster (from the POV of the local configuration), this function will:
|
||||
*
|
||||
* 1) Generate a new config epoch increment the current epoch.
|
||||
* 2) Assign the new epoch to this node, WITHOUT any consensus.
|
||||
* 3) Persist the configuration on disk before sending packets with the
|
||||
* new configuration.
|
||||
*
|
||||
* If the new config epoch is generated and assigend, REDIS_OK is returned,
|
||||
* otherwise REDIS_ERR is returned (since the node has already the greatest
|
||||
* configuration around) and no operation is performed.
|
||||
*
|
||||
* Important note: this function violates the principle that config epochs
|
||||
* should be generated with consensus and should be unique across the cluster.
|
||||
* However Redis Cluster uses this auto-generated new config epochs in two
|
||||
* cases:
|
||||
*
|
||||
* 1) When slots are closed after importing. Otherwise resharding would be
|
||||
* too exansive.
|
||||
* 2) When CLUSTER FAILOVER is called with options that force a slave to
|
||||
* failover its master even if there is not master majority able to
|
||||
* create a new configuration epoch.
|
||||
*
|
||||
* Redis Cluster does not explode using this function, even in the case of
|
||||
* a collision between this node and another node, generating the same
|
||||
* configuration epoch unilaterally, because the config epoch conflict
|
||||
* resolution algorithm will eventually move colliding nodes to different
|
||||
* config epochs. However usign this function may violate the "last failover
|
||||
* wins" rule, so should only be used with care. */
|
||||
int clusterBumpConfigEpochWithoutConsensus(void) {
|
||||
uint64_t maxEpoch = clusterGetMaxEpoch();
|
||||
|
||||
if (myself->configEpoch == 0 ||
|
||||
myself->configEpoch != maxEpoch)
|
||||
{
|
||||
server.cluster->currentEpoch++;
|
||||
myself->configEpoch = server.cluster->currentEpoch;
|
||||
clusterDoBeforeSleep(CLUSTER_TODO_SAVE_CONFIG|
|
||||
CLUSTER_TODO_FSYNC_CONFIG);
|
||||
redisLog(REDIS_WARNING,
|
||||
"New configEpoch set to %llu",
|
||||
(unsigned long long) myself->configEpoch);
|
||||
return REDIS_OK;
|
||||
} else {
|
||||
return REDIS_ERR;
|
||||
}
|
||||
}
|
||||
|
||||
/* This function is called when this node is a master, and we receive from
|
||||
* another master a configuration epoch that is equal to our configuration
|
||||
* epoch.
|
||||
*
|
||||
* BACKGROUND
|
||||
*
|
||||
* It is not possible that different slaves get the same config
|
||||
* epoch during a failover election, because the slaves need to get voted
|
||||
* by a majority. However when we perform a manual resharding of the cluster
|
||||
* the node will assign a configuration epoch to itself without to ask
|
||||
* for agreement. Usually resharding happens when the cluster is working well
|
||||
* and is supervised by the sysadmin, however it is possible for a failover
|
||||
* to happen exactly while the node we are resharding a slot to assigns itself
|
||||
* a new configuration epoch, but before it is able to propagate it.
|
||||
*
|
||||
* So technically it is possible in this condition that two nodes end with
|
||||
* the same configuration epoch.
|
||||
*
|
||||
* Another possibility is that there are bugs in the implementation causing
|
||||
* this to happen.
|
||||
*
|
||||
* Moreover when a new cluster is created, all the nodes start with the same
|
||||
* configEpoch. This collision resolution code allows nodes to automatically
|
||||
* end with a different configEpoch at startup automatically.
|
||||
*
|
||||
* In all the cases, we want a mechanism that resolves this issue automatically
|
||||
* as a safeguard. The same configuration epoch for masters serving different
|
||||
* set of slots is not harmful, but it is if the nodes end serving the same
|
||||
* slots for some reason (manual errors or software bugs) without a proper
|
||||
* failover procedure.
|
||||
*
|
||||
* In general we want a system that eventually always ends with different
|
||||
* masters having different configuration epochs whatever happened, since
|
||||
* nothign is worse than a split-brain condition in a distributed system.
|
||||
*
|
||||
* BEHAVIOR
|
||||
*
|
||||
* When this function gets called, what happens is that if this node
|
||||
* has the lexicographically smaller Node ID compared to the other node
|
||||
* with the conflicting epoch (the 'sender' node), it will assign itself
|
||||
* the greatest configuration epoch currently detected among nodes plus 1.
|
||||
*
|
||||
* This means that even if there are multiple nodes colliding, the node
|
||||
* with the greatest Node ID never moves forward, so eventually all the nodes
|
||||
* end with a different configuration epoch.
|
||||
*/
|
||||
void clusterHandleConfigEpochCollision(clusterNode *sender) {
|
||||
/* Prerequisites: nodes have the same configEpoch and are both masters. */
|
||||
if (sender->configEpoch != myself->configEpoch ||
|
||||
!nodeIsMaster(sender) || !nodeIsMaster(myself)) return;
|
||||
/* Don't act if the colliding node has a smaller Node ID. */
|
||||
if (memcmp(sender->name,myself->name,REDIS_CLUSTER_NAMELEN) <= 0) return;
|
||||
/* Get the next ID available at the best of this node knowledge. */
|
||||
server.cluster->currentEpoch++;
|
||||
myself->configEpoch = server.cluster->currentEpoch;
|
||||
clusterSaveConfigOrDie(1);
|
||||
redisLog(REDIS_VERBOSE,
|
||||
"WARNING: configEpoch collision with node %.40s."
|
||||
" configEpoch set to %llu",
|
||||
sender->name,
|
||||
(unsigned long long) myself->configEpoch);
|
||||
}
|
||||
|
||||
/* -----------------------------------------------------------------------------
|
||||
* CLUSTER nodes blacklist
|
||||
*
|
||||
@@ -1399,69 +1517,6 @@ void clusterUpdateSlotsConfigWith(clusterNode *sender, uint64_t senderConfigEpoc
|
||||
}
|
||||
}
|
||||
|
||||
/* This function is called when this node is a master, and we receive from
|
||||
* another master a configuration epoch that is equal to our configuration
|
||||
* epoch.
|
||||
*
|
||||
* BACKGROUND
|
||||
*
|
||||
* It is not possible that different slaves get the same config
|
||||
* epoch during a failover election, because the slaves need to get voted
|
||||
* by a majority. However when we perform a manual resharding of the cluster
|
||||
* the node will assign a configuration epoch to itself without to ask
|
||||
* for agreement. Usually resharding happens when the cluster is working well
|
||||
* and is supervised by the sysadmin, however it is possible for a failover
|
||||
* to happen exactly while the node we are resharding a slot to assigns itself
|
||||
* a new configuration epoch, but before it is able to propagate it.
|
||||
*
|
||||
* So technically it is possible in this condition that two nodes end with
|
||||
* the same configuration epoch.
|
||||
*
|
||||
* Another possibility is that there are bugs in the implementation causing
|
||||
* this to happen.
|
||||
*
|
||||
* Moreover when a new cluster is created, all the nodes start with the same
|
||||
* configEpoch. This collision resolution code allows nodes to automatically
|
||||
* end with a different configEpoch at startup automatically.
|
||||
*
|
||||
* In all the cases, we want a mechanism that resolves this issue automatically
|
||||
* as a safeguard. The same configuration epoch for masters serving different
|
||||
* set of slots is not harmful, but it is if the nodes end serving the same
|
||||
* slots for some reason (manual errors or software bugs) without a proper
|
||||
* failover procedure.
|
||||
*
|
||||
* In general we want a system that eventually always ends with different
|
||||
* masters having different configuration epochs whatever happened, since
|
||||
* nothign is worse than a split-brain condition in a distributed system.
|
||||
*
|
||||
* BEHAVIOR
|
||||
*
|
||||
* When this function gets called, what happens is that if this node
|
||||
* has the lexicographically smaller Node ID compared to the other node
|
||||
* with the conflicting epoch (the 'sender' node), it will assign itself
|
||||
* the greatest configuration epoch currently detected among nodes plus 1.
|
||||
*
|
||||
* This means that even if there are multiple nodes colliding, the node
|
||||
* with the greatest Node ID never moves forward, so eventually all the nodes
|
||||
* end with a different configuration epoch.
|
||||
*/
|
||||
void clusterHandleConfigEpochCollision(clusterNode *sender) {
|
||||
/* Prerequisites: nodes have the same configEpoch and are both masters. */
|
||||
if (sender->configEpoch != myself->configEpoch ||
|
||||
!nodeIsMaster(sender) || !nodeIsMaster(myself)) return;
|
||||
/* Don't act if the colliding node has a smaller Node ID. */
|
||||
if (memcmp(sender->name,myself->name,REDIS_CLUSTER_NAMELEN) <= 0) return;
|
||||
/* Get the next ID available at the best of this node knowledge. */
|
||||
server.cluster->currentEpoch++;
|
||||
myself->configEpoch = server.cluster->currentEpoch;
|
||||
clusterSaveConfigOrDie(1);
|
||||
redisLog(REDIS_VERBOSE,
|
||||
"WARNING: configEpoch collision with node %.40s."
|
||||
" configEpoch set to %llu",
|
||||
sender->name,
|
||||
(unsigned long long) myself->configEpoch);
|
||||
}
|
||||
|
||||
/* When this function is called, there is a packet to process starting
|
||||
* at node->rcvbuf. Releasing the buffer is up to the caller, so this
|
||||
* function should just handle the higher level stuff of processing the
|
||||
@@ -2582,6 +2637,42 @@ void clusterLogCantFailover(int reason) {
|
||||
redisLog(REDIS_WARNING,"Currently unable to failover: %s", msg);
|
||||
}
|
||||
|
||||
/* This function implements the final part of automatic and manual failovers,
|
||||
* where the slave grabs its master's hash slots, and propagates the new
|
||||
* configuration.
|
||||
*
|
||||
* Note that it's up to the caller to be sure that the node got a new
|
||||
* configuration epoch already. */
|
||||
void clusterFailoverReplaceYourMaster(void) {
|
||||
int j;
|
||||
clusterNode *oldmaster = myself->slaveof;
|
||||
|
||||
if (nodeIsMaster(myself) || oldmaster == NULL) return;
|
||||
|
||||
/* 1) Turn this node into a master. */
|
||||
clusterSetNodeAsMaster(myself);
|
||||
replicationUnsetMaster();
|
||||
|
||||
/* 2) Claim all the slots assigned to our master. */
|
||||
for (j = 0; j < REDIS_CLUSTER_SLOTS; j++) {
|
||||
if (clusterNodeGetSlotBit(oldmaster,j)) {
|
||||
clusterDelSlot(j);
|
||||
clusterAddSlot(myself,j);
|
||||
}
|
||||
}
|
||||
|
||||
/* 3) Update state and save config. */
|
||||
clusterUpdateState();
|
||||
clusterSaveConfigOrDie(1);
|
||||
|
||||
/* 4) Pong all the other nodes so that they can update the state
|
||||
* accordingly and detect that we switched to master role. */
|
||||
clusterBroadcastPong(CLUSTER_BROADCAST_ALL);
|
||||
|
||||
/* 5) If there was a manual failover in progress, clear the state. */
|
||||
resetManualFailover();
|
||||
}
|
||||
|
||||
/* This function is called if we are a slave node and our master serving
|
||||
* a non-zero amount of hash slots is in FAIL state.
|
||||
*
|
||||
@@ -2596,7 +2687,6 @@ void clusterHandleSlaveFailover(void) {
|
||||
int needed_quorum = (server.cluster->size / 2) + 1;
|
||||
int manual_failover = server.cluster->mf_end != 0 &&
|
||||
server.cluster->mf_can_start;
|
||||
int j;
|
||||
mstime_t auth_timeout, auth_retry_time;
|
||||
|
||||
server.cluster->todo_before_sleep &= ~CLUSTER_TODO_HANDLE_FAILOVER;
|
||||
@@ -2738,26 +2828,12 @@ void clusterHandleSlaveFailover(void) {
|
||||
|
||||
/* Check if we reached the quorum. */
|
||||
if (server.cluster->failover_auth_count >= needed_quorum) {
|
||||
clusterNode *oldmaster = myself->slaveof;
|
||||
/* We have the quorum, we can finally failover the master. */
|
||||
|
||||
redisLog(REDIS_WARNING,
|
||||
"Failover election won: I'm the new master.");
|
||||
/* We have the quorum, perform all the steps to correctly promote
|
||||
* this slave to a master.
|
||||
*
|
||||
* 1) Turn this node into a master. */
|
||||
clusterSetNodeAsMaster(myself);
|
||||
replicationUnsetMaster();
|
||||
|
||||
/* 2) Claim all the slots assigned to our master. */
|
||||
for (j = 0; j < REDIS_CLUSTER_SLOTS; j++) {
|
||||
if (clusterNodeGetSlotBit(oldmaster,j)) {
|
||||
clusterDelSlot(j);
|
||||
clusterAddSlot(myself,j);
|
||||
}
|
||||
}
|
||||
|
||||
/* 3) Update my configEpoch to the epoch of the election. */
|
||||
/* Update my configEpoch to the epoch of the election. */
|
||||
if (myself->configEpoch < server.cluster->failover_auth_epoch) {
|
||||
myself->configEpoch = server.cluster->failover_auth_epoch;
|
||||
redisLog(REDIS_WARNING,
|
||||
@@ -2765,16 +2841,8 @@ void clusterHandleSlaveFailover(void) {
|
||||
(unsigned long long) myself->configEpoch);
|
||||
}
|
||||
|
||||
/* 4) Update state and save config. */
|
||||
clusterUpdateState();
|
||||
clusterSaveConfigOrDie(1);
|
||||
|
||||
/* 5) Pong all the other nodes so that they can update the state
|
||||
* accordingly and detect that we switched to master role. */
|
||||
clusterBroadcastPong(CLUSTER_BROADCAST_ALL);
|
||||
|
||||
/* 6) If there was a manual failover in progress, clear the state. */
|
||||
resetManualFailover();
|
||||
/* Take responsability for the cluster slots. */
|
||||
clusterFailoverReplaceYourMaster();
|
||||
} else {
|
||||
clusterLogCantFailover(REDIS_CLUSTER_CANT_FAILOVER_WAITING_VOTES);
|
||||
}
|
||||
@@ -3567,7 +3635,7 @@ sds clusterGenNodeDescription(clusterNode *node) {
|
||||
else
|
||||
ci = sdscatlen(ci," - ",3);
|
||||
|
||||
/* Latency from the POV of this node, link status */
|
||||
/* Latency from the POV of this node, config epoch, link status */
|
||||
ci = sdscatprintf(ci,"%lld %lld %llu %s",
|
||||
(long long) node->ping_sent,
|
||||
(long long) node->pong_received,
|
||||
@@ -3902,17 +3970,9 @@ void clusterCommand(redisClient *c) {
|
||||
* failover happens at the same time we close the slot, the
|
||||
* configEpoch collision resolution will fix it assigning
|
||||
* a different epoch to each node. */
|
||||
uint64_t maxEpoch = clusterGetMaxEpoch();
|
||||
|
||||
if (myself->configEpoch == 0 ||
|
||||
myself->configEpoch != maxEpoch)
|
||||
{
|
||||
server.cluster->currentEpoch++;
|
||||
myself->configEpoch = server.cluster->currentEpoch;
|
||||
clusterDoBeforeSleep(CLUSTER_TODO_FSYNC_CONFIG);
|
||||
if (clusterBumpConfigEpochWithoutConsensus() == REDIS_OK) {
|
||||
redisLog(REDIS_WARNING,
|
||||
"configEpoch set to %llu after importing slot %d",
|
||||
(unsigned long long) myself->configEpoch, slot);
|
||||
"configEpoch updated after importing slot %d", slot);
|
||||
}
|
||||
server.cluster->importing_slots_from[slot] = NULL;
|
||||
}
|
||||
@@ -4115,24 +4175,31 @@ void clusterCommand(redisClient *c) {
|
||||
} else if (!strcasecmp(c->argv[1]->ptr,"failover") &&
|
||||
(c->argc == 2 || c->argc == 3))
|
||||
{
|
||||
/* CLUSTER FAILOVER [FORCE] */
|
||||
int force = 0;
|
||||
/* CLUSTER FAILOVER [FORCE|TAKEOVER] */
|
||||
int force = 0, takeover = 0;
|
||||
|
||||
if (c->argc == 3) {
|
||||
if (!strcasecmp(c->argv[2]->ptr,"force")) {
|
||||
force = 1;
|
||||
} else if (!strcasecmp(c->argv[2]->ptr,"takeover")) {
|
||||
takeover = 1;
|
||||
force = 1; /* Takeover also implies force. */
|
||||
} else {
|
||||
addReply(c,shared.syntaxerr);
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
||||
/* Check preconditions. */
|
||||
if (nodeIsMaster(myself)) {
|
||||
addReplyError(c,"You should send CLUSTER FAILOVER to a slave");
|
||||
return;
|
||||
} else if (myself->slaveof == NULL) {
|
||||
addReplyError(c,"I'm a slave but my master is unknown to me");
|
||||
return;
|
||||
} else if (!force &&
|
||||
(myself->slaveof == NULL || nodeFailed(myself->slaveof) ||
|
||||
myself->slaveof->link == NULL))
|
||||
(nodeFailed(myself->slaveof) ||
|
||||
myself->slaveof->link == NULL))
|
||||
{
|
||||
addReplyError(c,"Master is down or failed, "
|
||||
"please use CLUSTER FAILOVER FORCE");
|
||||
@@ -4141,15 +4208,24 @@ void clusterCommand(redisClient *c) {
|
||||
resetManualFailover();
|
||||
server.cluster->mf_end = mstime() + REDIS_CLUSTER_MF_TIMEOUT;
|
||||
|
||||
/* If this is a forced failover, we don't need to talk with our master
|
||||
* to agree about the offset. We just failover taking over it without
|
||||
* coordination. */
|
||||
if (force) {
|
||||
if (takeover) {
|
||||
/* A takeover does not perform any initial check. It just
|
||||
* generates a new configuration epoch for this node without
|
||||
* consensus, claims the master's slots, and broadcast the new
|
||||
* configuration. */
|
||||
redisLog(REDIS_WARNING,"Taking over the master (user request).");
|
||||
clusterBumpConfigEpochWithoutConsensus();
|
||||
clusterFailoverReplaceYourMaster();
|
||||
} else if (force) {
|
||||
/* If this is a forced failover, we don't need to talk with our
|
||||
* master to agree about the offset. We just failover taking over
|
||||
* it without coordination. */
|
||||
redisLog(REDIS_WARNING,"Forced failover user request accepted.");
|
||||
server.cluster->mf_can_start = 1;
|
||||
} else {
|
||||
redisLog(REDIS_WARNING,"Manual failover user request accepted.");
|
||||
clusterSendMFStart(myself->slaveof);
|
||||
}
|
||||
redisLog(REDIS_WARNING,"Manual failover user request accepted.");
|
||||
addReply(c,shared.ok);
|
||||
} else if (!strcasecmp(c->argv[1]->ptr,"set-config-epoch") && c->argc == 3)
|
||||
{
|
||||
@@ -4692,10 +4768,10 @@ void readwriteCommand(redisClient *c) {
|
||||
* belonging to the same slot, but the slot is not stable (in migration or
|
||||
* importing state, likely because a resharding is in progress).
|
||||
*
|
||||
* REDIS_CLUSTER_REDIR_DOWN if the request addresses a slot which is not
|
||||
* bound to any node. In this case the cluster global state should be already
|
||||
* "down" but it is fragile to rely on the update of the global state, so
|
||||
* we also handle it here. */
|
||||
* REDIS_CLUSTER_REDIR_DOWN_UNBOUND if the request addresses a slot which is
|
||||
* not bound to any node. In this case the cluster global state should be
|
||||
* already "down" but it is fragile to rely on the update of the global state,
|
||||
* so we also handle it here. */
|
||||
clusterNode *getNodeByQuery(redisClient *c, struct redisCommand *cmd, robj **argv, int argc, int *hashslot, int *error_code) {
|
||||
clusterNode *n = NULL;
|
||||
robj *firstkey = NULL;
|
||||
@@ -4757,7 +4833,7 @@ clusterNode *getNodeByQuery(redisClient *c, struct redisCommand *cmd, robj **arg
|
||||
if (n == NULL) {
|
||||
getKeysFreeResult(keyindex);
|
||||
if (error_code)
|
||||
*error_code = REDIS_CLUSTER_REDIR_DOWN;
|
||||
*error_code = REDIS_CLUSTER_REDIR_DOWN_UNBOUND;
|
||||
return NULL;
|
||||
}
|
||||
|
||||
@@ -4849,3 +4925,83 @@ clusterNode *getNodeByQuery(redisClient *c, struct redisCommand *cmd, robj **arg
|
||||
if (n != myself && error_code) *error_code = REDIS_CLUSTER_REDIR_MOVED;
|
||||
return n;
|
||||
}
|
||||
|
||||
/* Send the client the right redirection code, according to error_code
|
||||
* that should be set to one of REDIS_CLUSTER_REDIR_* macros.
|
||||
*
|
||||
* If REDIS_CLUSTER_REDIR_ASK or REDIS_CLUSTER_REDIR_MOVED error codes
|
||||
* are used, then the node 'n' should not be NULL, but should be the
|
||||
* node we want to mention in the redirection. Moreover hashslot should
|
||||
* be set to the hash slot that caused the redirection. */
|
||||
void clusterRedirectClient(redisClient *c, clusterNode *n, int hashslot, int error_code) {
|
||||
if (error_code == REDIS_CLUSTER_REDIR_CROSS_SLOT) {
|
||||
addReplySds(c,sdsnew("-CROSSSLOT Keys in request don't hash to the same slot\r\n"));
|
||||
} else if (error_code == REDIS_CLUSTER_REDIR_UNSTABLE) {
|
||||
/* The request spawns mutliple keys in the same slot,
|
||||
* but the slot is not "stable" currently as there is
|
||||
* a migration or import in progress. */
|
||||
addReplySds(c,sdsnew("-TRYAGAIN Multiple keys request during rehashing of slot\r\n"));
|
||||
} else if (error_code == REDIS_CLUSTER_REDIR_DOWN_STATE) {
|
||||
addReplySds(c,sdsnew("-CLUSTERDOWN The cluster is down\r\n"));
|
||||
} else if (error_code == REDIS_CLUSTER_REDIR_DOWN_UNBOUND) {
|
||||
addReplySds(c,sdsnew("-CLUSTERDOWN Hash slot not served\r\n"));
|
||||
} else if (error_code == REDIS_CLUSTER_REDIR_MOVED ||
|
||||
error_code == REDIS_CLUSTER_REDIR_ASK)
|
||||
{
|
||||
addReplySds(c,sdscatprintf(sdsempty(),
|
||||
"-%s %d %s:%d\r\n",
|
||||
(error_code == REDIS_CLUSTER_REDIR_ASK) ? "ASK" : "MOVED",
|
||||
hashslot,n->ip,n->port));
|
||||
} else {
|
||||
redisPanic("getNodeByQuery() unknown error.");
|
||||
}
|
||||
}
|
||||
|
||||
/* This function is called by the function processing clients incrementally
|
||||
* to detect timeouts, in order to handle the following case:
|
||||
*
|
||||
* 1) A client blocks with BLPOP or similar blocking operation.
|
||||
* 2) The master migrates the hash slot elsewhere or turns into a slave.
|
||||
* 3) The client may remain blocked forever (or up to the max timeout time)
|
||||
* waiting for a key change that will never happen.
|
||||
*
|
||||
* If the client is found to be blocked into an hash slot this node no
|
||||
* longer handles, the client is sent a redirection error, and the function
|
||||
* returns 1. Otherwise 0 is returned and no operation is performed. */
|
||||
int clusterRedirectBlockedClientIfNeeded(redisClient *c) {
|
||||
if (c->flags & REDIS_BLOCKED && c->btype == REDIS_BLOCKED_LIST) {
|
||||
dictEntry *de;
|
||||
dictIterator *di;
|
||||
|
||||
/* If the cluster is down, unblock the client with the right error. */
|
||||
if (server.cluster->state == REDIS_CLUSTER_FAIL) {
|
||||
clusterRedirectClient(c,NULL,0,REDIS_CLUSTER_REDIR_DOWN_STATE);
|
||||
return 1;
|
||||
}
|
||||
|
||||
di = dictGetIterator(c->bpop.keys);
|
||||
while((de = dictNext(di)) != NULL) {
|
||||
robj *key = dictGetKey(de);
|
||||
int slot = keyHashSlot((char*)key->ptr, sdslen(key->ptr));
|
||||
clusterNode *node = server.cluster->slots[slot];
|
||||
|
||||
/* We send an error and unblock the client if:
|
||||
* 1) The slot is unassigned, emitting a cluster down error.
|
||||
* 2) The slot is not handled by this node, nor being imported. */
|
||||
if (node != myself &&
|
||||
server.cluster->importing_slots_from[slot] == NULL)
|
||||
{
|
||||
if (node == NULL) {
|
||||
clusterRedirectClient(c,NULL,0,
|
||||
REDIS_CLUSTER_REDIR_DOWN_UNBOUND);
|
||||
} else {
|
||||
clusterRedirectClient(c,node,slot,
|
||||
REDIS_CLUSTER_REDIR_MOVED);
|
||||
}
|
||||
return 1;
|
||||
}
|
||||
}
|
||||
dictReleaseIterator(di);
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
|
||||
+6
-3
@@ -26,11 +26,12 @@
|
||||
|
||||
/* Redirection errors returned by getNodeByQuery(). */
|
||||
#define REDIS_CLUSTER_REDIR_NONE 0 /* Node can serve the request. */
|
||||
#define REDIS_CLUSTER_REDIR_CROSS_SLOT 1 /* Keys in different slots. */
|
||||
#define REDIS_CLUSTER_REDIR_UNSTABLE 2 /* Keys in slot resharding. */
|
||||
#define REDIS_CLUSTER_REDIR_CROSS_SLOT 1 /* -CROSSSLOT request. */
|
||||
#define REDIS_CLUSTER_REDIR_UNSTABLE 2 /* -TRYAGAIN redirection required */
|
||||
#define REDIS_CLUSTER_REDIR_ASK 3 /* -ASK redirection required. */
|
||||
#define REDIS_CLUSTER_REDIR_MOVED 4 /* -MOVED redirection required. */
|
||||
#define REDIS_CLUSTER_REDIR_DOWN 5 /* -CLUSTERDOWN error. */
|
||||
#define REDIS_CLUSTER_REDIR_DOWN_STATE 5 /* -CLUSTERDOWN, global state. */
|
||||
#define REDIS_CLUSTER_REDIR_DOWN_UNBOUND 6 /* -CLUSTERDOWN, unbound slot. */
|
||||
|
||||
struct clusterNode;
|
||||
|
||||
@@ -249,5 +250,7 @@ typedef struct {
|
||||
|
||||
/* ---------------------- API exported outside cluster.c -------------------- */
|
||||
clusterNode *getNodeByQuery(redisClient *c, struct redisCommand *cmd, robj **argv, int argc, int *hashslot, int *ask);
|
||||
int clusterRedirectBlockedClientIfNeeded(redisClient *c);
|
||||
void clusterRedirectClient(redisClient *c, clusterNode *n, int hashslot, int error_code);
|
||||
|
||||
#endif /* __REDIS_CLUSTER_H */
|
||||
|
||||
+1
-76
@@ -661,81 +661,6 @@ dictEntry *dictGetRandomKey(dict *d)
|
||||
return he;
|
||||
}
|
||||
|
||||
/* XXX: This is going to be removed soon and SPOP internals
|
||||
* reimplemented.
|
||||
*
|
||||
* This is a version of dictGetRandomKey() that is modified in order to
|
||||
* return multiple entries by jumping at a random place of the hash table
|
||||
* and scanning linearly for entries.
|
||||
*
|
||||
* Returned pointers to hash table entries are stored into 'des' that
|
||||
* points to an array of dictEntry pointers. The array must have room for
|
||||
* at least 'count' elements, that is the argument we pass to the function
|
||||
* to tell how many random elements we need.
|
||||
*
|
||||
* The function returns the number of items stored into 'des', that may
|
||||
* be less than 'count' if the hash table has less than 'count' elements
|
||||
* inside.
|
||||
*
|
||||
* Note that this function is not suitable when you need a good distribution
|
||||
* of the returned items, but only when you need to "sample" a given number
|
||||
* of continuous elements to run some kind of algorithm or to produce
|
||||
* statistics. However the function is much faster than dictGetRandomKey()
|
||||
* at producing N elements, and the elements are guaranteed to be non
|
||||
* repeating. */
|
||||
unsigned int dictGetRandomKeys(dict *d, dictEntry **des, unsigned int count) {
|
||||
unsigned int j; /* internal hash table id, 0 or 1. */
|
||||
unsigned int tables; /* 1 or 2 tables? */
|
||||
unsigned int stored = 0, maxsizemask;
|
||||
|
||||
if (dictSize(d) < count) count = dictSize(d);
|
||||
|
||||
/* Try to do a rehashing work proportional to 'count'. */
|
||||
for (j = 0; j < count; j++) {
|
||||
if (dictIsRehashing(d))
|
||||
_dictRehashStep(d);
|
||||
else
|
||||
break;
|
||||
}
|
||||
|
||||
tables = dictIsRehashing(d) ? 2 : 1;
|
||||
maxsizemask = d->ht[0].sizemask;
|
||||
if (tables > 1 && maxsizemask < d->ht[1].sizemask)
|
||||
maxsizemask = d->ht[1].sizemask;
|
||||
|
||||
/* Pick a random point inside the larger table. */
|
||||
unsigned int i = random() & maxsizemask;
|
||||
while(stored < count) {
|
||||
for (j = 0; j < tables; j++) {
|
||||
/* Invariant of the dict.c rehashing: up to the indexes already
|
||||
* visited in ht[0] during the rehashing, there are no populated
|
||||
* buckets, so we can skip ht[0] for indexes between 0 and idx-1. */
|
||||
if (tables == 2 && j == 0 && i < d->rehashidx) {
|
||||
/* Moreover, if we are currently out of range in the second
|
||||
* table, there will be no elements in both tables up to
|
||||
* the current rehashing index, so we jump if possible.
|
||||
* (this happens when going from big to small table). */
|
||||
if (i >= d->ht[1].size) i = d->rehashidx;
|
||||
continue;
|
||||
}
|
||||
if (i >= d->ht[j].size) continue; /* Out of range for this table. */
|
||||
dictEntry *he = d->ht[j].table[i];
|
||||
while (he) {
|
||||
/* Collect all the elements of the buckets found non
|
||||
* empty while iterating. */
|
||||
*des = he;
|
||||
des++;
|
||||
he = he->next;
|
||||
stored++;
|
||||
if (stored == count) return stored;
|
||||
}
|
||||
}
|
||||
i = (i+1) & maxsizemask;
|
||||
}
|
||||
return stored; /* Never reached. */
|
||||
}
|
||||
|
||||
|
||||
/* This function samples the dictionary to return a few keys from random
|
||||
* locations.
|
||||
*
|
||||
@@ -788,7 +713,7 @@ unsigned int dictGetSomeKeys(dict *d, dictEntry **des, unsigned int count) {
|
||||
/* Invariant of the dict.c rehashing: up to the indexes already
|
||||
* visited in ht[0] during the rehashing, there are no populated
|
||||
* buckets, so we can skip ht[0] for indexes between 0 and idx-1. */
|
||||
if (tables == 2 && j == 0 && i < d->rehashidx) {
|
||||
if (tables == 2 && j == 0 && i < (unsigned int) d->rehashidx) {
|
||||
/* Moreover, if we are currently out of range in the second
|
||||
* table, there will be no elements in both tables up to
|
||||
* the current rehashing index, so we jump if possible.
|
||||
|
||||
@@ -165,7 +165,6 @@ dictEntry *dictNext(dictIterator *iter);
|
||||
void dictReleaseIterator(dictIterator *iter);
|
||||
dictEntry *dictGetRandomKey(dict *d);
|
||||
unsigned int dictGetSomeKeys(dict *d, dictEntry **des, unsigned int count);
|
||||
unsigned int dictGetRandomKeys(dict *d, dictEntry **des, unsigned int count);
|
||||
void dictPrintStats(dict *d);
|
||||
unsigned int dictGenHashFunction(const void *key, int len);
|
||||
unsigned int dictGenCaseHashFunction(const unsigned char *buf, int len);
|
||||
|
||||
+42
-11
@@ -135,23 +135,49 @@ redisClient *createClient(int fd) {
|
||||
* returns REDIS_OK, and make sure to install the write handler in our event
|
||||
* loop so that when the socket is writable new data gets written.
|
||||
*
|
||||
* If the client should not receive new data, because it is a fake client,
|
||||
* a master, a slave not yet online, or because the setup of the write handler
|
||||
* failed, the function returns REDIS_ERR.
|
||||
* If the client should not receive new data, because it is a fake client
|
||||
* (used to load AOF in memory), a master or because the setup of the write
|
||||
* handler failed, the function returns REDIS_ERR.
|
||||
*
|
||||
* The function may return REDIS_OK without actually installing the write
|
||||
* event handler in the following cases:
|
||||
*
|
||||
* 1) The event handler should already be installed since the output buffer
|
||||
* already contained something.
|
||||
* 2) The client is a slave but not yet online, so we want to just accumulate
|
||||
* writes in the buffer but not actually sending them yet.
|
||||
*
|
||||
* Typically gets called every time a reply is built, before adding more
|
||||
* data to the clients output buffers. If the function returns REDIS_ERR no
|
||||
* data should be appended to the output buffers. */
|
||||
int prepareClientToWrite(redisClient *c) {
|
||||
/* If it's the Lua client we always return ok without installing any
|
||||
* handler since there is no socket at all. */
|
||||
if (c->flags & REDIS_LUA_CLIENT) return REDIS_OK;
|
||||
|
||||
/* Masters don't receive replies, unless REDIS_MASTER_FORCE_REPLY flag
|
||||
* is set. */
|
||||
if ((c->flags & REDIS_MASTER) &&
|
||||
!(c->flags & REDIS_MASTER_FORCE_REPLY)) return REDIS_ERR;
|
||||
if (c->fd <= 0) return REDIS_ERR; /* Fake client */
|
||||
|
||||
if (c->fd <= 0) return REDIS_ERR; /* Fake client for AOF loading. */
|
||||
|
||||
/* Only install the handler if not already installed and, in case of
|
||||
* slaves, if the client can actually receive writes. */
|
||||
if (c->bufpos == 0 && listLength(c->reply) == 0 &&
|
||||
(c->replstate == REDIS_REPL_NONE ||
|
||||
c->replstate == REDIS_REPL_ONLINE) &&
|
||||
aeCreateFileEvent(server.el, c->fd, AE_WRITABLE,
|
||||
sendReplyToClient, c) == AE_ERR) return REDIS_ERR;
|
||||
(c->replstate == REDIS_REPL_ONLINE && !c->repl_put_online_on_ack)))
|
||||
{
|
||||
/* Try to install the write handler. */
|
||||
if (aeCreateFileEvent(server.el, c->fd, AE_WRITABLE,
|
||||
sendReplyToClient, c) == AE_ERR)
|
||||
{
|
||||
freeClientAsync(c);
|
||||
return REDIS_ERR;
|
||||
}
|
||||
}
|
||||
|
||||
/* Authorize the caller to queue in the output buffer of this client. */
|
||||
return REDIS_OK;
|
||||
}
|
||||
|
||||
@@ -771,7 +797,7 @@ void freeClient(redisClient *c) {
|
||||
* a context where calling freeClient() is not possible, because the client
|
||||
* should be valid for the continuation of the flow of the program. */
|
||||
void freeClientAsync(redisClient *c) {
|
||||
if (c->flags & REDIS_CLOSE_ASAP) return;
|
||||
if (c->flags & REDIS_CLOSE_ASAP || c->flags & REDIS_LUA_CLIENT) return;
|
||||
c->flags |= REDIS_CLOSE_ASAP;
|
||||
listAddNodeTail(server.clients_to_close,c);
|
||||
}
|
||||
@@ -948,7 +974,7 @@ int processInlineBuffer(redisClient *c) {
|
||||
/* Helper function. Trims query buffer to make the function that processes
|
||||
* multi bulk requests idempotent. */
|
||||
static void setProtocolError(redisClient *c, int pos) {
|
||||
if (server.verbosity >= REDIS_VERBOSE) {
|
||||
if (server.verbosity <= REDIS_VERBOSE) {
|
||||
sds client = catClientInfoString(sdsempty(),c);
|
||||
redisLog(REDIS_VERBOSE,
|
||||
"Protocol error from client: %s", client);
|
||||
@@ -1689,7 +1715,9 @@ void pauseClients(mstime_t end) {
|
||||
/* Return non-zero if clients are currently paused. As a side effect the
|
||||
* function checks if the pause time was reached and clear it. */
|
||||
int clientsArePaused(void) {
|
||||
if (server.clients_paused && server.clients_pause_end_time < server.mstime) {
|
||||
if (server.clients_paused &&
|
||||
server.clients_pause_end_time < server.mstime)
|
||||
{
|
||||
listNode *ln;
|
||||
listIter li;
|
||||
redisClient *c;
|
||||
@@ -1702,7 +1730,10 @@ int clientsArePaused(void) {
|
||||
while ((ln = listNext(&li)) != NULL) {
|
||||
c = listNodeValue(ln);
|
||||
|
||||
if (c->flags & REDIS_SLAVE) continue;
|
||||
/* Don't touch slaves and blocked clients. The latter pending
|
||||
* requests be processed when unblocked. */
|
||||
if (c->flags & (REDIS_SLAVE|REDIS_BLOCKED)) continue;
|
||||
c->flags |= REDIS_UNBLOCKED;
|
||||
listAddNodeTail(server.unblocked_clients,c);
|
||||
}
|
||||
}
|
||||
|
||||
+12
-22
@@ -923,8 +923,14 @@ int clientsCronHandleTimeout(redisClient *c) {
|
||||
mstime_t now_ms = mstime();
|
||||
|
||||
if (c->bpop.timeout != 0 && c->bpop.timeout < now_ms) {
|
||||
/* Handle blocking operation specific timeout. */
|
||||
replyToBlockedClientTimedOut(c);
|
||||
unblockClient(c);
|
||||
} else if (server.cluster_enabled) {
|
||||
/* Cluster: handle unblock & redirect of clients blocked
|
||||
* into keys no longer served by this server. */
|
||||
if (clusterRedirectBlockedClientIfNeeded(c))
|
||||
unblockClient(c);
|
||||
}
|
||||
}
|
||||
return 0;
|
||||
@@ -1258,7 +1264,7 @@ void beforeSleep(struct aeEventLoop *eventLoop) {
|
||||
REDIS_NOTUSED(eventLoop);
|
||||
|
||||
/* Call the Redis Cluster before sleep function. Note that this function
|
||||
* may change the state of Redis Cluster (frok ok to fail or vice versa),
|
||||
* may change the state of Redis Cluster (from ok to fail or vice versa),
|
||||
* so it's a good idea to call it before serving the unblocked clients
|
||||
* later in this function. */
|
||||
if (server.cluster_enabled) clusterBeforeSleep();
|
||||
@@ -2158,38 +2164,22 @@ int processCommand(redisClient *c) {
|
||||
* 2) The command has no key arguments. */
|
||||
if (server.cluster_enabled &&
|
||||
!(c->flags & REDIS_MASTER) &&
|
||||
!(c->flags & REDIS_LUA_CLIENT &&
|
||||
server.lua_caller->flags & REDIS_MASTER) &&
|
||||
!(c->cmd->getkeys_proc == NULL && c->cmd->firstkey == 0))
|
||||
{
|
||||
int hashslot;
|
||||
|
||||
if (server.cluster->state != REDIS_CLUSTER_OK) {
|
||||
flagTransaction(c);
|
||||
addReplySds(c,sdsnew("-CLUSTERDOWN The cluster is down. Use CLUSTER INFO for more information\r\n"));
|
||||
clusterRedirectClient(c,NULL,0,REDIS_CLUSTER_REDIR_DOWN_STATE);
|
||||
return REDIS_OK;
|
||||
} else {
|
||||
int error_code;
|
||||
clusterNode *n = getNodeByQuery(c,c->cmd,c->argv,c->argc,&hashslot,&error_code);
|
||||
if (n == NULL) {
|
||||
if (n == NULL || n != server.cluster->myself) {
|
||||
flagTransaction(c);
|
||||
if (error_code == REDIS_CLUSTER_REDIR_CROSS_SLOT) {
|
||||
addReplySds(c,sdsnew("-CROSSSLOT Keys in request don't hash to the same slot\r\n"));
|
||||
} else if (error_code == REDIS_CLUSTER_REDIR_UNSTABLE) {
|
||||
/* The request spawns mutliple keys in the same slot,
|
||||
* but the slot is not "stable" currently as there is
|
||||
* a migration or import in progress. */
|
||||
addReplySds(c,sdsnew("-TRYAGAIN Multiple keys request during rehashing of slot\r\n"));
|
||||
} else if (error_code == REDIS_CLUSTER_REDIR_DOWN) {
|
||||
addReplySds(c,sdsnew("-CLUSTERDOWN The cluster is down. Hash slot is unbound\r\n"));
|
||||
} else {
|
||||
redisPanic("getNodeByQuery() unknown error.");
|
||||
}
|
||||
return REDIS_OK;
|
||||
} else if (n != server.cluster->myself) {
|
||||
flagTransaction(c);
|
||||
addReplySds(c,sdscatprintf(sdsempty(),
|
||||
"-%s %d %s:%d\r\n",
|
||||
(error_code == REDIS_CLUSTER_REDIR_ASK) ? "ASK" : "MOVED",
|
||||
hashslot,n->ip,n->port));
|
||||
clusterRedirectClient(c,n,hashslot,error_code);
|
||||
return REDIS_OK;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1366,6 +1366,7 @@ void blockClient(redisClient *c, int btype);
|
||||
void unblockClient(redisClient *c);
|
||||
void replyToBlockedClientTimedOut(redisClient *c);
|
||||
int getTimeoutFromObjectOrReply(redisClient *c, robj *object, mstime_t *timeout, int unit);
|
||||
void disconnectAllBlockedClients(void);
|
||||
|
||||
/* Git SHA1 */
|
||||
char *redisGitSHA1(void);
|
||||
|
||||
+5
-2
@@ -652,7 +652,8 @@ void replconfCommand(redisClient *c) {
|
||||
*
|
||||
* It does a few things:
|
||||
*
|
||||
* 1) Put the slave in ONLINE state.
|
||||
* 1) Put the slave in ONLINE state (useless when the function is called
|
||||
* because state is already ONLINE but repl_put_online_on_ack is true).
|
||||
* 2) Make sure the writable event is re-installed, since calling the SYNC
|
||||
* command disables it, so that we can accumulate output buffer without
|
||||
* sending it to the slave.
|
||||
@@ -660,7 +661,7 @@ void replconfCommand(redisClient *c) {
|
||||
void putSlaveOnline(redisClient *slave) {
|
||||
slave->replstate = REDIS_REPL_ONLINE;
|
||||
slave->repl_put_online_on_ack = 0;
|
||||
slave->repl_ack_time = server.unixtime;
|
||||
slave->repl_ack_time = server.unixtime; /* Prevent false timeout. */
|
||||
if (aeCreateFileEvent(server.el, slave->fd, AE_WRITABLE,
|
||||
sendReplyToClient, slave) == AE_ERR) {
|
||||
redisLog(REDIS_WARNING,"Unable to register writable event for slave bulk transfer: %s", strerror(errno));
|
||||
@@ -773,6 +774,7 @@ void updateSlavesWaitingBgsave(int bgsaveerr, int type) {
|
||||
* is technically online now. */
|
||||
slave->replstate = REDIS_REPL_ONLINE;
|
||||
slave->repl_put_online_on_ack = 1;
|
||||
slave->repl_ack_time = server.unixtime; /* Timeout otherwise. */
|
||||
} else {
|
||||
if (bgsaveerr != REDIS_OK) {
|
||||
freeClient(slave);
|
||||
@@ -1437,6 +1439,7 @@ void replicationSetMaster(char *ip, int port) {
|
||||
server.masterhost = sdsnew(ip);
|
||||
server.masterport = port;
|
||||
if (server.master) freeClient(server.master);
|
||||
disconnectAllBlockedClients(); /* Clients blocked in master, now slave. */
|
||||
disconnectSlaves(); /* Force our slaves to resync with us as well. */
|
||||
replicationDiscardCachedMaster(); /* Don't try a PSYNC. */
|
||||
freeReplicationBacklog(); /* Don't allow our chained slaves to PSYNC. */
|
||||
|
||||
+12
-8
@@ -357,8 +357,9 @@ int luaRedisGenericCommand(lua_State *lua, int raise_error) {
|
||||
if (cmd->flags & REDIS_CMD_WRITE) server.lua_write_dirty = 1;
|
||||
|
||||
/* If this is a Redis Cluster node, we need to make sure Lua is not
|
||||
* trying to access non-local keys. */
|
||||
if (server.cluster_enabled) {
|
||||
* trying to access non-local keys, with the exception of commands
|
||||
* received from our master. */
|
||||
if (server.cluster_enabled && !(server.lua_caller->flags & REDIS_MASTER)) {
|
||||
/* Duplicate relevant flags in the lua client. */
|
||||
c->flags &= ~(REDIS_READONLY|REDIS_ASKING);
|
||||
c->flags |= server.lua_caller->flags & (REDIS_READONLY|REDIS_ASKING);
|
||||
@@ -611,11 +612,12 @@ void scriptingEnableGlobalsProtection(lua_State *lua) {
|
||||
|
||||
/* strict.lua from: http://metalua.luaforge.net/src/lib/strict.lua.html.
|
||||
* Modified to be adapted to Redis. */
|
||||
s[j++]="local dbg=debug\n";
|
||||
s[j++]="local mt = {}\n";
|
||||
s[j++]="setmetatable(_G, mt)\n";
|
||||
s[j++]="mt.__newindex = function (t, n, v)\n";
|
||||
s[j++]=" if debug.getinfo(2) then\n";
|
||||
s[j++]=" local w = debug.getinfo(2, \"S\").what\n";
|
||||
s[j++]=" if dbg.getinfo(2) then\n";
|
||||
s[j++]=" local w = dbg.getinfo(2, \"S\").what\n";
|
||||
s[j++]=" if w ~= \"main\" and w ~= \"C\" then\n";
|
||||
s[j++]=" error(\"Script attempted to create global variable '\"..tostring(n)..\"'\", 2)\n";
|
||||
s[j++]=" end\n";
|
||||
@@ -623,11 +625,12 @@ void scriptingEnableGlobalsProtection(lua_State *lua) {
|
||||
s[j++]=" rawset(t, n, v)\n";
|
||||
s[j++]="end\n";
|
||||
s[j++]="mt.__index = function (t, n)\n";
|
||||
s[j++]=" if debug.getinfo(2) and debug.getinfo(2, \"S\").what ~= \"C\" then\n";
|
||||
s[j++]=" if dbg.getinfo(2) and dbg.getinfo(2, \"S\").what ~= \"C\" then\n";
|
||||
s[j++]=" error(\"Script attempted to access unexisting global variable '\"..tostring(n)..\"'\", 2)\n";
|
||||
s[j++]=" end\n";
|
||||
s[j++]=" return rawget(t, n)\n";
|
||||
s[j++]="end\n";
|
||||
s[j++]="debug = nil\n";
|
||||
s[j++]=NULL;
|
||||
|
||||
for (j = 0; s[j] != NULL; j++) code = sdscatlen(code,s[j],strlen(s[j]));
|
||||
@@ -731,10 +734,11 @@ void scriptingInit(void) {
|
||||
* information about the caller, that's what makes sense from the point
|
||||
* of view of the user debugging a script. */
|
||||
{
|
||||
char *errh_func = "function __redis__err__handler(err)\n"
|
||||
" local i = debug.getinfo(2,'nSl')\n"
|
||||
char *errh_func = "local dbg = debug\n"
|
||||
"function __redis__err__handler(err)\n"
|
||||
" local i = dbg.getinfo(2,'nSl')\n"
|
||||
" if i and i.what == 'C' then\n"
|
||||
" i = debug.getinfo(3,'nSl')\n"
|
||||
" i = dbg.getinfo(3,'nSl')\n"
|
||||
" end\n"
|
||||
" if i then\n"
|
||||
" return i.source .. ':' .. i.currentline .. ': ' .. err\n"
|
||||
|
||||
@@ -71,7 +71,7 @@ sds sdsempty(void) {
|
||||
return sdsnewlen("",0);
|
||||
}
|
||||
|
||||
/* Create a new sds string starting from a null termined C string. */
|
||||
/* Create a new sds string starting from a null terminated C string. */
|
||||
sds sdsnew(const char *init) {
|
||||
size_t initlen = (init == NULL) ? 0 : strlen(init);
|
||||
return sdsnewlen(init, initlen);
|
||||
@@ -557,7 +557,7 @@ sds sdscatfmt(sds s, char const *fmt, ...) {
|
||||
* Example:
|
||||
*
|
||||
* s = sdsnew("AA...AA.a.aa.aHelloWorld :::");
|
||||
* s = sdstrim(s,"A. :");
|
||||
* s = sdstrim(s,"Aa. :");
|
||||
* printf("%s\n", s);
|
||||
*
|
||||
* Output will be just "Hello World".
|
||||
@@ -1083,6 +1083,7 @@ int main(void) {
|
||||
int oldfree;
|
||||
|
||||
sdsfree(x);
|
||||
sdsfree(y);
|
||||
x = sdsnew("0");
|
||||
sh = (void*) (x-(sizeof(struct sdshdr)));
|
||||
test_cond("sdsnew() free/len buffers", sh->len == 1 && sh->free == 0);
|
||||
@@ -1095,6 +1096,8 @@ int main(void) {
|
||||
test_cond("sdsIncrLen() -- content", x[0] == '0' && x[1] == '1');
|
||||
test_cond("sdsIncrLen() -- len", sh->len == 2);
|
||||
test_cond("sdsIncrLen() -- free", sh->free == oldfree-1);
|
||||
|
||||
sdsfree(x);
|
||||
}
|
||||
}
|
||||
test_report()
|
||||
|
||||
+58
-4
@@ -923,6 +923,7 @@ sentinelRedisInstance *createSentinelRedisInstance(char *name, int flags, char *
|
||||
else if (flags & SRI_SENTINEL) table = master->sentinels;
|
||||
sdsname = sdsnew(name);
|
||||
if (dictFind(table,sdsname)) {
|
||||
releaseSentinelAddr(addr);
|
||||
sdsfree(sdsname);
|
||||
errno = EBUSY;
|
||||
return NULL;
|
||||
@@ -1269,10 +1270,7 @@ int sentinelResetMasterAndChangeAddress(sentinelRedisInstance *master, char *ip,
|
||||
slave = createSentinelRedisInstance(NULL,SRI_SLAVE,slaves[j]->ip,
|
||||
slaves[j]->port, master->quorum, master);
|
||||
releaseSentinelAddr(slaves[j]);
|
||||
if (slave) {
|
||||
sentinelEvent(REDIS_NOTICE,"+slave",slave,"%@");
|
||||
sentinelFlushConfig();
|
||||
}
|
||||
if (slave) sentinelEvent(REDIS_NOTICE,"+slave",slave,"%@");
|
||||
}
|
||||
zfree(slaves);
|
||||
|
||||
@@ -1844,6 +1842,7 @@ void sentinelRefreshInstanceInfo(sentinelRedisInstance *ri, const char *info) {
|
||||
atoi(port), ri->quorum, ri)) != NULL)
|
||||
{
|
||||
sentinelEvent(REDIS_NOTICE,"+slave",slave,"%@");
|
||||
sentinelFlushConfig();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -2615,6 +2614,31 @@ sentinelRedisInstance *sentinelGetMasterByNameOrReplyError(redisClient *c,
|
||||
return ri;
|
||||
}
|
||||
|
||||
#define SENTINEL_ISQR_OK 0
|
||||
#define SENTINEL_ISQR_NOQUORUM (1<<0)
|
||||
#define SENTINEL_ISQR_NOAUTH (1<<1)
|
||||
int sentinelIsQuorumReachable(sentinelRedisInstance *master, int *usableptr) {
|
||||
dictIterator *di;
|
||||
dictEntry *de;
|
||||
int usable = 1; /* Number of usable Sentinels. Init to 1 to count myself. */
|
||||
int result = SENTINEL_ISQR_OK;
|
||||
int voters = dictSize(master->sentinels)+1; /* Known Sentinels + myself. */
|
||||
|
||||
di = dictGetIterator(master->sentinels);
|
||||
while((de = dictNext(di)) != NULL) {
|
||||
sentinelRedisInstance *ri = dictGetVal(de);
|
||||
|
||||
if (ri->flags & (SRI_S_DOWN|SRI_O_DOWN)) continue;
|
||||
usable++;
|
||||
}
|
||||
dictReleaseIterator(di);
|
||||
|
||||
if (usable < (int)master->quorum) result |= SENTINEL_ISQR_NOQUORUM;
|
||||
if (usable < voters/2+1) result |= SENTINEL_ISQR_NOAUTH;
|
||||
if (usableptr) *usableptr = usable;
|
||||
return result;
|
||||
}
|
||||
|
||||
void sentinelCommand(redisClient *c) {
|
||||
if (!strcasecmp(c->argv[1]->ptr,"masters")) {
|
||||
/* SENTINEL MASTERS */
|
||||
@@ -2765,6 +2789,10 @@ void sentinelCommand(redisClient *c) {
|
||||
sentinelEvent(REDIS_WARNING,"+monitor",ri,"%@ quorum %d",ri->quorum);
|
||||
addReply(c,shared.ok);
|
||||
}
|
||||
} else if (!strcasecmp(c->argv[1]->ptr,"flushconfig")) {
|
||||
sentinelFlushConfig();
|
||||
addReply(c,shared.ok);
|
||||
return;
|
||||
} else if (!strcasecmp(c->argv[1]->ptr,"remove")) {
|
||||
/* SENTINEL REMOVE <name> */
|
||||
sentinelRedisInstance *ri;
|
||||
@@ -2775,6 +2803,32 @@ void sentinelCommand(redisClient *c) {
|
||||
dictDelete(sentinel.masters,c->argv[2]->ptr);
|
||||
sentinelFlushConfig();
|
||||
addReply(c,shared.ok);
|
||||
} else if (!strcasecmp(c->argv[1]->ptr,"ckquorum")) {
|
||||
/* SENTINEL CKQUORUM <name> */
|
||||
sentinelRedisInstance *ri;
|
||||
int usable;
|
||||
|
||||
if ((ri = sentinelGetMasterByNameOrReplyError(c,c->argv[2]))
|
||||
== NULL) return;
|
||||
int result = sentinelIsQuorumReachable(ri,&usable);
|
||||
if (result == SENTINEL_ISQR_OK) {
|
||||
addReplySds(c, sdscatfmt(sdsempty(),
|
||||
"+OK %i usable Sentinels. Quorum and failover authorization "
|
||||
"can be reached\r\n",usable));
|
||||
} else {
|
||||
sds e = sdscatfmt(sdsempty(),
|
||||
"-NOQUORUM %i usable Sentinels. ",usable);
|
||||
if (result & SENTINEL_ISQR_NOQUORUM)
|
||||
e = sdscat(e,"Not enough available Sentinels to reach the"
|
||||
" specified quorum for this master");
|
||||
if (result & SENTINEL_ISQR_NOAUTH) {
|
||||
if (result & SENTINEL_ISQR_NOQUORUM) e = sdscat(e,". ");
|
||||
e = sdscat(e, "Not enough available Sentinels to reach the"
|
||||
" majority and authorize a failover");
|
||||
}
|
||||
e = sdscat(e,"\r\n");
|
||||
addReplySds(c,e);
|
||||
}
|
||||
} else if (!strcasecmp(c->argv[1]->ptr,"set")) {
|
||||
if (c->argc < 3 || c->argc % 2 == 0) goto numargserr;
|
||||
sentinelSetCommand(c);
|
||||
|
||||
+7
-7
@@ -23,7 +23,7 @@ A million repetitions of "a"
|
||||
|
||||
#include <stdio.h>
|
||||
#include <string.h>
|
||||
#include <sys/types.h> /* for u_int*_t */
|
||||
#include <stdint.h>
|
||||
#include "solarisfixes.h"
|
||||
#include "sha1.h"
|
||||
#include "config.h"
|
||||
@@ -53,12 +53,12 @@ A million repetitions of "a"
|
||||
|
||||
/* Hash a single 512-bit block. This is the core of the algorithm. */
|
||||
|
||||
void SHA1Transform(u_int32_t state[5], const unsigned char buffer[64])
|
||||
void SHA1Transform(uint32_t state[5], const unsigned char buffer[64])
|
||||
{
|
||||
u_int32_t a, b, c, d, e;
|
||||
uint32_t a, b, c, d, e;
|
||||
typedef union {
|
||||
unsigned char c[64];
|
||||
u_int32_t l[16];
|
||||
uint32_t l[16];
|
||||
} CHAR64LONG16;
|
||||
#ifdef SHA1HANDSOFF
|
||||
CHAR64LONG16 block[1]; /* use array to appear as a pointer */
|
||||
@@ -128,9 +128,9 @@ void SHA1Init(SHA1_CTX* context)
|
||||
|
||||
/* Run your data through this. */
|
||||
|
||||
void SHA1Update(SHA1_CTX* context, const unsigned char* data, u_int32_t len)
|
||||
void SHA1Update(SHA1_CTX* context, const unsigned char* data, uint32_t len)
|
||||
{
|
||||
u_int32_t i, j;
|
||||
uint32_t i, j;
|
||||
|
||||
j = context->count[0];
|
||||
if ((context->count[0] += len << 3) < j)
|
||||
@@ -168,7 +168,7 @@ void SHA1Final(unsigned char digest[20], SHA1_CTX* context)
|
||||
|
||||
for (i = 0; i < 2; i++)
|
||||
{
|
||||
u_int32_t t = context->count[i];
|
||||
uint32_t t = context->count[i];
|
||||
int j;
|
||||
|
||||
for (j = 0; j < 4; t >>= 8, j++)
|
||||
|
||||
+4
-4
@@ -6,12 +6,12 @@ By Steve Reid <steve@edmweb.com>
|
||||
*/
|
||||
|
||||
typedef struct {
|
||||
u_int32_t state[5];
|
||||
u_int32_t count[2];
|
||||
uint32_t state[5];
|
||||
uint32_t count[2];
|
||||
unsigned char buffer[64];
|
||||
} SHA1_CTX;
|
||||
|
||||
void SHA1Transform(u_int32_t state[5], const unsigned char buffer[64]);
|
||||
void SHA1Transform(uint32_t state[5], const unsigned char buffer[64]);
|
||||
void SHA1Init(SHA1_CTX* context);
|
||||
void SHA1Update(SHA1_CTX* context, const unsigned char* data, u_int32_t len);
|
||||
void SHA1Update(SHA1_CTX* context, const unsigned char* data, uint32_t len);
|
||||
void SHA1Final(unsigned char digest[20], SHA1_CTX* context);
|
||||
|
||||
+1
-1
@@ -209,7 +209,7 @@ void sortCommand(redisClient *c) {
|
||||
}
|
||||
|
||||
/* Create a list of operations to perform for every sorted element.
|
||||
* Operations can be GET/DEL/INCR/DECR */
|
||||
* Operations can be GET */
|
||||
operations = listCreate();
|
||||
listSetFreeMethod(operations,zfree);
|
||||
j = 2; /* options start at argv[2] */
|
||||
|
||||
+1
-1
@@ -321,7 +321,7 @@ void smoveCommand(redisClient *c) {
|
||||
|
||||
/* If srcset and dstset are equal, SMOVE is a no-op */
|
||||
if (srcset == dstset) {
|
||||
addReply(c,shared.cone);
|
||||
addReply(c,setTypeIsMember(srcset,ele) ? shared.cone : shared.czero);
|
||||
return;
|
||||
}
|
||||
|
||||
|
||||
+77
-15
@@ -1171,33 +1171,82 @@ void zsetConvert(robj *zobj, int encoding) {
|
||||
*----------------------------------------------------------------------------*/
|
||||
|
||||
/* This generic command implements both ZADD and ZINCRBY. */
|
||||
void zaddGenericCommand(redisClient *c, int incr) {
|
||||
#define ZADD_NONE 0
|
||||
#define ZADD_INCR (1<<0) /* Increment the score instead of setting it. */
|
||||
#define ZADD_NX (1<<1) /* Don't touch elements not already existing. */
|
||||
#define ZADD_XX (1<<2) /* Only touch elements already exisitng. */
|
||||
#define ZADD_CH (1<<3) /* Return num of elements added or updated. */
|
||||
void zaddGenericCommand(redisClient *c, int flags) {
|
||||
static char *nanerr = "resulting score is not a number (NaN)";
|
||||
robj *key = c->argv[1];
|
||||
robj *ele;
|
||||
robj *zobj;
|
||||
robj *curobj;
|
||||
double score = 0, *scores = NULL, curscore = 0.0;
|
||||
int j, elements = (c->argc-2)/2;
|
||||
int added = 0, updated = 0;
|
||||
int j, elements;
|
||||
int scoreidx = 0;
|
||||
/* The following vars are used in order to track what the command actually
|
||||
* did during the execution, to reply to the client and to trigger the
|
||||
* notification of keyspace change. */
|
||||
int added = 0; /* Number of new elements added. */
|
||||
int updated = 0; /* Number of elements with updated score. */
|
||||
int processed = 0; /* Number of elements processed, may remain zero with
|
||||
options like XX. */
|
||||
|
||||
if (c->argc % 2) {
|
||||
/* Parse options. At the end 'scoreidx' is set to the argument position
|
||||
* of the score of the first score-element pair. */
|
||||
scoreidx = 2;
|
||||
while(scoreidx < c->argc) {
|
||||
char *opt = c->argv[scoreidx]->ptr;
|
||||
if (!strcasecmp(opt,"nx")) flags |= ZADD_NX;
|
||||
else if (!strcasecmp(opt,"xx")) flags |= ZADD_XX;
|
||||
else if (!strcasecmp(opt,"ch")) flags |= ZADD_CH;
|
||||
else if (!strcasecmp(opt,"incr")) flags |= ZADD_INCR;
|
||||
else break;
|
||||
scoreidx++;
|
||||
}
|
||||
|
||||
/* Turn options into simple to check vars. */
|
||||
int incr = (flags & ZADD_INCR) != 0;
|
||||
int nx = (flags & ZADD_NX) != 0;
|
||||
int xx = (flags & ZADD_XX) != 0;
|
||||
int ch = (flags & ZADD_CH) != 0;
|
||||
|
||||
/* After the options, we expect to have an even number of args, since
|
||||
* we expect any number of score-element pairs. */
|
||||
elements = c->argc-scoreidx;
|
||||
if (elements % 2) {
|
||||
addReply(c,shared.syntaxerr);
|
||||
return;
|
||||
}
|
||||
elements /= 2; /* Now this holds the number of score-element pairs. */
|
||||
|
||||
/* Check for incompatible options. */
|
||||
if (nx && xx) {
|
||||
addReplyError(c,
|
||||
"XX and NX options at the same time are not compatible");
|
||||
return;
|
||||
}
|
||||
|
||||
if (incr && elements > 1) {
|
||||
addReplyError(c,
|
||||
"INCR option supports a single increment-element pair");
|
||||
return;
|
||||
}
|
||||
|
||||
/* Start parsing all the scores, we need to emit any syntax error
|
||||
* before executing additions to the sorted set, as the command should
|
||||
* either execute fully or nothing at all. */
|
||||
scores = zmalloc(sizeof(double)*elements);
|
||||
for (j = 0; j < elements; j++) {
|
||||
if (getDoubleFromObjectOrReply(c,c->argv[2+j*2],&scores[j],NULL)
|
||||
if (getDoubleFromObjectOrReply(c,c->argv[scoreidx+j*2],&scores[j],NULL)
|
||||
!= REDIS_OK) goto cleanup;
|
||||
}
|
||||
|
||||
/* Lookup the key and create the sorted set if does not exist. */
|
||||
zobj = lookupKeyWrite(c->db,key);
|
||||
if (zobj == NULL) {
|
||||
if (xx) goto reply_to_client; /* No key + XX option: nothing to do. */
|
||||
if (server.zset_max_ziplist_entries == 0 ||
|
||||
server.zset_max_ziplist_value < sdslen(c->argv[3]->ptr))
|
||||
{
|
||||
@@ -1220,8 +1269,9 @@ void zaddGenericCommand(redisClient *c, int incr) {
|
||||
unsigned char *eptr;
|
||||
|
||||
/* Prefer non-encoded element when dealing with ziplists. */
|
||||
ele = c->argv[3+j*2];
|
||||
ele = c->argv[scoreidx+1+j*2];
|
||||
if ((eptr = zzlFind(zobj->ptr,ele,&curscore)) != NULL) {
|
||||
if (nx) continue;
|
||||
if (incr) {
|
||||
score += curscore;
|
||||
if (isnan(score)) {
|
||||
@@ -1237,7 +1287,8 @@ void zaddGenericCommand(redisClient *c, int incr) {
|
||||
server.dirty++;
|
||||
updated++;
|
||||
}
|
||||
} else {
|
||||
processed++;
|
||||
} else if (!xx) {
|
||||
/* Optimize: check if the element is too large or the list
|
||||
* becomes too long *before* executing zzlInsert. */
|
||||
zobj->ptr = zzlInsert(zobj->ptr,ele,score);
|
||||
@@ -1247,15 +1298,18 @@ void zaddGenericCommand(redisClient *c, int incr) {
|
||||
zsetConvert(zobj,REDIS_ENCODING_SKIPLIST);
|
||||
server.dirty++;
|
||||
added++;
|
||||
processed++;
|
||||
}
|
||||
} else if (zobj->encoding == REDIS_ENCODING_SKIPLIST) {
|
||||
zset *zs = zobj->ptr;
|
||||
zskiplistNode *znode;
|
||||
dictEntry *de;
|
||||
|
||||
ele = c->argv[3+j*2] = tryObjectEncoding(c->argv[3+j*2]);
|
||||
ele = c->argv[scoreidx+1+j*2] =
|
||||
tryObjectEncoding(c->argv[scoreidx+1+j*2]);
|
||||
de = dictFind(zs->dict,ele);
|
||||
if (de != NULL) {
|
||||
if (nx) continue;
|
||||
curobj = dictGetKey(de);
|
||||
curscore = *(double*)dictGetVal(de);
|
||||
|
||||
@@ -1280,22 +1334,30 @@ void zaddGenericCommand(redisClient *c, int incr) {
|
||||
server.dirty++;
|
||||
updated++;
|
||||
}
|
||||
} else {
|
||||
processed++;
|
||||
} else if (!xx) {
|
||||
znode = zslInsert(zs->zsl,score,ele);
|
||||
incrRefCount(ele); /* Inserted in skiplist. */
|
||||
redisAssertWithInfo(c,NULL,dictAdd(zs->dict,ele,&znode->score) == DICT_OK);
|
||||
incrRefCount(ele); /* Added to dictionary. */
|
||||
server.dirty++;
|
||||
added++;
|
||||
processed++;
|
||||
}
|
||||
} else {
|
||||
redisPanic("Unknown sorted set encoding");
|
||||
}
|
||||
}
|
||||
if (incr) /* ZINCRBY */
|
||||
addReplyDouble(c,score);
|
||||
else /* ZADD */
|
||||
addReplyLongLong(c,added);
|
||||
|
||||
reply_to_client:
|
||||
if (incr) { /* ZINCRBY or INCR option. */
|
||||
if (processed)
|
||||
addReplyDouble(c,score);
|
||||
else
|
||||
addReply(c,shared.nullbulk);
|
||||
} else { /* ZADD. */
|
||||
addReplyLongLong(c,ch ? added+updated : added);
|
||||
}
|
||||
|
||||
cleanup:
|
||||
zfree(scores);
|
||||
@@ -1307,11 +1369,11 @@ cleanup:
|
||||
}
|
||||
|
||||
void zaddCommand(redisClient *c) {
|
||||
zaddGenericCommand(c,0);
|
||||
zaddGenericCommand(c,ZADD_NONE);
|
||||
}
|
||||
|
||||
void zincrbyCommand(redisClient *c) {
|
||||
zaddGenericCommand(c,1);
|
||||
zaddGenericCommand(c,ZADD_INCR);
|
||||
}
|
||||
|
||||
void zremCommand(redisClient *c) {
|
||||
|
||||
+1
-1
@@ -1 +1 @@
|
||||
#define REDIS_VERSION "2.9.105"
|
||||
#define REDIS_VERSION "3.0.2"
|
||||
|
||||
@@ -17,6 +17,7 @@ proc main {} {
|
||||
}
|
||||
run_tests
|
||||
cleanup
|
||||
end_tests
|
||||
}
|
||||
|
||||
if {[catch main e]} {
|
||||
|
||||
@@ -0,0 +1,192 @@
|
||||
# Check the manual failover
|
||||
|
||||
source "../tests/includes/init-tests.tcl"
|
||||
|
||||
test "Create a 5 nodes cluster" {
|
||||
create_cluster 5 5
|
||||
}
|
||||
|
||||
test "Cluster is up" {
|
||||
assert_cluster_state ok
|
||||
}
|
||||
|
||||
test "Cluster is writable" {
|
||||
cluster_write_test 0
|
||||
}
|
||||
|
||||
test "Instance #5 is a slave" {
|
||||
assert {[RI 5 role] eq {slave}}
|
||||
}
|
||||
|
||||
test "Instance #5 synced with the master" {
|
||||
wait_for_condition 1000 50 {
|
||||
[RI 5 master_link_status] eq {up}
|
||||
} else {
|
||||
fail "Instance #5 master link status is not up"
|
||||
}
|
||||
}
|
||||
|
||||
set current_epoch [CI 1 cluster_current_epoch]
|
||||
|
||||
set numkeys 50000
|
||||
set numops 10000
|
||||
set cluster [redis_cluster 127.0.0.1:[get_instance_attrib redis 0 port]]
|
||||
catch {unset content}
|
||||
array set content {}
|
||||
|
||||
test "Send CLUSTER FAILOVER to #5, during load" {
|
||||
for {set j 0} {$j < $numops} {incr j} {
|
||||
# Write random data to random list.
|
||||
set listid [randomInt $numkeys]
|
||||
set key "key:$listid"
|
||||
set ele [randomValue]
|
||||
# We write both with Lua scripts and with plain commands.
|
||||
# This way we are able to stress Lua -> Redis command invocation
|
||||
# as well, that has tests to prevent Lua to write into wrong
|
||||
# hash slots.
|
||||
if {$listid % 2} {
|
||||
$cluster rpush $key $ele
|
||||
} else {
|
||||
$cluster eval {redis.call("rpush",KEYS[1],ARGV[1])} 1 $key $ele
|
||||
}
|
||||
lappend content($key) $ele
|
||||
|
||||
if {($j % 1000) == 0} {
|
||||
puts -nonewline W; flush stdout
|
||||
}
|
||||
|
||||
if {$j == $numops/2} {R 5 cluster failover}
|
||||
}
|
||||
}
|
||||
|
||||
test "Wait for failover" {
|
||||
wait_for_condition 1000 50 {
|
||||
[CI 1 cluster_current_epoch] > $current_epoch
|
||||
} else {
|
||||
fail "No failover detected"
|
||||
}
|
||||
}
|
||||
|
||||
test "Cluster should eventually be up again" {
|
||||
assert_cluster_state ok
|
||||
}
|
||||
|
||||
test "Cluster is writable" {
|
||||
cluster_write_test 1
|
||||
}
|
||||
|
||||
test "Instance #5 is now a master" {
|
||||
assert {[RI 5 role] eq {master}}
|
||||
}
|
||||
|
||||
test "Verify $numkeys keys for consistency with logical content" {
|
||||
# Check that the Redis Cluster content matches our logical content.
|
||||
foreach {key value} [array get content] {
|
||||
assert {[$cluster lrange $key 0 -1] eq $value}
|
||||
}
|
||||
}
|
||||
|
||||
test "Instance #0 gets converted into a slave" {
|
||||
wait_for_condition 1000 50 {
|
||||
[RI 0 role] eq {slave}
|
||||
} else {
|
||||
fail "Old master was not converted into slave"
|
||||
}
|
||||
}
|
||||
|
||||
## Check that manual failover does not happen if we can't talk with the master.
|
||||
|
||||
source "../tests/includes/init-tests.tcl"
|
||||
|
||||
test "Create a 5 nodes cluster" {
|
||||
create_cluster 5 5
|
||||
}
|
||||
|
||||
test "Cluster is up" {
|
||||
assert_cluster_state ok
|
||||
}
|
||||
|
||||
test "Cluster is writable" {
|
||||
cluster_write_test 0
|
||||
}
|
||||
|
||||
test "Instance #5 is a slave" {
|
||||
assert {[RI 5 role] eq {slave}}
|
||||
}
|
||||
|
||||
test "Instance #5 synced with the master" {
|
||||
wait_for_condition 1000 50 {
|
||||
[RI 5 master_link_status] eq {up}
|
||||
} else {
|
||||
fail "Instance #5 master link status is not up"
|
||||
}
|
||||
}
|
||||
|
||||
test "Make instance #0 unreachable without killing it" {
|
||||
R 0 deferred 1
|
||||
R 0 DEBUG SLEEP 10
|
||||
}
|
||||
|
||||
test "Send CLUSTER FAILOVER to instance #5" {
|
||||
R 5 cluster failover
|
||||
}
|
||||
|
||||
test "Instance #5 is still a slave after some time (no failover)" {
|
||||
after 5000
|
||||
assert {[RI 5 role] eq {master}}
|
||||
}
|
||||
|
||||
test "Wait for instance #0 to return back alive" {
|
||||
R 0 deferred 0
|
||||
assert {[R 0 read] eq {OK}}
|
||||
}
|
||||
|
||||
## Check with "force" failover happens anyway.
|
||||
|
||||
source "../tests/includes/init-tests.tcl"
|
||||
|
||||
test "Create a 5 nodes cluster" {
|
||||
create_cluster 5 5
|
||||
}
|
||||
|
||||
test "Cluster is up" {
|
||||
assert_cluster_state ok
|
||||
}
|
||||
|
||||
test "Cluster is writable" {
|
||||
cluster_write_test 0
|
||||
}
|
||||
|
||||
test "Instance #5 is a slave" {
|
||||
assert {[RI 5 role] eq {slave}}
|
||||
}
|
||||
|
||||
test "Instance #5 synced with the master" {
|
||||
wait_for_condition 1000 50 {
|
||||
[RI 5 master_link_status] eq {up}
|
||||
} else {
|
||||
fail "Instance #5 master link status is not up"
|
||||
}
|
||||
}
|
||||
|
||||
test "Make instance #0 unreachable without killing it" {
|
||||
R 0 deferred 1
|
||||
R 0 DEBUG SLEEP 10
|
||||
}
|
||||
|
||||
test "Send CLUSTER FAILOVER to instance #5" {
|
||||
R 5 cluster failover force
|
||||
}
|
||||
|
||||
test "Instance #5 is a master after some time" {
|
||||
wait_for_condition 1000 50 {
|
||||
[RI 5 role] eq {master}
|
||||
} else {
|
||||
fail "Instance #5 is not a master after some time regardless of FORCE"
|
||||
}
|
||||
}
|
||||
|
||||
test "Wait for instance #0 to return back alive" {
|
||||
R 0 deferred 0
|
||||
assert {[R 0 read] eq {OK}}
|
||||
}
|
||||
@@ -0,0 +1,59 @@
|
||||
# Manual takeover test
|
||||
|
||||
source "../tests/includes/init-tests.tcl"
|
||||
|
||||
test "Create a 5 nodes cluster" {
|
||||
create_cluster 5 5
|
||||
}
|
||||
|
||||
test "Cluster is up" {
|
||||
assert_cluster_state ok
|
||||
}
|
||||
|
||||
test "Cluster is writable" {
|
||||
cluster_write_test 0
|
||||
}
|
||||
|
||||
test "Killing majority of master nodes" {
|
||||
kill_instance redis 0
|
||||
kill_instance redis 1
|
||||
kill_instance redis 2
|
||||
}
|
||||
|
||||
test "Cluster should eventually be down" {
|
||||
assert_cluster_state fail
|
||||
}
|
||||
|
||||
test "Use takeover to bring slaves back" {
|
||||
R 5 cluster failover takeover
|
||||
R 6 cluster failover takeover
|
||||
R 7 cluster failover takeover
|
||||
}
|
||||
|
||||
test "Cluster should eventually be up again" {
|
||||
assert_cluster_state ok
|
||||
}
|
||||
|
||||
test "Cluster is writable" {
|
||||
cluster_write_test 4
|
||||
}
|
||||
|
||||
test "Instance #5, #6, #7 are now masters" {
|
||||
assert {[RI 5 role] eq {master}}
|
||||
assert {[RI 6 role] eq {master}}
|
||||
assert {[RI 7 role] eq {master}}
|
||||
}
|
||||
|
||||
test "Restarting the previously killed master nodes" {
|
||||
restart_instance redis 0
|
||||
restart_instance redis 1
|
||||
restart_instance redis 2
|
||||
}
|
||||
|
||||
test "Instance #0, #1, #2 gets converted into a slaves" {
|
||||
wait_for_condition 1000 50 {
|
||||
[RI 0 role] eq {slave} && [RI 1 role] eq {slave} && [RI 2 role] eq {slave}
|
||||
} else {
|
||||
fail "Old masters not converted into slaves"
|
||||
}
|
||||
}
|
||||
@@ -19,6 +19,7 @@ set ::verbose 0
|
||||
set ::valgrind 0
|
||||
set ::pause_on_error 0
|
||||
set ::simulate_error 0
|
||||
set ::failed 0
|
||||
set ::sentinel_instances {}
|
||||
set ::redis_instances {}
|
||||
set ::sentinel_base_port 20000
|
||||
@@ -231,6 +232,7 @@ proc test {descr code} {
|
||||
flush stdout
|
||||
|
||||
if {[catch {set retval [uplevel 1 $code]} error]} {
|
||||
incr ::failed
|
||||
if {[string match "assertion:*" $error]} {
|
||||
set msg [string range $error 10 end]
|
||||
puts [colorstr red $msg]
|
||||
@@ -246,6 +248,7 @@ proc test {descr code} {
|
||||
}
|
||||
}
|
||||
|
||||
# Execute all the units inside the 'tests' directory.
|
||||
proc run_tests {} {
|
||||
set tests [lsort [glob ../tests/*]]
|
||||
foreach test $tests {
|
||||
@@ -258,6 +261,17 @@ proc run_tests {} {
|
||||
}
|
||||
}
|
||||
|
||||
# Print a message and exists with 0 / 1 according to zero or more failures.
|
||||
proc end_tests {} {
|
||||
if {$::failed == 0} {
|
||||
puts "GOOD! No errors."
|
||||
exit 0
|
||||
} else {
|
||||
puts "WARNING $::failed tests faield."
|
||||
exit 1
|
||||
}
|
||||
}
|
||||
|
||||
# The "S" command is used to interact with the N-th Sentinel.
|
||||
# The general form is:
|
||||
#
|
||||
|
||||
@@ -1,10 +1,17 @@
|
||||
start_server {tags {"repl"}} {
|
||||
set A [srv 0 client]
|
||||
set A_host [srv 0 host]
|
||||
set A_port [srv 0 port]
|
||||
start_server {} {
|
||||
test {First server should have role slave after SLAVEOF} {
|
||||
r -1 slaveof [srv 0 host] [srv 0 port]
|
||||
set B [srv 0 client]
|
||||
set B_host [srv 0 host]
|
||||
set B_port [srv 0 port]
|
||||
|
||||
test {Set instance A as slave of B} {
|
||||
$A slaveof $B_host $B_port
|
||||
wait_for_condition 50 100 {
|
||||
[s -1 role] eq {slave} &&
|
||||
[string match {*master_link_status:up*} [r -1 info replication]]
|
||||
[lindex [$A role] 0] eq {slave} &&
|
||||
[string match {*master_link_status:up*} [$A info replication]]
|
||||
} else {
|
||||
fail "Can't turn the instance into a slave"
|
||||
}
|
||||
@@ -15,9 +22,9 @@ start_server {tags {"repl"}} {
|
||||
$rd brpoplpush a b 5
|
||||
r lpush a foo
|
||||
wait_for_condition 50 100 {
|
||||
[r debug digest] eq [r -1 debug digest]
|
||||
[$A debug digest] eq [$B debug digest]
|
||||
} else {
|
||||
fail "Master and slave have different digest: [r debug digest] VS [r -1 debug digest]"
|
||||
fail "Master and slave have different digest: [$A debug digest] VS [$B debug digest]"
|
||||
}
|
||||
}
|
||||
|
||||
@@ -28,7 +35,36 @@ start_server {tags {"repl"}} {
|
||||
r lpush c 3
|
||||
$rd brpoplpush c d 5
|
||||
after 1000
|
||||
assert_equal [r debug digest] [r -1 debug digest]
|
||||
assert_equal [$A debug digest] [$B debug digest]
|
||||
}
|
||||
|
||||
test {BLPOP followed by role change, issue #2473} {
|
||||
set rd [redis_deferring_client]
|
||||
$rd blpop foo 0 ; # Block while B is a master
|
||||
|
||||
# Turn B into master of A
|
||||
$A slaveof no one
|
||||
$B slaveof $A_host $A_port
|
||||
wait_for_condition 50 100 {
|
||||
[lindex [$B role] 0] eq {slave} &&
|
||||
[string match {*master_link_status:up*} [$B info replication]]
|
||||
} else {
|
||||
fail "Can't turn the instance into a slave"
|
||||
}
|
||||
|
||||
# Push elements into the "foo" list of the new slave.
|
||||
# If the client is still attached to the instance, we'll get
|
||||
# a desync between the two instances.
|
||||
$A rpush foo a b c
|
||||
after 100
|
||||
|
||||
wait_for_condition 50 100 {
|
||||
[$A debug digest] eq [$B debug digest] &&
|
||||
[$A lrange foo 0 -1] eq {a b c} &&
|
||||
[$B lrange foo 0 -1] eq {a b c}
|
||||
} else {
|
||||
fail "Master and slave have different digest: [$A debug digest] VS [$B debug digest]"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -113,7 +149,7 @@ foreach dl {no yes} {
|
||||
start_server {} {
|
||||
lappend slaves [srv 0 client]
|
||||
test "Connect multiple slaves at the same time (issue #141), diskless=$dl" {
|
||||
# Send SALVEOF commands to slaves
|
||||
# Send SLAVEOF commands to slaves
|
||||
[lindex $slaves 0] slaveof $master_host $master_port
|
||||
[lindex $slaves 1] slaveof $master_host $master_port
|
||||
[lindex $slaves 2] slaveof $master_host $master_port
|
||||
|
||||
@@ -13,6 +13,7 @@ proc main {} {
|
||||
spawn_instance redis $::redis_base_port $::instances_count
|
||||
run_tests
|
||||
cleanup
|
||||
end_tests
|
||||
}
|
||||
|
||||
if {[catch main e]} {
|
||||
|
||||
@@ -0,0 +1,34 @@
|
||||
# Test for the SENTINEL CKQUORUM command
|
||||
|
||||
source "../tests/includes/init-tests.tcl"
|
||||
set num_sentinels [llength $::sentinel_instances]
|
||||
|
||||
test "CKQUORUM reports OK and the right amount of Sentinels" {
|
||||
foreach_sentinel_id id {
|
||||
assert_match "*OK $num_sentinels usable*" [S $id SENTINEL CKQUORUM mymaster]
|
||||
}
|
||||
}
|
||||
|
||||
test "CKQUORUM detects quorum cannot be reached" {
|
||||
set orig_quorum [expr {$num_sentinels/2+1}]
|
||||
S 0 SENTINEL SET mymaster quorum [expr {$num_sentinels+1}]
|
||||
catch {[S 0 SENTINEL CKQUORUM mymaster]} err
|
||||
assert_match "*NOQUORUM*" $err
|
||||
S 0 SENTINEL SET mymaster quorum $orig_quorum
|
||||
}
|
||||
|
||||
test "CKQUORUM detects failover authorization cannot be reached" {
|
||||
set orig_quorum [expr {$num_sentinels/2+1}]
|
||||
S 0 SENTINEL SET mymaster quorum 1
|
||||
kill_instance sentinel 1
|
||||
kill_instance sentinel 2
|
||||
kill_instance sentinel 3
|
||||
after 5000
|
||||
catch {[S 0 SENTINEL CKQUORUM mymaster]} err
|
||||
assert_match "*NOQUORUM*" $err
|
||||
S 0 SENTINEL SET mymaster quorum $orig_quorum
|
||||
restart_instance sentinel 1
|
||||
restart_instance sentinel 2
|
||||
restart_instance sentinel 3
|
||||
}
|
||||
|
||||
@@ -54,10 +54,15 @@ proc kill_server config {
|
||||
|
||||
# kill server and wait for the process to be totally exited
|
||||
catch {exec kill $pid}
|
||||
if {$::valgrind} {
|
||||
set max_wait 60000
|
||||
} else {
|
||||
set max_wait 10000
|
||||
}
|
||||
while {[is_alive $config]} {
|
||||
incr wait 10
|
||||
|
||||
if {$wait >= 5000} {
|
||||
if {$wait >= $max_wait} {
|
||||
puts "Forcing process $pid to exit..."
|
||||
catch {exec kill -KILL $pid}
|
||||
} elseif {$wait % 1000 == 0} {
|
||||
|
||||
@@ -450,6 +450,7 @@ start_server {
|
||||
test "SMOVE non existing key" {
|
||||
setup_move
|
||||
assert_equal 0 [r smove myset1 myset2 foo]
|
||||
assert_equal 0 [r smove myset1 myset1 foo]
|
||||
assert_equal {1 a b} [lsort [r smembers myset1]]
|
||||
assert_equal {2 3 4} [lsort [r smembers myset2]]
|
||||
}
|
||||
|
||||
@@ -43,6 +43,84 @@ start_server {tags {"zset"}} {
|
||||
assert_error "*not*float*" {r zadd myzset nan abc}
|
||||
}
|
||||
|
||||
test "ZADD with options syntax error with incomplete pair" {
|
||||
r del ztmp
|
||||
catch {r zadd ztmp xx 10 x 20} err
|
||||
set err
|
||||
} {ERR*}
|
||||
|
||||
test "ZADD XX option without key - $encoding" {
|
||||
r del ztmp
|
||||
assert {[r zadd ztmp xx 10 x] == 0}
|
||||
assert {[r type ztmp] eq {none}}
|
||||
}
|
||||
|
||||
test "ZADD XX existing key - $encoding" {
|
||||
r del ztmp
|
||||
r zadd ztmp 10 x
|
||||
assert {[r zadd ztmp xx 20 y] == 0}
|
||||
assert {[r zcard ztmp] == 1}
|
||||
}
|
||||
|
||||
test "ZADD XX returns the number of elements actually added" {
|
||||
r del ztmp
|
||||
r zadd ztmp 10 x
|
||||
set retval [r zadd ztmp 10 x 20 y 30 z]
|
||||
assert {$retval == 2}
|
||||
}
|
||||
|
||||
test "ZADD XX updates existing elements score" {
|
||||
r del ztmp
|
||||
r zadd ztmp 10 x 20 y 30 z
|
||||
r zadd ztmp xx 5 foo 11 x 21 y 40 zap
|
||||
assert {[r zcard ztmp] == 3}
|
||||
assert {[r zscore ztmp x] == 11}
|
||||
assert {[r zscore ztmp y] == 21}
|
||||
}
|
||||
|
||||
test "ZADD XX and NX are not compatible" {
|
||||
r del ztmp
|
||||
catch {r zadd ztmp xx nx 10 x} err
|
||||
set err
|
||||
} {ERR*}
|
||||
|
||||
test "ZADD NX with non exisitng key" {
|
||||
r del ztmp
|
||||
r zadd ztmp nx 10 x 20 y 30 z
|
||||
assert {[r zcard ztmp] == 3}
|
||||
}
|
||||
|
||||
test "ZADD NX only add new elements without updating old ones" {
|
||||
r del ztmp
|
||||
r zadd ztmp 10 x 20 y 30 z
|
||||
assert {[r zadd ztmp nx 11 x 21 y 100 a 200 b] == 2}
|
||||
assert {[r zscore ztmp x] == 10}
|
||||
assert {[r zscore ztmp y] == 20}
|
||||
assert {[r zscore ztmp a] == 100}
|
||||
assert {[r zscore ztmp b] == 200}
|
||||
}
|
||||
|
||||
test "ZADD INCR works like ZINCRBY" {
|
||||
r del ztmp
|
||||
r zadd ztmp 10 x 20 y 30 z
|
||||
r zadd ztmp INCR 15 x
|
||||
assert {[r zscore ztmp x] == 25}
|
||||
}
|
||||
|
||||
test "ZADD INCR works with a single score-elemenet pair" {
|
||||
r del ztmp
|
||||
r zadd ztmp 10 x 20 y 30 z
|
||||
catch {r zadd ztmp INCR 15 x 10 y} err
|
||||
set err
|
||||
} {ERR*}
|
||||
|
||||
test "ZADD CH option changes return value to all changed elements" {
|
||||
r del ztmp
|
||||
r zadd ztmp 10 x 20 y 30 z
|
||||
assert {[r zadd ztmp 11 x 21 y 30 z] == 0}
|
||||
assert {[r zadd ztmp ch 12 x 22 y 30 z] == 2}
|
||||
}
|
||||
|
||||
test "ZINCRBY calls leading to NaN result in error" {
|
||||
r zincrby myzset +inf abc
|
||||
assert_error "*NaN*" {r zincrby myzset -inf abc}
|
||||
@@ -77,6 +155,8 @@ start_server {tags {"zset"}} {
|
||||
}
|
||||
|
||||
test "ZCARD basics - $encoding" {
|
||||
r del ztmp
|
||||
r zadd ztmp 10 a 20 b 30 c
|
||||
assert_equal 3 [r zcard ztmp]
|
||||
assert_equal 0 [r zcard zdoesntexist]
|
||||
}
|
||||
|
||||
@@ -43,7 +43,7 @@ then
|
||||
while [ $((PORT < ENDPORT)) != "0" ]; do
|
||||
PORT=$((PORT+1))
|
||||
echo "Stopping $PORT"
|
||||
redis-cli -p $PORT shutdown nosave
|
||||
../../src/redis-cli -p $PORT shutdown nosave
|
||||
done
|
||||
exit 0
|
||||
fi
|
||||
@@ -54,7 +54,7 @@ then
|
||||
while [ 1 ]; do
|
||||
clear
|
||||
date
|
||||
redis-cli -p $PORT cluster nodes | head -30
|
||||
../../src/redis-cli -p $PORT cluster nodes | head -30
|
||||
sleep 1
|
||||
done
|
||||
exit 0
|
||||
|
||||
Reference in New Issue
Block a user