X-Git-Url: https://git.saurik.com/redis.git/blobdiff_plain/33c08b39164eb669313f4871ce1b16a8e4cb2e5c..913e9d6bcacfa28ace3b7ee1c75f8ec4e146974b:/redis.c diff --git a/redis.c b/redis.c index 8b3f3be2..7158de60 100644 --- a/redis.c +++ b/redis.c @@ -27,16 +27,24 @@ * POSSIBILITY OF SUCH DAMAGE. */ -#define REDIS_VERSION "0.101" +#define REDIS_VERSION "1.050" #include "fmacros.h" +#include "config.h" #include #include #include #include #include +#define __USE_POSIX199309 #include + +#ifdef HAVE_BACKTRACE +#include +#include +#endif /* HAVE_BACKTRACE */ + #include #include #include @@ -49,8 +57,8 @@ #include #include #include -#include +#include "redis.h" #include "ae.h" /* Event driven programming library */ #include "sds.h" /* Dynamic safe strings */ #include "anet.h" /* Networking the easy way */ @@ -60,8 +68,6 @@ #include "lzf.h" /* LZF compression library */ #include "pqsort.h" /* Partial qsort for SORT+LIMIT */ -#include "config.h" - /* Error codes */ #define REDIS_OK 0 #define REDIS_ERR -1 @@ -77,10 +83,11 @@ #define REDIS_OBJFREELIST_MAX 1000000 /* Max number of objects to cache */ #define REDIS_MAX_SYNC_TIME 60 /* Slave can't take more to sync */ #define REDIS_EXPIRELOOKUPS_PER_CRON 100 /* try to expire 100 keys/second */ +#define REDIS_MAX_WRITE_PER_EVENT (1024*64) +#define REDIS_REQUEST_MAX_SIZE (1024*1024*256) /* max bytes in inline command */ /* Hash table parameters */ #define REDIS_HT_MINFILL 10 /* Minimal hash table fill 10% */ -#define REDIS_HT_MINSLOTS 16384 /* Never resize the HT under this */ /* Command flags */ #define REDIS_CMD_BULK 1 /* Bulk write command */ @@ -95,7 +102,12 @@ #define REDIS_STRING 0 #define REDIS_LIST 1 #define REDIS_SET 2 -#define REDIS_HASH 3 +#define REDIS_ZSET 3 +#define REDIS_HASH 4 + +/* Objects encoding */ +#define REDIS_ENCODING_RAW 0 /* Raw representation */ +#define REDIS_ENCODING_INT 1 /* Encoded as integer */ /* Object types only used for dumping to disk */ #define REDIS_EXPIRETIME 253 @@ -170,12 +182,17 @@ /* Anti-warning macro... */ #define REDIS_NOTUSED(V) ((void) V) +#define ZSKIPLIST_MAXLEVEL 32 /* Should be enough for 2^32 elements */ +#define ZSKIPLIST_P 0.25 /* Skiplist P = 1/4 */ + /*================================= Data types ============================== */ /* A redis object, that is a type able to hold a string / list / set */ typedef struct redisObject { void *ptr; - int type; + unsigned char type; + unsigned char encoding; + unsigned char notused[2]; int refcount; } robj; @@ -192,9 +209,10 @@ typedef struct redisClient { redisDb *db; int dictid; sds querybuf; - robj **argv; - int argc; + robj **argv, **mbargv; + int argc, mbargc; int bulklen; /* bulk read len. -1 if not in bulk read mode */ + int multibulk; /* multi bulk command format active */ list *reply; int sentlen; time_t lastinteraction; /* time of the last interaction, used for timeout */ @@ -240,6 +258,7 @@ struct redisServer { int daemonize; char *pidfile; int bgsaveinprogress; + pid_t bgsavechildpid; struct saveparam *saveparams; int saveparamslen; char *logfile; @@ -254,7 +273,7 @@ struct redisServer { redisClient *master; /* client that is master for this slave */ int replstate; unsigned int maxclients; - unsigned int maxmemory; + unsigned long maxmemory; /* Sort parameters - qsort_r() is only available under BSD so we * have to take this state global, in order to pass it to sortCompare() */ int sort_desc; @@ -270,6 +289,11 @@ struct redisCommand { int flags; }; +struct redisFunctionSym { + char *name; + unsigned long pointer; +}; + typedef struct _redisSortObject { robj *obj; union { @@ -283,6 +307,27 @@ typedef struct _redisSortOperation { robj *pattern; } redisSortOperation; +/* ZSETs use a specialized version of Skiplists */ + +typedef struct zskiplistNode { + struct zskiplistNode **forward; + double score; + robj *obj; +} zskiplistNode; + +typedef struct zskiplist { + struct zskiplistNode *header; + long length; + int level; +} zskiplist; + +typedef struct zset { + dict *dict; + zskiplist *zsl; +} zset; + +/* Our shared "common" objects */ + struct sharedObjectsStruct { robj *crlf, *ok, *err, *emptybulk, *czero, *cone, *pong, *space, *colon, *nullbulk, *nullmultibulk, @@ -309,14 +354,23 @@ static robj *createStringObject(char *ptr, size_t len); static void replicationFeedSlaves(list *slaves, struct redisCommand *cmd, int dictid, robj **argv, int argc); static int syncWithMaster(void); static robj *tryObjectSharing(robj *o); +static int tryObjectEncoding(robj *o); +static robj *getDecodedObject(const robj *o); static int removeExpire(redisDb *db, robj *key); static int expireIfNeeded(redisDb *db, robj *key); static int deleteIfVolatile(redisDb *db, robj *key); static int deleteKey(redisDb *db, robj *key); static time_t getExpire(redisDb *db, robj *key); static int setExpire(redisDb *db, robj *key, time_t when); -static void updateSalvesWaitingBgsave(int bgsaveerr); +static void updateSlavesWaitingBgsave(int bgsaveerr); static void freeMemoryIfNeeded(void); +static int processCommand(redisClient *c); +static void setupSigSegvAction(void); +static void rdbRemoveTempFile(pid_t childpid); +static size_t stringObjectLen(robj *o); +static void processInputBuffer(redisClient *c); +static zskiplist *zslCreate(void); +static void zslFree(zskiplist *zsl); static void authCommand(redisClient *c); static void pingCommand(redisClient *c); @@ -356,6 +410,8 @@ static void sremCommand(redisClient *c); static void smoveCommand(redisClient *c); static void sismemberCommand(redisClient *c); static void scardCommand(redisClient *c); +static void spopCommand(redisClient *c); +static void srandmemberCommand(redisClient *c); static void sinterCommand(redisClient *c); static void sinterstoreCommand(redisClient *c); static void sunionCommand(redisClient *c); @@ -371,10 +427,13 @@ static void infoCommand(redisClient *c); static void mgetCommand(redisClient *c); static void monitorCommand(redisClient *c); static void expireCommand(redisClient *c); -static void getSetCommand(redisClient *c); +static void getsetCommand(redisClient *c); static void ttlCommand(redisClient *c); static void slaveofCommand(redisClient *c); static void debugCommand(redisClient *c); +static void msetCommand(redisClient *c); +static void msetnxCommand(redisClient *c); +static void zaddCommand(redisClient *c); /*================================= Globals ================================= */ @@ -404,6 +463,8 @@ static struct redisCommand cmdTable[] = { {"smove",smoveCommand,4,REDIS_CMD_BULK}, {"sismember",sismemberCommand,3,REDIS_CMD_BULK}, {"scard",scardCommand,2,REDIS_CMD_INLINE}, + {"spop",spopCommand,2,REDIS_CMD_INLINE}, + {"srandmember",srandmemberCommand,2,REDIS_CMD_INLINE}, {"sinter",sinterCommand,-2,REDIS_CMD_INLINE|REDIS_CMD_DENYOOM}, {"sinterstore",sinterstoreCommand,-3,REDIS_CMD_INLINE|REDIS_CMD_DENYOOM}, {"sunion",sunionCommand,-2,REDIS_CMD_INLINE|REDIS_CMD_DENYOOM}, @@ -411,9 +472,12 @@ static struct redisCommand cmdTable[] = { {"sdiff",sdiffCommand,-2,REDIS_CMD_INLINE|REDIS_CMD_DENYOOM}, {"sdiffstore",sdiffstoreCommand,-3,REDIS_CMD_INLINE|REDIS_CMD_DENYOOM}, {"smembers",sinterCommand,2,REDIS_CMD_INLINE}, + {"zadd",zaddCommand,4,REDIS_CMD_BULK|REDIS_CMD_DENYOOM}, {"incrby",incrbyCommand,3,REDIS_CMD_INLINE|REDIS_CMD_DENYOOM}, {"decrby",decrbyCommand,3,REDIS_CMD_INLINE|REDIS_CMD_DENYOOM}, - {"getset",getSetCommand,3,REDIS_CMD_BULK|REDIS_CMD_DENYOOM}, + {"getset",getsetCommand,3,REDIS_CMD_BULK|REDIS_CMD_DENYOOM}, + {"mset",msetCommand,-3,REDIS_CMD_BULK|REDIS_CMD_DENYOOM}, + {"msetnx",msetnxCommand,-3,REDIS_CMD_BULK|REDIS_CMD_DENYOOM}, {"randomkey",randomkeyCommand,1,REDIS_CMD_INLINE}, {"select",selectCommand,2,REDIS_CMD_INLINE}, {"move",moveCommand,3,REDIS_CMD_INLINE}, @@ -441,7 +505,6 @@ static struct redisCommand cmdTable[] = { {"debug",debugCommand,-2,REDIS_CMD_INLINE}, {NULL,NULL,0,0} }; - /*============================ Utility functions ============================ */ /* Glob-style pattern matching. */ @@ -567,8 +630,7 @@ int stringmatchlen(const char *pattern, int patternLen, return 0; } -void redisLog(int level, const char *fmt, ...) -{ +static void redisLog(int level, const char *fmt, ...) { va_list ap; FILE *fp; @@ -599,6 +661,12 @@ void redisLog(int level, const char *fmt, ...) * keys and radis objects as values (objects can hold SDS strings, * lists, sets). */ +static void dictVanillaFree(void *privdata, void *val) +{ + DICT_NOTUSED(privdata); + zfree(val); +} + static int sdsDictKeyCompare(void *privdata, const void *key1, const void *key2) { @@ -618,32 +686,77 @@ static void dictRedisObjectDestructor(void *privdata, void *val) decrRefCount(val); } -static int dictSdsKeyCompare(void *privdata, const void *key1, +static int dictObjKeyCompare(void *privdata, const void *key1, const void *key2) { const robj *o1 = key1, *o2 = key2; return sdsDictKeyCompare(privdata,o1->ptr,o2->ptr); } -static unsigned int dictSdsHash(const void *key) { +static unsigned int dictObjHash(const void *key) { const robj *o = key; return dictGenHashFunction(o->ptr, sdslen((sds)o->ptr)); } +static int dictEncObjKeyCompare(void *privdata, const void *key1, + const void *key2) +{ + const robj *o1 = key1, *o2 = key2; + + if (o1->encoding == REDIS_ENCODING_RAW && + o2->encoding == REDIS_ENCODING_RAW) + return sdsDictKeyCompare(privdata,o1->ptr,o2->ptr); + else { + robj *dec1, *dec2; + int cmp; + + dec1 = o1->encoding != REDIS_ENCODING_RAW ? + getDecodedObject(o1) : (robj*)o1; + dec2 = o2->encoding != REDIS_ENCODING_RAW ? + getDecodedObject(o2) : (robj*)o2; + cmp = sdsDictKeyCompare(privdata,dec1->ptr,dec2->ptr); + if (dec1 != o1) decrRefCount(dec1); + if (dec2 != o2) decrRefCount(dec2); + return cmp; + } +} + +static unsigned int dictEncObjHash(const void *key) { + const robj *o = key; + + if (o->encoding == REDIS_ENCODING_RAW) + return dictGenHashFunction(o->ptr, sdslen((sds)o->ptr)); + else { + robj *dec = getDecodedObject(o); + unsigned int hash = dictGenHashFunction(dec->ptr, sdslen((sds)dec->ptr)); + decrRefCount(dec); + return hash; + } +} + static dictType setDictType = { - dictSdsHash, /* hash function */ + dictEncObjHash, /* hash function */ NULL, /* key dup */ NULL, /* val dup */ - dictSdsKeyCompare, /* key compare */ + dictEncObjKeyCompare, /* key compare */ dictRedisObjectDestructor, /* key destructor */ NULL /* val destructor */ }; +static dictType zsetDictType = { + dictEncObjHash, /* hash function */ + NULL, /* key dup */ + NULL, /* val dup */ + dictEncObjKeyCompare, /* key compare */ + dictRedisObjectDestructor, /* key destructor */ + dictVanillaFree /* val destructor */ +}; + static dictType hashDictType = { - dictSdsHash, /* hash function */ + dictObjHash, /* hash function */ NULL, /* key dup */ NULL, /* val dup */ - dictSdsKeyCompare, /* key compare */ + dictObjKeyCompare, /* key compare */ dictRedisObjectDestructor, /* key destructor */ dictRedisObjectDestructor /* val destructor */ }; @@ -663,7 +776,7 @@ static void oom(const char *msg) { } /* ====================== Redis server networking stuff ===================== */ -void closeTimedoutClients(void) { +static void closeTimedoutClients(void) { redisClient *c; listNode *ln; time_t now = time(NULL); @@ -680,26 +793,32 @@ void closeTimedoutClients(void) { } } +static int htNeedsResize(dict *dict) { + long long size, used; + + size = dictSlots(dict); + used = dictSize(dict); + return (size && used && size > DICT_HT_INITIAL_SIZE && + (used*100/size < REDIS_HT_MINFILL)); +} + /* If the percentage of used slots in the HT reaches REDIS_HT_MINFILL * we resize the hash table to save memory */ -void tryResizeHashTables(void) { +static void tryResizeHashTables(void) { int j; for (j = 0; j < server.dbnum; j++) { - long long size, used; - - size = dictSlots(server.db[j].dict); - used = dictSize(server.db[j].dict); - if (size && used && size > REDIS_HT_MINSLOTS && - (used*100/size < REDIS_HT_MINFILL)) { - redisLog(REDIS_NOTICE,"The hash table %d is too sparse, resize it...",j); + if (htNeedsResize(server.db[j].dict)) { + redisLog(REDIS_DEBUG,"The hash table %d is too sparse, resize it...",j); dictResize(server.db[j].dict); - redisLog(REDIS_NOTICE,"Hash table %d resized.",j); + redisLog(REDIS_DEBUG,"Hash table %d resized.",j); } + if (htNeedsResize(server.db[j].expires)) + dictResize(server.db[j].expires); } } -int serverCron(struct aeEventLoop *eventLoop, long long id, void *clientData) { +static int serverCron(struct aeEventLoop *eventLoop, long long id, void *clientData) { int j, loops = server.cronloops++; REDIS_NOTUSED(eventLoop); REDIS_NOTUSED(id); @@ -715,8 +834,8 @@ int serverCron(struct aeEventLoop *eventLoop, long long id, void *clientData) { size = dictSlots(server.db[j].dict); used = dictSize(server.db[j].dict); vkeys = dictSize(server.db[j].expires); - if (!(loops % 5) && used > 0) { - redisLog(REDIS_DEBUG,"DB %d: %d keys (%d volatile) in %d slots HT.",j,used,vkeys,size); + if (!(loops % 5) && (used || vkeys)) { + redisLog(REDIS_DEBUG,"DB %d: %lld keys (%lld volatile) in %lld slots HT.",j,used,vkeys,size); /* dictPrintStats(server.dict); */ } } @@ -731,7 +850,7 @@ int serverCron(struct aeEventLoop *eventLoop, long long id, void *clientData) { /* Show information about connected clients */ if (!(loops % 5)) { - redisLog(REDIS_DEBUG,"%d clients connected (%d slaves), %zu bytes in use", + redisLog(REDIS_DEBUG,"%d clients connected (%d slaves), %zu bytes in use, %d shared objects", listLength(server.clients)-listLength(server.slaves), listLength(server.slaves), server.usedmemory, @@ -745,20 +864,25 @@ int serverCron(struct aeEventLoop *eventLoop, long long id, void *clientData) { /* Check if a background saving in progress terminated */ if (server.bgsaveinprogress) { int statloc; - /* XXX: TODO handle the case of the saving child killed */ if (wait4(-1,&statloc,WNOHANG,NULL)) { int exitcode = WEXITSTATUS(statloc); - if (exitcode == 0) { + int bysignal = WIFSIGNALED(statloc); + + if (!bysignal && exitcode == 0) { redisLog(REDIS_NOTICE, "Background saving terminated with success"); server.dirty = 0; server.lastsave = time(NULL); + } else if (!bysignal && exitcode != 0) { + redisLog(REDIS_WARNING, "Background saving error"); } else { redisLog(REDIS_WARNING, - "Background saving error"); + "Background saving terminated by signal"); + rdbRemoveTempFile(server.bgsavechildpid); } server.bgsaveinprogress = 0; - updateSalvesWaitingBgsave(exitcode == 0 ? REDIS_OK : REDIS_ERR); + server.bgsavechildpid = -1; + updateSlavesWaitingBgsave(exitcode == 0 ? REDIS_OK : REDIS_ERR); } } else { /* If there is not a background saving in progress check if @@ -849,7 +973,6 @@ static void createSharedObjects(void) { static void appendServerSaveParams(time_t seconds, int changes) { server.saveparams = zrealloc(server.saveparams,sizeof(struct saveparam)*(server.saveparamslen+1)); - if (server.saveparams == NULL) oom("appendServerSaveParams"); server.saveparams[server.saveparamslen].seconds = seconds; server.saveparams[server.saveparamslen].changes = changes; server.saveparamslen++; @@ -875,6 +998,7 @@ static void initServerConfig() { server.dbfilename = "dump.rdb"; server.requirepass = NULL; server.shareobjects = 0; + server.sharingpoolsize = 1024; server.maxclients = 0; server.maxmemory = 0; ResetServerSaveParams(); @@ -895,6 +1019,7 @@ static void initServer() { signal(SIGHUP, SIG_IGN); signal(SIGPIPE, SIG_IGN); + setupSigSegvAction(); server.clients = listCreate(); server.slaves = listCreate(); @@ -904,9 +1029,6 @@ static void initServer() { server.el = aeCreateEventLoop(); server.db = zmalloc(sizeof(redisDb)*server.dbnum); server.sharingpool = dictCreate(&setDictType,NULL); - server.sharingpoolsize = 1024; - if (!server.db || !server.clients || !server.slaves || !server.monitors || !server.el || !server.objfreelist) - oom("server initialization"); /* Fatal OOM */ server.fd = anetTcpServer(server.neterr, server.port, server.bindaddr); if (server.fd == -1) { redisLog(REDIS_WARNING, "Opening TCP port: %s", server.neterr); @@ -919,6 +1041,7 @@ static void initServer() { } server.cronloops = 0; server.bgsaveinprogress = 0; + server.bgsavechildpid = -1; server.lastsave = time(NULL); server.dirty = 0; server.usedmemory = 0; @@ -950,15 +1073,20 @@ static int yesnotoi(char *s) { /* I agree, this is a very rudimental way to load a configuration... will improve later if the config gets more complex */ static void loadServerConfig(char *filename) { - FILE *fp = fopen(filename,"r"); + FILE *fp; char buf[REDIS_CONFIGLINE_MAX+1], *err = NULL; int linenum = 0; sds line = NULL; - - if (!fp) { - redisLog(REDIS_WARNING,"Fatal error, can't open config file"); - exit(1); + + if (filename[0] == '-' && filename[1] == '\0') + fp = stdin; + else { + if ((fp = fopen(filename,"r")) == NULL) { + redisLog(REDIS_WARNING,"Fatal error, can't open config file"); + exit(1); + } } + while(fgets(buf,REDIS_CONFIGLINE_MAX+1,fp) != NULL) { sds *argv; int argc, j; @@ -1012,7 +1140,7 @@ static void loadServerConfig(char *filename) { goto loaderr; } } else if (!strcasecmp(argv[0],"logfile") && argc == 2) { - FILE *fp; + FILE *logfp; server.logfile = zstrdup(argv[1]); if (!strcasecmp(server.logfile,"stdout")) { @@ -1022,13 +1150,13 @@ static void loadServerConfig(char *filename) { if (server.logfile) { /* Test if we are able to open the file. The server will not * be able to abort just for this problem later... */ - fp = fopen(server.logfile,"a"); - if (fp == NULL) { + logfp = fopen(server.logfile,"a"); + if (logfp == NULL) { err = sdscatprintf(sdsempty(), "Can't open the log file: %s", strerror(errno)); goto loaderr; } - fclose(fp); + fclose(logfp); } } else if (!strcasecmp(argv[0],"databases") && argc == 2) { server.dbnum = atoi(argv[1]); @@ -1038,7 +1166,7 @@ static void loadServerConfig(char *filename) { } else if (!strcasecmp(argv[0],"maxclients") && argc == 2) { server.maxclients = atoi(argv[1]); } else if (!strcasecmp(argv[0],"maxmemory") && argc == 2) { - server.maxmemory = atoi(argv[1]); + server.maxmemory = strtoll(argv[1], NULL, 10); } else if (!strcasecmp(argv[0],"slaveof") && argc == 3) { server.masterhost = sdsnew(argv[1]); server.masterport = atoi(argv[2]); @@ -1074,7 +1202,7 @@ static void loadServerConfig(char *filename) { zfree(argv); sdsfree(line); } - fclose(fp); + if (fp != stdin) fclose(fp); return; loaderr: @@ -1090,7 +1218,10 @@ static void freeClientArgv(redisClient *c) { for (j = 0; j < c->argc; j++) decrRefCount(c->argv[j]); + for (j = 0; j < c->mbargc; j++) + decrRefCount(c->mbargv[j]); c->argc = 0; + c->mbargc = 0; } static void freeClient(redisClient *c) { @@ -1118,6 +1249,7 @@ static void freeClient(redisClient *c) { server.replstate = REDIS_REPL_CONNECT; } zfree(c->argv); + zfree(c->mbargv); zfree(c); } @@ -1147,7 +1279,7 @@ static void glueReplyBuffersIfNeeded(redisClient *c) { } /* Now the output buffer is empty, add the new single element */ o = createObject(REDIS_STRING,sdsnewlen(buf,totlen)); - if (!listAddNodeTail(c->reply,o)) oom("listAddNodeTail"); + listAddNodeTail(c->reply,o); } } @@ -1170,6 +1302,7 @@ static void sendReplyToClient(aeEventLoop *el, int fd, void *privdata, int mask) } if (c->flags & REDIS_MASTER) { + /* Don't reply to a master */ nwritten = objlen - c->sentlen; } else { nwritten = write(fd, ((char*)o->ptr)+c->sentlen, objlen - c->sentlen); @@ -1182,6 +1315,12 @@ static void sendReplyToClient(aeEventLoop *el, int fd, void *privdata, int mask) listDelNode(c->reply,listFirst(c->reply)); c->sentlen = 0; } + /* Note that we avoid to send more thank REDIS_MAX_WRITE_PER_EVENT + * bytes, in a single threaded server it's a good idea to server + * other clients as well, even if a very large request comes from + * super fast link that is always able to accept data (in real world + * terms think to 'KEYS *' against the loopback interfae) */ + if (totwritten > REDIS_MAX_WRITE_PER_EVENT) break; } if (nwritten == -1) { if (errno == EAGAIN) { @@ -1213,6 +1352,7 @@ static struct redisCommand *lookupCommand(char *name) { static void resetClient(redisClient *c) { freeClientArgv(c); c->bulklen = -1; + c->multibulk = 0; } /* If this function gets called we already read a whole @@ -1230,6 +1370,74 @@ static int processCommand(redisClient *c) { /* Free some memory if needed (maxmemory setting) */ if (server.maxmemory) freeMemoryIfNeeded(); + /* Handle the multi bulk command type. This is an alternative protocol + * supported by Redis in order to receive commands that are composed of + * multiple binary-safe "bulk" arguments. The latency of processing is + * a bit higher but this allows things like multi-sets, so if this + * protocol is used only for MSET and similar commands this is a big win. */ + if (c->multibulk == 0 && c->argc == 1 && ((char*)(c->argv[0]->ptr))[0] == '*') { + c->multibulk = atoi(((char*)c->argv[0]->ptr)+1); + if (c->multibulk <= 0) { + resetClient(c); + return 1; + } else { + decrRefCount(c->argv[c->argc-1]); + c->argc--; + return 1; + } + } else if (c->multibulk) { + if (c->bulklen == -1) { + if (((char*)c->argv[0]->ptr)[0] != '$') { + addReplySds(c,sdsnew("-ERR multi bulk protocol error\r\n")); + resetClient(c); + return 1; + } else { + int bulklen = atoi(((char*)c->argv[0]->ptr)+1); + decrRefCount(c->argv[0]); + if (bulklen < 0 || bulklen > 1024*1024*1024) { + c->argc--; + addReplySds(c,sdsnew("-ERR invalid bulk write count\r\n")); + resetClient(c); + return 1; + } + c->argc--; + c->bulklen = bulklen+2; /* add two bytes for CR+LF */ + return 1; + } + } else { + c->mbargv = zrealloc(c->mbargv,(sizeof(robj*))*(c->mbargc+1)); + c->mbargv[c->mbargc] = c->argv[0]; + c->mbargc++; + c->argc--; + c->multibulk--; + if (c->multibulk == 0) { + robj **auxargv; + int auxargc; + + /* Here we need to swap the multi-bulk argc/argv with the + * normal argc/argv of the client structure. */ + auxargv = c->argv; + c->argv = c->mbargv; + c->mbargv = auxargv; + + auxargc = c->argc; + c->argc = c->mbargc; + c->mbargc = auxargc; + + /* We need to set bulklen to something different than -1 + * in order for the code below to process the command without + * to try to read the last argument of a bulk command as + * a special argument. */ + c->bulklen = 0; + /* continue below and process the command */ + } else { + c->bulklen = -1; + return 1; + } + } + } + /* -- end of multi bulk commands processing -- */ + /* The QUIT command is handled as a special case. Normal command * procs are unable to close the client connection safely */ if (!strcasecmp(c->argv[0]->ptr,"quit")) { @@ -1263,7 +1471,10 @@ static int processCommand(redisClient *c) { c->argc--; c->bulklen = bulklen+2; /* add two bytes for CR+LF */ /* It is possible that the bulk read is already in the - * buffer. Check this condition and handle it accordingly */ + * buffer. Check this condition and handle it accordingly. + * This is just a fast path, alternative to call processInputBuffer(). + * It's a good idea since the code is small and this condition + * happens most of the times. */ if ((signed)sdslen(c->querybuf) >= c->bulklen) { c->argv[c->argc] = createStringObject(c->querybuf,c->bulklen-2); c->argc++; @@ -1278,6 +1489,10 @@ static int processCommand(redisClient *c) { for(j = 1; j < c->argc; j++) c->argv[j] = tryObjectSharing(c->argv[j]); } + /* Let's try to encode the bulk object to save space. */ + if (cmd->flags & REDIS_CMD_BULK) + tryObjectEncoding(c->argv[c->argc-1]); + /* Check if the user is authenticated */ if (server.requirepass && !c->authenticated && cmd->proc != authCommand) { addReplySds(c,sdsnew("-ERR operation not permitted\r\n")); @@ -1314,7 +1529,6 @@ static void replicationFeedSlaves(list *slaves, struct redisCommand *cmd, int di outv = static_outv; } else { outv = zmalloc(sizeof(robj*)*(argc*2+1)); - if (!outv) oom("replicationFeedSlaves"); } for (j = 0; j < argc; j++) { @@ -1323,7 +1537,8 @@ static void replicationFeedSlaves(list *slaves, struct redisCommand *cmd, int di robj *lenobj; lenobj = createObject(REDIS_STRING, - sdscatprintf(sdsempty(),"%d\r\n",sdslen(argv[j]->ptr))); + sdscatprintf(sdsempty(),"%d\r\n", + stringObjectLen(argv[j]))); lenobj->refcount = 0; outv[outc++] = lenobj; } @@ -1372,39 +1587,13 @@ static void replicationFeedSlaves(list *slaves, struct redisCommand *cmd, int di if (outv != static_outv) zfree(outv); } -static void readQueryFromClient(aeEventLoop *el, int fd, void *privdata, int mask) { - redisClient *c = (redisClient*) privdata; - char buf[REDIS_IOBUF_LEN]; - int nread; - REDIS_NOTUSED(el); - REDIS_NOTUSED(mask); - - nread = read(fd, buf, REDIS_IOBUF_LEN); - if (nread == -1) { - if (errno == EAGAIN) { - nread = 0; - } else { - redisLog(REDIS_DEBUG, "Reading from client: %s",strerror(errno)); - freeClient(c); - return; - } - } else if (nread == 0) { - redisLog(REDIS_DEBUG, "Client closed connection"); - freeClient(c); - return; - } - if (nread) { - c->querybuf = sdscatlen(c->querybuf, buf, nread); - c->lastinteraction = time(NULL); - } else { - return; - } - +static void processInputBuffer(redisClient *c) { again: if (c->bulklen == -1) { /* Read the first line of the query */ char *p = strchr(c->querybuf,'\n'); size_t querylen; + if (p) { sds query, *argv; int argc, j; @@ -1427,12 +1616,10 @@ again: return; } argv = sdssplitlen(query,sdslen(query)," ",1,&argc); - if (argv == NULL) oom("sdssplitlen"); sdsfree(query); if (c->argv) zfree(c->argv); c->argv = zmalloc(sizeof(robj*)*argc); - if (c->argv == NULL) oom("allocating arguments list for client"); for (j = 0; j < argc; j++) { if (sdslen(argv[j])) { @@ -1446,9 +1633,9 @@ again: /* Execute the command. If the client is still valid * after processCommand() return and there is something * on the query buffer try to process the next command. */ - if (processCommand(c) && sdslen(c->querybuf)) goto again; + if (c->argc && processCommand(c) && sdslen(c->querybuf)) goto again; return; - } else if (sdslen(c->querybuf) >= 1024*32) { + } else if (sdslen(c->querybuf) >= REDIS_REQUEST_MAX_SIZE) { redisLog(REDIS_DEBUG, "Client protocol error"); freeClient(c); return; @@ -1465,12 +1652,45 @@ again: c->argv[c->argc] = createStringObject(c->querybuf,c->bulklen-2); c->argc++; c->querybuf = sdsrange(c->querybuf,c->bulklen,-1); - processCommand(c); + /* Process the command. If the client is still valid after + * the processing and there is more data in the buffer + * try to parse it. */ + if (processCommand(c) && sdslen(c->querybuf)) goto again; return; } } } +static void readQueryFromClient(aeEventLoop *el, int fd, void *privdata, int mask) { + redisClient *c = (redisClient*) privdata; + char buf[REDIS_IOBUF_LEN]; + int nread; + REDIS_NOTUSED(el); + REDIS_NOTUSED(mask); + + nread = read(fd, buf, REDIS_IOBUF_LEN); + if (nread == -1) { + if (errno == EAGAIN) { + nread = 0; + } else { + redisLog(REDIS_DEBUG, "Reading from client: %s",strerror(errno)); + freeClient(c); + return; + } + } else if (nread == 0) { + redisLog(REDIS_DEBUG, "Client closed connection"); + freeClient(c); + return; + } + if (nread) { + c->querybuf = sdscatlen(c->querybuf, buf, nread); + c->lastinteraction = time(NULL); + } else { + return; + } + processInputBuffer(c); +} + static int selectDb(redisClient *c, int id) { if (id < 0 || id >= server.dbnum) return REDIS_ERR; @@ -1495,12 +1715,15 @@ static redisClient *createClient(int fd) { c->argc = 0; c->argv = NULL; c->bulklen = -1; + c->multibulk = 0; + c->mbargc = 0; + c->mbargv = NULL; c->sentlen = 0; c->flags = 0; c->lastinteraction = time(NULL); c->authenticated = 0; c->replstate = REDIS_REPL_NONE; - if ((c->reply = listCreate()) == NULL) oom("listCreate"); + c->reply = listCreate(); listSetFreeMethod(c->reply,decrRefCount); listSetDupMethod(c->reply,dupClientReplyValue); if (aeCreateFileEvent(server.el, c->fd, AE_READABLE, @@ -1508,7 +1731,7 @@ static redisClient *createClient(int fd) { freeClient(c); return NULL; } - if (!listAddNodeTail(server.clients,c)) oom("listAddNodeTail"); + listAddNodeTail(server.clients,c); return c; } @@ -1518,8 +1741,12 @@ static void addReply(redisClient *c, robj *obj) { c->replstate == REDIS_REPL_ONLINE) && aeCreateFileEvent(server.el, c->fd, AE_WRITABLE, sendReplyToClient, c, NULL) == AE_ERR) return; - if (!listAddNodeTail(c->reply,obj)) oom("listAddNodeTail"); - incrRefCount(obj); + if (obj->encoding != REDIS_ENCODING_RAW) { + obj = getDecodedObject(obj); + } else { + incrRefCount(obj); + } + listAddNodeTail(c->reply,obj); } static void addReplySds(redisClient *c, sds s) { @@ -1528,6 +1755,26 @@ static void addReplySds(redisClient *c, sds s) { decrRefCount(o); } +static void addReplyBulkLen(redisClient *c, robj *obj) { + size_t len; + + if (obj->encoding == REDIS_ENCODING_RAW) { + len = sdslen(obj->ptr); + } else { + long n = (long)obj->ptr; + + len = 1; + if (n < 0) { + len++; + n = -n; + } + while((n = n/10) != 0) { + len++; + } + } + addReplySds(c,sdscatprintf(sdsempty(),"$%d\r\n",len)); +} + static void acceptHandler(aeEventLoop *el, int fd, void *privdata, int mask) { int cport, cfd; char cip[128]; @@ -1574,8 +1821,8 @@ static robj *createObject(int type, void *ptr) { } else { o = zmalloc(sizeof(*o)); } - if (!o) oom("createObject"); o->type = type; + o->encoding = REDIS_ENCODING_RAW; o->ptr = ptr; o->refcount = 1; return o; @@ -1588,19 +1835,27 @@ static robj *createStringObject(char *ptr, size_t len) { static robj *createListObject(void) { list *l = listCreate(); - if (!l) oom("listCreate"); listSetFreeMethod(l,decrRefCount); return createObject(REDIS_LIST,l); } static robj *createSetObject(void) { dict *d = dictCreate(&setDictType,NULL); - if (!d) oom("dictCreate"); return createObject(REDIS_SET,d); } +static robj *createZsetObject(void) { + zset *zs = zmalloc(sizeof(*zs)); + + zs->dict = dictCreate(&zsetDictType,NULL); + zs->zsl = zslCreate(); + return createObject(REDIS_ZSET,zs); +} + static void freeStringObject(robj *o) { - sdsfree(o->ptr); + if (o->encoding == REDIS_ENCODING_RAW) { + sdsfree(o->ptr); + } } static void freeListObject(robj *o) { @@ -1611,6 +1866,14 @@ static void freeSetObject(robj *o) { dictRelease((dict*) o->ptr); } +static void freeZsetObject(robj *o) { + zset *zs = o->ptr; + + dictRelease(zs->dict); + zslFree(zs->zsl); + zfree(zs); +} + static void freeHashObject(robj *o) { dictRelease((dict*) o->ptr); } @@ -1635,6 +1898,7 @@ static void decrRefCount(void *obj) { case REDIS_STRING: freeStringObject(o); break; case REDIS_LIST: freeListObject(o); break; case REDIS_SET: freeSetObject(o); break; + case REDIS_ZSET: freeZsetObject(o); break; case REDIS_HASH: freeHashObject(o); break; default: assert(0 != 0); break; } @@ -1644,6 +1908,36 @@ static void decrRefCount(void *obj) { } } +static robj *lookupKey(redisDb *db, robj *key) { + dictEntry *de = dictFind(db->dict,key); + return de ? dictGetEntryVal(de) : NULL; +} + +static robj *lookupKeyRead(redisDb *db, robj *key) { + expireIfNeeded(db,key); + return lookupKey(db,key); +} + +static robj *lookupKeyWrite(redisDb *db, robj *key) { + deleteIfVolatile(db,key); + return lookupKey(db,key); +} + +static int deleteKey(redisDb *db, robj *key) { + int retval; + + /* We need to protect key from destruction: after the first dictDelete() + * it may happen that 'key' is no longer valid if we don't increment + * it's count. This may happen when we get the object reference directly + * from the hash table with dictRandomKey() or dict iterators */ + incrRefCount(key); + if (dictSize(db->expires)) dictDelete(db->expires,key); + retval = dictDelete(db->dict,key); + decrRefCount(key); + + return retval == DICT_OK; +} + /* Try to share an object against the shared objects pool */ static robj *tryObjectSharing(robj *o) { struct dictEntry *de; @@ -1689,34 +1983,98 @@ static robj *tryObjectSharing(robj *o) { } } -static robj *lookupKey(redisDb *db, robj *key) { - dictEntry *de = dictFind(db->dict,key); - return de ? dictGetEntryVal(de) : NULL; +/* Check if the nul-terminated string 's' can be represented by a long + * (that is, is a number that fits into long without any other space or + * character before or after the digits). + * + * If so, the function returns REDIS_OK and *longval is set to the value + * of the number. Otherwise REDIS_ERR is returned */ +static int isStringRepresentableAsLong(sds s, long *longval) { + char buf[32], *endptr; + long value; + int slen; + + value = strtol(s, &endptr, 10); + if (endptr[0] != '\0') return REDIS_ERR; + slen = snprintf(buf,32,"%ld",value); + + /* If the number converted back into a string is not identical + * then it's not possible to encode the string as integer */ + if (sdslen(s) != (unsigned)slen || memcmp(buf,s,slen)) return REDIS_ERR; + if (longval) *longval = value; + return REDIS_OK; } -static robj *lookupKeyRead(redisDb *db, robj *key) { - expireIfNeeded(db,key); - return lookupKey(db,key); +/* Try to encode a string object in order to save space */ +static int tryObjectEncoding(robj *o) { + long value; + sds s = o->ptr; + + if (o->encoding != REDIS_ENCODING_RAW) + return REDIS_ERR; /* Already encoded */ + + /* It's not save to encode shared objects: shared objects can be shared + * everywhere in the "object space" of Redis. Encoded objects can only + * appear as "values" (and not, for instance, as keys) */ + if (o->refcount > 1) return REDIS_ERR; + + /* Currently we try to encode only strings */ + assert(o->type == REDIS_STRING); + + /* Check if we can represent this string as a long integer */ + if (isStringRepresentableAsLong(s,&value) == REDIS_ERR) return REDIS_ERR; + + /* Ok, this object can be encoded */ + o->encoding = REDIS_ENCODING_INT; + sdsfree(o->ptr); + o->ptr = (void*) value; + return REDIS_OK; } -static robj *lookupKeyWrite(redisDb *db, robj *key) { - deleteIfVolatile(db,key); - return lookupKey(db,key); +/* Get a decoded version of an encoded object (returned as a new object) */ +static robj *getDecodedObject(const robj *o) { + robj *dec; + + assert(o->encoding != REDIS_ENCODING_RAW); + if (o->type == REDIS_STRING && o->encoding == REDIS_ENCODING_INT) { + char buf[32]; + + snprintf(buf,32,"%ld",(long)o->ptr); + dec = createStringObject(buf,strlen(buf)); + return dec; + } else { + assert(1 != 1); + } } -static int deleteKey(redisDb *db, robj *key) { - int retval; +static int compareStringObjects(robj *a, robj *b) { + assert(a->type == REDIS_STRING && b->type == REDIS_STRING); - /* We need to protect key from destruction: after the first dictDelete() - * it may happen that 'key' is no longer valid if we don't increment - * it's count. This may happen when we get the object reference directly - * from the hash table with dictRandomKey() or dict iterators */ - incrRefCount(key); - if (dictSize(db->expires)) dictDelete(db->expires,key); - retval = dictDelete(db->dict,key); - decrRefCount(key); + if (a->encoding == REDIS_ENCODING_INT && b->encoding == REDIS_ENCODING_INT){ + return (long)a->ptr - (long)b->ptr; + } else { + int retval; - return retval == DICT_OK; + incrRefCount(a); + incrRefCount(b); + if (a->encoding != REDIS_ENCODING_RAW) a = getDecodedObject(a); + if (b->encoding != REDIS_ENCODING_RAW) b = getDecodedObject(a); + retval = sdscmp(a->ptr,b->ptr); + decrRefCount(a); + decrRefCount(b); + return retval; + } +} + +static size_t stringObjectLen(robj *o) { + assert(o->type == REDIS_STRING); + if (o->encoding == REDIS_ENCODING_RAW) { + return sdslen(o->ptr); + } else { + char buf[32]; + + return snprintf(buf,32,"%ld",(long)o->ptr); + } } /*============================ DB saving/loading ============================ */ @@ -1758,7 +2116,7 @@ static int rdbSaveLen(FILE *fp, uint32_t len) { /* String objects in the form "2391" "-100" without any space and with a * range of values that can fit in an 8, 16 or 32 bit signed value can be * encoded as integers to save space */ -int rdbTryIntegerEncoding(sds s, unsigned char *enc) { +static int rdbTryIntegerEncoding(sds s, unsigned char *enc) { long long value; char *endptr, buf[32]; @@ -1823,10 +2181,12 @@ writeerr: /* Save a string objet as [len][data] on disk. If the object is a string * representation of an integer value we try to safe it in a special form */ -static int rdbSaveStringObject(FILE *fp, robj *obj) { - size_t len = sdslen(obj->ptr); +static int rdbSaveStringObjectRaw(FILE *fp, robj *obj) { + size_t len; int enclen; + len = sdslen(obj->ptr); + /* Try integer encoding */ if (len <= 11) { unsigned char buf[5]; @@ -1838,7 +2198,7 @@ static int rdbSaveStringObject(FILE *fp, robj *obj) { /* Try LZF compression - under 20 bytes it's unable to compress even * aaaaaaaaaaaaaaaaaa so skip it */ - if (1 && len > 20) { + if (len > 20) { int retval; retval = rdbSaveLzfStringObject(fp,obj); @@ -1853,6 +2213,21 @@ static int rdbSaveStringObject(FILE *fp, robj *obj) { return 0; } +/* Like rdbSaveStringObjectRaw() but handle encoded objects */ +static int rdbSaveStringObject(FILE *fp, robj *obj) { + int retval; + robj *dec; + + if (obj->encoding != REDIS_ENCODING_RAW) { + dec = getDecodedObject(obj); + retval = rdbSaveStringObjectRaw(fp,dec); + decrRefCount(dec); + return retval; + } else { + return rdbSaveStringObjectRaw(fp,obj); + } +} + /* Save the DB on disk. Return REDIS_ERR on error, REDIS_OK on success */ static int rdbSave(char *filename) { dictIterator *di = NULL; @@ -1862,7 +2237,7 @@ static int rdbSave(char *filename) { int j; time_t now = time(NULL); - snprintf(tmpfile,256,"temp-%d.%ld.rdb",(int)time(NULL),(long int)random()); + snprintf(tmpfile,256,"temp-%d.rdb", (int) getpid()); fp = fopen(tmpfile,"w"); if (!fp) { redisLog(REDIS_WARNING, "Failed saving the DB: %s", strerror(errno)); @@ -1920,7 +2295,6 @@ static int rdbSave(char *filename) { dictIterator *di = dictGetIterator(set); dictEntry *de; - if (!set) oom("dictGetIteraotr"); if (rdbSaveLen(fp,dictSize(set)) == -1) goto werr; while((de = dictNext(di)) != NULL) { robj *eleobj = dictGetEntryKey(de); @@ -1983,11 +2357,19 @@ static int rdbSaveBackground(char *filename) { } redisLog(REDIS_NOTICE,"Background saving started by pid %d",childpid); server.bgsaveinprogress = 1; + server.bgsavechildpid = childpid; return REDIS_OK; } return REDIS_OK; /* unreached */ } +static void rdbRemoveTempFile(pid_t childpid) { + char tmpfile[256]; + + snprintf(tmpfile,256,"temp-%d.rdb", (int) childpid); + unlink(tmpfile); +} + static int rdbLoadType(FILE *fp) { unsigned char type; if (fread(&type,1,1,fp) == 0) return -1; @@ -2162,6 +2544,7 @@ static int rdbLoad(char *filename) { if (type == REDIS_STRING) { /* Read string value */ if ((o = rdbLoadStringObject(fp,rdbver)) == NULL) goto eoferr; + tryObjectEncoding(o); } else if (type == REDIS_LIST || type == REDIS_SET) { /* Read list/set value */ uint32_t listlen; @@ -2174,12 +2557,11 @@ static int rdbLoad(char *filename) { robj *ele; if ((ele = rdbLoadStringObject(fp,rdbver)) == NULL) goto eoferr; + tryObjectEncoding(ele); if (type == REDIS_LIST) { - if (!listAddNodeTail((list*)o->ptr,ele)) - oom("listAddNodeTail"); + listAddNodeTail((list*)o->ptr,ele); } else { - if (dictAdd((dict*)o->ptr,ele,NULL) == DICT_ERR) - oom("dictAdd"); + dictAdd((dict*)o->ptr,ele,NULL); } } } else { @@ -2227,8 +2609,7 @@ static void pingCommand(redisClient *c) { } static void echoCommand(redisClient *c) { - addReplySds(c,sdscatprintf(sdsempty(),"$%d\r\n", - (int)sdslen(c->argv[1]->ptr))); + addReplyBulkLen(c,c->argv[1]); addReply(c,c->argv[1]); addReply(c,shared.crlf); } @@ -2273,14 +2654,14 @@ static void getCommand(redisClient *c) { if (o->type != REDIS_STRING) { addReply(c,shared.wrongtypeerr); } else { - addReplySds(c,sdscatprintf(sdsempty(),"$%d\r\n",(int)sdslen(o->ptr))); + addReplyBulkLen(c,o); addReply(c,o); addReply(c,shared.crlf); } } } -static void getSetCommand(redisClient *c) { +static void getsetCommand(redisClient *c) { getCommand(c); if (dictAdd(c->db->dict,c->argv[1],c->argv[2]) == DICT_ERR) { dictReplace(c->db->dict,c->argv[1],c->argv[2]); @@ -2304,7 +2685,7 @@ static void mgetCommand(redisClient *c) { if (o->type != REDIS_STRING) { addReply(c,shared.nullbulk); } else { - addReplySds(c,sdscatprintf(sdsempty(),"$%d\r\n",(int)sdslen(o->ptr))); + addReplyBulkLen(c,o); addReply(c,o); addReply(c,shared.crlf); } @@ -2326,12 +2707,18 @@ static void incrDecrCommand(redisClient *c, long long incr) { } else { char *eptr; - value = strtoll(o->ptr, &eptr, 10); + if (o->encoding == REDIS_ENCODING_RAW) + value = strtoll(o->ptr, &eptr, 10); + else if (o->encoding == REDIS_ENCODING_INT) + value = (long)o->ptr; + else + assert(1 != 1); } } value += incr; o = createObject(REDIS_STRING,sdscatprintf(sdsempty(),"%lld",value)); + tryObjectEncoding(o); retval = dictAdd(c->db->dict,c->argv[1],o); if (retval == DICT_ERR) { dictReplace(c->db->dict,c->argv[1],o); @@ -2427,7 +2814,6 @@ static void keysCommand(redisClient *c) { robj *lenobj = createObject(REDIS_STRING,NULL); di = dictGetIterator(c->db->dict); - if (!di) oom("dictGetIterator"); addReply(c,lenobj); decrRefCount(lenobj); while((de = dictNext(di)) != NULL) { @@ -2505,15 +2891,26 @@ static void bgsaveCommand(redisClient *c) { static void shutdownCommand(redisClient *c) { redisLog(REDIS_WARNING,"User requested shutdown, saving DB..."); - /* XXX: TODO kill the child if there is a bgsave in progress */ + /* Kill the saving child if there is a background saving in progress. + We want to avoid race conditions, for instance our saving child may + overwrite the synchronous saving did by SHUTDOWN. */ + if (server.bgsaveinprogress) { + redisLog(REDIS_WARNING,"There is a live saving child. Killing it!"); + kill(server.bgsavechildpid,SIGKILL); + rdbRemoveTempFile(server.bgsavechildpid); + } + /* SYNC SAVE */ if (rdbSave(server.dbfilename) == REDIS_OK) { - if (server.daemonize) { + if (server.daemonize) unlink(server.pidfile); - } redisLog(REDIS_WARNING,"%zu bytes used at exit",zmalloc_used_memory()); redisLog(REDIS_WARNING,"Server exit now, bye bye..."); exit(1); } else { + /* Ooops.. error saving! The best we can do is to continue operating. + * Note that if there was a background saving process, in the next + * cron() Redis will be notified that the background saving aborted, + * handling special stuff like slaves pending for synchronization... */ redisLog(REDIS_WARNING,"Error trying to save the DB, can't exit"); addReplySds(c,sdsnew("-ERR can't quit, problems saving the DB\r\n")); } @@ -2612,9 +3009,9 @@ static void pushGenericCommand(redisClient *c, int where) { lobj = createListObject(); list = lobj->ptr; if (where == REDIS_HEAD) { - if (!listAddNodeHead(list,c->argv[2])) oom("listAddNodeHead"); + listAddNodeHead(list,c->argv[2]); } else { - if (!listAddNodeTail(list,c->argv[2])) oom("listAddNodeTail"); + listAddNodeTail(list,c->argv[2]); } dictAdd(c->db->dict,c->argv[1],lobj); incrRefCount(c->argv[1]); @@ -2626,9 +3023,9 @@ static void pushGenericCommand(redisClient *c, int where) { } list = lobj->ptr; if (where == REDIS_HEAD) { - if (!listAddNodeHead(list,c->argv[2])) oom("listAddNodeHead"); + listAddNodeHead(list,c->argv[2]); } else { - if (!listAddNodeTail(list,c->argv[2])) oom("listAddNodeTail"); + listAddNodeTail(list,c->argv[2]); } incrRefCount(c->argv[2]); } @@ -2681,7 +3078,7 @@ static void lindexCommand(redisClient *c) { addReply(c,shared.nullbulk); } else { robj *ele = listNodeValue(ln); - addReplySds(c,sdscatprintf(sdsempty(),"$%d\r\n",(int)sdslen(ele->ptr))); + addReplyBulkLen(c,ele); addReply(c,ele); addReply(c,shared.crlf); } @@ -2741,7 +3138,7 @@ static void popGenericCommand(redisClient *c, int where) { addReply(c,shared.nullbulk); } else { robj *ele = listNodeValue(ln); - addReplySds(c,sdscatprintf(sdsempty(),"$%d\r\n",(int)sdslen(ele->ptr))); + addReplyBulkLen(c,ele); addReply(c,ele); addReply(c,shared.crlf); listDelNode(list,ln); @@ -2797,7 +3194,7 @@ static void lrangeCommand(redisClient *c) { addReplySds(c,sdscatprintf(sdsempty(),"*%d\r\n",rangelen)); for (j = 0; j < rangelen; j++) { ele = listNodeValue(ln); - addReplySds(c,sdscatprintf(sdsempty(),"$%d\r\n",(int)sdslen(ele->ptr))); + addReplyBulkLen(c,ele); addReply(c,ele); addReply(c,shared.crlf); ln = ln->next; @@ -2849,8 +3246,8 @@ static void ltrimCommand(redisClient *c) { ln = listLast(list); listDelNode(list,ln); } - addReply(c,shared.ok); server.dirty++; + addReply(c,shared.ok); } } } @@ -2880,7 +3277,7 @@ static void lremCommand(redisClient *c) { robj *ele = listNodeValue(ln); next = fromtail ? ln->prev : ln->next; - if (sdscmp(ele->ptr,c->argv[3]->ptr) == 0) { + if (compareStringObjects(ele,c->argv[3]) == 0) { listDelNode(list,ln); server.dirty++; removed++; @@ -2931,6 +3328,7 @@ static void sremCommand(redisClient *c) { } if (dictDelete(set->ptr,c->argv[2]) == DICT_OK) { server.dirty++; + if (htNeedsResize(set->ptr)) dictResize(set->ptr); addReply(c,shared.cone); } else { addReply(c,shared.czero); @@ -3010,6 +3408,59 @@ static void scardCommand(redisClient *c) { } } +static void spopCommand(redisClient *c) { + robj *set; + dictEntry *de; + + set = lookupKeyWrite(c->db,c->argv[1]); + if (set == NULL) { + addReply(c,shared.nullbulk); + } else { + if (set->type != REDIS_SET) { + addReply(c,shared.wrongtypeerr); + return; + } + de = dictGetRandomKey(set->ptr); + if (de == NULL) { + addReply(c,shared.nullbulk); + } else { + robj *ele = dictGetEntryKey(de); + + addReplyBulkLen(c,ele); + addReply(c,ele); + addReply(c,shared.crlf); + dictDelete(set->ptr,ele); + if (htNeedsResize(set->ptr)) dictResize(set->ptr); + server.dirty++; + } + } +} + +static void srandmemberCommand(redisClient *c) { + robj *set; + dictEntry *de; + + set = lookupKeyRead(c->db,c->argv[1]); + if (set == NULL) { + addReply(c,shared.nullbulk); + } else { + if (set->type != REDIS_SET) { + addReply(c,shared.wrongtypeerr); + return; + } + de = dictGetRandomKey(set->ptr); + if (de == NULL) { + addReply(c,shared.nullbulk); + } else { + robj *ele = dictGetEntryKey(de); + + addReplyBulkLen(c,ele); + addReply(c,ele); + addReply(c,shared.crlf); + } + } +} + static int qsortCompareSetsByCardinality(const void *s1, const void *s2) { dict **d1 = (void*) s1, **d2 = (void*) s2; @@ -3023,7 +3474,6 @@ static void sinterGenericCommand(redisClient *c, robj **setskeys, int setsnum, r robj *lenobj = NULL, *dstset = NULL; int j, cardinality = 0; - if (!dv) oom("sinterGenericCommand"); for (j = 0; j < setsnum; j++) { robj *setobj; @@ -3070,7 +3520,6 @@ static void sinterGenericCommand(redisClient *c, robj **setskeys, int setsnum, r * the element against all the other sets, if at least one set does * not include the element it is discarded */ di = dictGetIterator(dv[0]); - if (!di) oom("dictGetIterator"); while((de = dictNext(di)) != NULL) { robj *ele; @@ -3081,7 +3530,7 @@ static void sinterGenericCommand(redisClient *c, robj **setskeys, int setsnum, r continue; /* at least one set does not contain the member */ ele = dictGetEntryKey(de); if (!dstkey) { - addReplySds(c,sdscatprintf(sdsempty(),"$%d\r\n",sdslen(ele->ptr))); + addReplyBulkLen(c,ele); addReply(c,ele); addReply(c,shared.crlf); cardinality++; @@ -3127,7 +3576,6 @@ static void sunionDiffGenericCommand(redisClient *c, robj **setskeys, int setsnu robj *dstset = NULL; int j, cardinality = 0; - if (!dv) oom("sunionDiffGenericCommand"); for (j = 0; j < setsnum; j++) { robj *setobj; @@ -3158,7 +3606,6 @@ static void sunionDiffGenericCommand(redisClient *c, robj **setskeys, int setsnu if (!dv[j]) continue; /* non existing keys are like empty sets */ di = dictGetIterator(dv[j]); - if (!di) oom("dictGetIterator"); while((de = dictNext(di)) != NULL) { robj *ele; @@ -3185,13 +3632,11 @@ static void sunionDiffGenericCommand(redisClient *c, robj **setskeys, int setsnu if (!dstkey) { addReplySds(c,sdscatprintf(sdsempty(),"*%d\r\n",cardinality)); di = dictGetIterator(dstset->ptr); - if (!di) oom("dictGetIterator"); while((de = dictNext(di)) != NULL) { robj *ele; ele = dictGetEntryKey(de); - addReplySds(c,sdscatprintf(sdsempty(), - "$%d\r\n",sdslen(ele->ptr))); + addReplyBulkLen(c,ele); addReply(c,ele); addReply(c,shared.crlf); } @@ -3231,6 +3676,151 @@ static void sdiffstoreCommand(redisClient *c) { sunionDiffGenericCommand(c,c->argv+2,c->argc-2,c->argv[1],REDIS_OP_DIFF); } +/* ==================================== ZSets =============================== */ + +/* ZSETs are ordered sets using two data structures to hold the same elements + * in order to get O(log(N)) INSERT and REMOVE operations into a sorted + * data structure. + * + * The elements are added to an hash table mapping Redis objects to scores. + * At the same time the elements are added to a skip list mapping scores + * to Redis objects (so objects are sorted by scores in this "view"). */ + +/* This skiplist implementation is almost a C translation of the original + * algorithm described by William Pugh in "Skip Lists: A Probabilistic + * Alternative to Balanced Trees", modified in three ways: + * a) this implementation allows for repeated values. + * b) the comparison is not just by key (our 'score') but by satellite data. + * c) there is a back pointer, so it's a doubly linked list with the back + * pointers being only at "level 1". This allows to traverse the list + * from tail to head, useful for ZREVRANGE. */ + +static zskiplistNode *zslCreateNode(int level, double score, robj *obj) { + zskiplistNode *zn = zmalloc(sizeof(*zn)); + + zn->forward = zmalloc(sizeof(zskiplistNode*) * level); + zn->score = score; + zn->obj = obj; + return zn; +} + +static zskiplist *zslCreate(void) { + int j; + zskiplist *zsl; + + zsl = zmalloc(sizeof(*zsl)); + zsl->level = 1; + zsl->header = zslCreateNode(ZSKIPLIST_MAXLEVEL,0,NULL); + for (j = 0; j < ZSKIPLIST_MAXLEVEL; j++) + zsl->header->forward[j] = NULL; + return zsl; +} + +static void zslFreeNode(zskiplistNode *node) { + decrRefCount(node->obj); + zfree(node); +} + +static void zslFree(zskiplist *zsl) { + zskiplistNode *node = zsl->header->forward[1], *next; + + while(node) { + next = node->forward[1]; + zslFreeNode(node); + node = next; + } +} + +static int zslRandomLevel(void) { + int level = 1; + while ((random()&0xFFFF) < (ZSKIPLIST_P * 0xFFFF)) + level += 1; + return level; +} + +static void zslInsert(zskiplist *zsl, double score, robj *obj) { + zskiplistNode *update[ZSKIPLIST_MAXLEVEL], *x; + int i, level; + + x = zsl->header; + for (i = zsl->level-1; i >= 0; i--) { + while (x->forward[i] && x->forward[i]->score < score) + x = x->forward[i]; + update[i] = x; + } + x = x->forward[1]; + /* we assume the key is not already inside, since we allow duplicated + * scores, and the re-insertion of score and redis object should never + * happpen since the caller of zslInsert() should test in the hash table + * if the element is already inside or not. */ + level = zslRandomLevel(); + if (level > zsl->level) { + for (i = zsl->level; i < level; i++) + update[i] = zsl->header; + zsl->level = level; + } + x = zslCreateNode(level,score,obj); + for (i = 0; i < level; i++) { + x->forward[i] = update[i]->forward[i]; + update[i]->forward[i] = x; + } +} + +static int zslDelete(zskiplist *zsl, double score, robj *obj) { + return 1; +} + +/* The actual Z-commands implementations */ + +static void zaddCommand(redisClient *c) { + robj *zsetobj; + zset *zs; + double *score; + + zsetobj = lookupKeyWrite(c->db,c->argv[1]); + if (zsetobj == NULL) { + zsetobj = createZsetObject(); + dictAdd(c->db->dict,c->argv[1],zsetobj); + incrRefCount(c->argv[1]); + } else { + if (zsetobj->type != REDIS_ZSET) { + addReply(c,shared.wrongtypeerr); + return; + } + } + score = zmalloc(sizeof(double)); + *score = strtod(c->argv[2]->ptr,NULL); + zs = zsetobj->ptr; + if (dictAdd(zs->dict,c->argv[3],score) == DICT_OK) { + /* case 1: New element */ + incrRefCount(c->argv[3]); /* added to hash */ + zslInsert(zs->zsl,*score,c->argv[3]); + incrRefCount(c->argv[3]); /* added to skiplist */ + server.dirty++; + addReply(c,shared.cone); + } else { + dictEntry *de; + double *oldscore; + + /* case 2: Score update operation */ + de = dictFind(zs->dict,c->argv[3]); + assert(de != NULL); + oldscore = dictGetEntryVal(de); + if (*score != *oldscore) { + int deleted; + + deleted = zslDelete(zs->zsl,*score,c->argv[3]); + assert(deleted != 0); + zslInsert(zs->zsl,*score,c->argv[3]); + incrRefCount(c->argv[3]); + server.dirty++; + } + addReply(c,shared.czero); + } +} + +/* ========================= Non type-specific commands ==================== */ + static void flushdbCommand(redisClient *c) { server.dirty += dictSize(c->db->dict); dictEmpty(c->db->dict); @@ -3245,9 +3835,8 @@ static void flushallCommand(redisClient *c) { server.dirty++; } -redisSortOperation *createSortOperation(int type, robj *pattern) { +static redisSortOperation *createSortOperation(int type, robj *pattern) { redisSortOperation *so = zmalloc(sizeof(*so)); - if (!so) oom("createSortOperation"); so->type = type; so->pattern = pattern; return so; @@ -3255,7 +3844,7 @@ redisSortOperation *createSortOperation(int type, robj *pattern) { /* Return the value associated to the key with a name obtained * substituting the first occurence of '*' in 'pattern' with 'subst' */ -robj *lookupKeyByPattern(redisDb *db, robj *pattern, robj *subst) { +static robj *lookupKeyByPattern(redisDb *db, robj *pattern, robj *subst) { char *p; sds spat, ssub; robj keyobj; @@ -3267,6 +3856,12 @@ robj *lookupKeyByPattern(redisDb *db, robj *pattern, robj *subst) { char buf[REDIS_SORTKEY_MAX+1]; } keyname; + if (subst->encoding == REDIS_ENCODING_RAW) + incrRefCount(subst); + else { + subst = getDecodedObject(subst); + } + spat = pattern->ptr; ssub = subst->ptr; if (sdslen(spat)+sdslen(ssub)-1 > REDIS_SORTKEY_MAX) return NULL; @@ -3286,6 +3881,8 @@ robj *lookupKeyByPattern(redisDb *db, robj *pattern, robj *subst) { keyobj.type = REDIS_STRING; keyobj.ptr = ((char*)&keyname)+(sizeof(long)*2); + decrRefCount(subst); + /* printf("lookup '%s' => %p\n", keyname.buf,de); */ return lookupKeyRead(db,&keyobj); } @@ -3323,7 +3920,20 @@ static int sortCompare(const void *s1, const void *s2) { } } else { /* Compare elements directly */ - cmp = strcoll(so1->obj->ptr,so2->obj->ptr); + if (so1->obj->encoding == REDIS_ENCODING_RAW && + so2->obj->encoding == REDIS_ENCODING_RAW) { + cmp = strcoll(so1->obj->ptr,so2->obj->ptr); + } else { + robj *dec1, *dec2; + + dec1 = so1->obj->encoding == REDIS_ENCODING_RAW ? + so1->obj : getDecodedObject(so1->obj); + dec2 = so2->obj->encoding == REDIS_ENCODING_RAW ? + so2->obj : getDecodedObject(so2->obj); + cmp = strcoll(dec1->ptr,dec2->ptr); + if (dec1 != so1->obj) decrRefCount(dec1); + if (dec2 != so2->obj) decrRefCount(dec2); + } } } return server.sort_desc ? -cmp : cmp; @@ -3413,7 +4023,6 @@ static void sortCommand(redisClient *c) { listLength((list*)sortval->ptr) : dictSize((dict*)sortval->ptr); vector = zmalloc(sizeof(redisSortObject)*vectorlen); - if (!vector) oom("allocating objects vector for SORT"); j = 0; if (sortval->type == REDIS_LIST) { list *list = sortval->ptr; @@ -3433,7 +4042,6 @@ static void sortCommand(redisClient *c) { dictEntry *setele; di = dictGetIterator(set); - if (!di) oom("dictGetIterator"); while((setele = dictNext(di)) != NULL) { vector[j].obj = dictGetEntryKey(setele); vector[j].u.score = 0; @@ -3453,13 +4061,33 @@ static void sortCommand(redisClient *c) { byval = lookupKeyByPattern(c->db,sortby,vector[j].obj); if (!byval || byval->type != REDIS_STRING) continue; if (alpha) { - vector[j].u.cmpobj = byval; - incrRefCount(byval); + if (byval->encoding == REDIS_ENCODING_RAW) { + vector[j].u.cmpobj = byval; + incrRefCount(byval); + } else { + vector[j].u.cmpobj = getDecodedObject(byval); + } } else { - vector[j].u.score = strtod(byval->ptr,NULL); + if (byval->encoding == REDIS_ENCODING_RAW) { + vector[j].u.score = strtod(byval->ptr,NULL); + } else { + if (byval->encoding == REDIS_ENCODING_INT) { + vector[j].u.score = (long)byval->ptr; + } else + assert(1 != 1); + } } } else { - if (!alpha) vector[j].u.score = strtod(vector[j].obj->ptr,NULL); + if (!alpha) { + if (vector[j].obj->encoding == REDIS_ENCODING_RAW) + vector[j].u.score = strtod(vector[j].obj->ptr,NULL); + else { + if (vector[j].obj->encoding == REDIS_ENCODING_INT) + vector[j].u.score = (long) vector[j].obj->ptr; + else + assert(1 != 1); + } + } } } } @@ -3491,8 +4119,7 @@ static void sortCommand(redisClient *c) { for (j = start; j <= end; j++) { listNode *ln; if (!getop) { - addReplySds(c,sdscatprintf(sdsempty(),"$%d\r\n", - sdslen(vector[j].obj->ptr))); + addReplyBulkLen(c,vector[j].obj); addReply(c,vector[j].obj); addReply(c,shared.crlf); } @@ -3506,8 +4133,7 @@ static void sortCommand(redisClient *c) { if (!val || val->type != REDIS_STRING) { addReply(c,shared.nullbulk); } else { - addReplySds(c,sdscatprintf(sdsempty(),"$%d\r\n", - sdslen(val->ptr))); + addReplyBulkLen(c,val); addReply(c,val); addReply(c,shared.crlf); } @@ -3530,9 +4156,11 @@ static void sortCommand(redisClient *c) { static void infoCommand(redisClient *c) { sds info; time_t uptime = time(NULL)-server.stat_starttime; + int j; info = sdscatprintf(sdsempty(), "redis_version:%s\r\n" + "arch_bits:%s\r\n" "uptime_in_seconds:%d\r\n" "uptime_in_days:%d\r\n" "connected_clients:%d\r\n" @@ -3545,6 +4173,7 @@ static void infoCommand(redisClient *c) { "total_commands_processed:%lld\r\n" "role:%s\r\n" ,REDIS_VERSION, + (sizeof(long) == 8) ? "64" : "32", uptime, uptime/(3600*24), listLength(server.clients)-listLength(server.slaves), @@ -3570,6 +4199,16 @@ static void infoCommand(redisClient *c) { (int)(time(NULL)-server.master->lastinteraction) ); } + for (j = 0; j < server.dbnum; j++) { + long long keys, vkeys; + + keys = dictSize(server.db[j].dict); + vkeys = dictSize(server.db[j].expires); + if (keys || vkeys) { + info = sdscatprintf(info, "db%d: keys=%lld,expires=%lld\r\n", + j, keys, vkeys); + } + } addReplySds(c,sdscatprintf(sdsempty(),"$%d\r\n",sdslen(info))); addReplySds(c,info); addReply(c,shared.crlf); @@ -3581,7 +4220,7 @@ static void monitorCommand(redisClient *c) { c->flags |= (REDIS_SLAVE|REDIS_MONITOR); c->slaveseldb = 0; - if (!listAddNodeTail(server.monitors,c)) oom("listAddNodeTail"); + listAddNodeTail(server.monitors,c); addReply(c,shared.ok); } @@ -3659,10 +4298,12 @@ static void expireCommand(redisClient *c) { return; } else { time_t when = time(NULL)+seconds; - if (setExpire(c->db,c->argv[1],when)) + if (setExpire(c->db,c->argv[1],when)) { addReply(c,shared.cone); - else + server.dirty++; + } else { addReply(c,shared.czero); + } return; } } @@ -3679,6 +4320,49 @@ static void ttlCommand(redisClient *c) { addReplySds(c,sdscatprintf(sdsempty(),":%d\r\n",ttl)); } +static void msetGenericCommand(redisClient *c, int nx) { + int j; + + if ((c->argc % 2) == 0) { + addReplySds(c,sdsnew("-ERR wrong number of arguments\r\n")); + return; + } + /* Handle the NX flag. The MSETNX semantic is to return zero and don't + * set nothing at all if at least one already key exists. */ + if (nx) { + for (j = 1; j < c->argc; j += 2) { + if (dictFind(c->db->dict,c->argv[j]) != NULL) { + addReply(c, shared.czero); + return; + } + } + } + + for (j = 1; j < c->argc; j += 2) { + int retval; + + retval = dictAdd(c->db->dict,c->argv[j],c->argv[j+1]); + if (retval == DICT_ERR) { + dictReplace(c->db->dict,c->argv[j],c->argv[j+1]); + incrRefCount(c->argv[j+1]); + } else { + incrRefCount(c->argv[j]); + incrRefCount(c->argv[j+1]); + } + removeExpire(c->db,c->argv[j]); + } + server.dirty += (c->argc-1)/2; + addReply(c, nx ? shared.cone : shared.ok); +} + +static void msetCommand(redisClient *c) { + msetGenericCommand(c,0); +} + +static void msetnxCommand(redisClient *c) { + msetGenericCommand(c,1); +} + /* =============================== Replication ============================= */ static int syncWrite(int fd, char *ptr, ssize_t size, int timeout) { @@ -3776,7 +4460,6 @@ static void syncCommand(redisClient *c) { * another slave. Set the right state, and copy the buffer. */ listRelease(c->reply); c->reply = listDup(slave->reply); - if (!c->reply) oom("listDup copying slave reply list"); c->replstate = REDIS_REPL_WAIT_BGSAVE_END; redisLog(REDIS_NOTICE,"Waiting for end of BGSAVE for SYNC"); } else { @@ -3798,7 +4481,7 @@ static void syncCommand(redisClient *c) { c->repldbfd = -1; c->flags |= REDIS_SLAVE; c->slaveseldb = 0; - if (!listAddNodeTail(server.slaves,c)) oom("listAddNodeTail"); + listAddNodeTail(server.slaves,c); return; } @@ -3856,7 +4539,13 @@ static void sendBulkToSlave(aeEventLoop *el, int fd, void *privdata, int mask) { } } -static void updateSalvesWaitingBgsave(int bgsaveerr) { +/* This function is called at the end of every backgrond saving. + * The argument bgsaveerr is REDIS_OK if the background saving succeeded + * otherwise REDIS_ERR is passed to the function. + * + * The goal of this function is to handle slaves waiting for a successful + * background saving in order to perform non-blocking synchronization. */ +static void updateSlavesWaitingBgsave(int bgsaveerr) { listNode *ln; int startbgsave = 0; @@ -4069,14 +4758,249 @@ static void debugCommand(redisClient *c) { key = dictGetEntryKey(de); val = dictGetEntryVal(de); addReplySds(c,sdscatprintf(sdsempty(), - "+Key at:%p refcount:%d, value at:%p refcount:%d\r\n", - key, key->refcount, val, val->refcount)); + "+Key at:%p refcount:%d, value at:%p refcount:%d encoding:%d\r\n", + key, key->refcount, val, val->refcount, val->encoding)); } else { addReplySds(c,sdsnew( "-ERR Syntax error, try DEBUG [SEGFAULT|OBJECT ]\r\n")); } } +#ifdef HAVE_BACKTRACE +static struct redisFunctionSym symsTable[] = { +{"compareStringObjects", (unsigned long)compareStringObjects}, +{"isStringRepresentableAsLong", (unsigned long)isStringRepresentableAsLong}, +{"dictEncObjKeyCompare", (unsigned long)dictEncObjKeyCompare}, +{"dictEncObjHash", (unsigned long)dictEncObjHash}, +{"incrDecrCommand", (unsigned long)incrDecrCommand}, +{"freeStringObject", (unsigned long)freeStringObject}, +{"freeListObject", (unsigned long)freeListObject}, +{"freeSetObject", (unsigned long)freeSetObject}, +{"decrRefCount", (unsigned long)decrRefCount}, +{"createObject", (unsigned long)createObject}, +{"freeClient", (unsigned long)freeClient}, +{"rdbLoad", (unsigned long)rdbLoad}, +{"rdbSaveStringObject", (unsigned long)rdbSaveStringObject}, +{"rdbSaveStringObjectRaw", (unsigned long)rdbSaveStringObjectRaw}, +{"addReply", (unsigned long)addReply}, +{"addReplySds", (unsigned long)addReplySds}, +{"incrRefCount", (unsigned long)incrRefCount}, +{"rdbSaveBackground", (unsigned long)rdbSaveBackground}, +{"createStringObject", (unsigned long)createStringObject}, +{"replicationFeedSlaves", (unsigned long)replicationFeedSlaves}, +{"syncWithMaster", (unsigned long)syncWithMaster}, +{"tryObjectSharing", (unsigned long)tryObjectSharing}, +{"tryObjectEncoding", (unsigned long)tryObjectEncoding}, +{"getDecodedObject", (unsigned long)getDecodedObject}, +{"removeExpire", (unsigned long)removeExpire}, +{"expireIfNeeded", (unsigned long)expireIfNeeded}, +{"deleteIfVolatile", (unsigned long)deleteIfVolatile}, +{"deleteKey", (unsigned long)deleteKey}, +{"getExpire", (unsigned long)getExpire}, +{"setExpire", (unsigned long)setExpire}, +{"updateSlavesWaitingBgsave", (unsigned long)updateSlavesWaitingBgsave}, +{"freeMemoryIfNeeded", (unsigned long)freeMemoryIfNeeded}, +{"authCommand", (unsigned long)authCommand}, +{"pingCommand", (unsigned long)pingCommand}, +{"echoCommand", (unsigned long)echoCommand}, +{"setCommand", (unsigned long)setCommand}, +{"setnxCommand", (unsigned long)setnxCommand}, +{"getCommand", (unsigned long)getCommand}, +{"delCommand", (unsigned long)delCommand}, +{"existsCommand", (unsigned long)existsCommand}, +{"incrCommand", (unsigned long)incrCommand}, +{"decrCommand", (unsigned long)decrCommand}, +{"incrbyCommand", (unsigned long)incrbyCommand}, +{"decrbyCommand", (unsigned long)decrbyCommand}, +{"selectCommand", (unsigned long)selectCommand}, +{"randomkeyCommand", (unsigned long)randomkeyCommand}, +{"keysCommand", (unsigned long)keysCommand}, +{"dbsizeCommand", (unsigned long)dbsizeCommand}, +{"lastsaveCommand", (unsigned long)lastsaveCommand}, +{"saveCommand", (unsigned long)saveCommand}, +{"bgsaveCommand", (unsigned long)bgsaveCommand}, +{"shutdownCommand", (unsigned long)shutdownCommand}, +{"moveCommand", (unsigned long)moveCommand}, +{"renameCommand", (unsigned long)renameCommand}, +{"renamenxCommand", (unsigned long)renamenxCommand}, +{"lpushCommand", (unsigned long)lpushCommand}, +{"rpushCommand", (unsigned long)rpushCommand}, +{"lpopCommand", (unsigned long)lpopCommand}, +{"rpopCommand", (unsigned long)rpopCommand}, +{"llenCommand", (unsigned long)llenCommand}, +{"lindexCommand", (unsigned long)lindexCommand}, +{"lrangeCommand", (unsigned long)lrangeCommand}, +{"ltrimCommand", (unsigned long)ltrimCommand}, +{"typeCommand", (unsigned long)typeCommand}, +{"lsetCommand", (unsigned long)lsetCommand}, +{"saddCommand", (unsigned long)saddCommand}, +{"sremCommand", (unsigned long)sremCommand}, +{"smoveCommand", (unsigned long)smoveCommand}, +{"sismemberCommand", (unsigned long)sismemberCommand}, +{"scardCommand", (unsigned long)scardCommand}, +{"spopCommand", (unsigned long)spopCommand}, +{"srandmemberCommand", (unsigned long)srandmemberCommand}, +{"sinterCommand", (unsigned long)sinterCommand}, +{"sinterstoreCommand", (unsigned long)sinterstoreCommand}, +{"sunionCommand", (unsigned long)sunionCommand}, +{"sunionstoreCommand", (unsigned long)sunionstoreCommand}, +{"sdiffCommand", (unsigned long)sdiffCommand}, +{"sdiffstoreCommand", (unsigned long)sdiffstoreCommand}, +{"syncCommand", (unsigned long)syncCommand}, +{"flushdbCommand", (unsigned long)flushdbCommand}, +{"flushallCommand", (unsigned long)flushallCommand}, +{"sortCommand", (unsigned long)sortCommand}, +{"lremCommand", (unsigned long)lremCommand}, +{"infoCommand", (unsigned long)infoCommand}, +{"mgetCommand", (unsigned long)mgetCommand}, +{"monitorCommand", (unsigned long)monitorCommand}, +{"expireCommand", (unsigned long)expireCommand}, +{"getsetCommand", (unsigned long)getsetCommand}, +{"ttlCommand", (unsigned long)ttlCommand}, +{"slaveofCommand", (unsigned long)slaveofCommand}, +{"debugCommand", (unsigned long)debugCommand}, +{"processCommand", (unsigned long)processCommand}, +{"setupSigSegvAction", (unsigned long)setupSigSegvAction}, +{"readQueryFromClient", (unsigned long)readQueryFromClient}, +{"rdbRemoveTempFile", (unsigned long)rdbRemoveTempFile}, +{"msetGenericCommand", (unsigned long)msetGenericCommand}, +{"msetCommand", (unsigned long)msetCommand}, +{"msetnxCommand", (unsigned long)msetnxCommand}, +{"zslCreateNode", (unsigned long)zslCreateNode}, +{"zslCreate", (unsigned long)zslCreate}, +{"zslFreeNode",(unsigned long)zslFreeNode}, +{"zslFree",(unsigned long)zslFree}, +{"zslRandomLevel",(unsigned long)zslRandomLevel}, +{"zslInsert",(unsigned long)zslInsert}, +{"zslDelete",(unsigned long)zslDelete}, +{"createZsetObject",(unsigned long)createZsetObject}, +{"zaddCommand",(unsigned long)zaddCommand}, +{NULL,0} +}; + +/* This function try to convert a pointer into a function name. It's used in + * oreder to provide a backtrace under segmentation fault that's able to + * display functions declared as static (otherwise the backtrace is useless). */ +static char *findFuncName(void *pointer, unsigned long *offset){ + int i, ret = -1; + unsigned long off, minoff = 0; + + /* Try to match against the Symbol with the smallest offset */ + for (i=0; symsTable[i].pointer; i++) { + unsigned long lp = (unsigned long) pointer; + + if (lp != (unsigned long)-1 && lp >= symsTable[i].pointer) { + off=lp-symsTable[i].pointer; + if (ret < 0 || off < minoff) { + minoff=off; + ret=i; + } + } + } + if (ret == -1) return NULL; + *offset = minoff; + return symsTable[ret].name; +} + +static void *getMcontextEip(ucontext_t *uc) { +#if defined(__FreeBSD__) + return (void*) uc->uc_mcontext.mc_eip; +#elif defined(__dietlibc__) + return (void*) uc->uc_mcontext.eip; +#elif defined(__APPLE__) && !defined(MAC_OS_X_VERSION_10_6) + return (void*) uc->uc_mcontext->__ss.__eip; +#elif defined(__APPLE__) && defined(MAC_OS_X_VERSION_10_6) + #if defined(_STRUCT_X86_THREAD_STATE64) && !defined(__i386__) + return (void*) uc->uc_mcontext->__ss.__rip; + #else + return (void*) uc->uc_mcontext->__ss.__eip; + #endif +#elif defined(__i386__) || defined(__X86_64__) /* Linux x86 */ + return (void*) uc->uc_mcontext.gregs[REG_EIP]; +#elif defined(__ia64__) /* Linux IA64 */ + return (void*) uc->uc_mcontext.sc_ip; +#else + return NULL; +#endif +} + +static void segvHandler(int sig, siginfo_t *info, void *secret) { + void *trace[100]; + char **messages = NULL; + int i, trace_size = 0; + unsigned long offset=0; + time_t uptime = time(NULL)-server.stat_starttime; + ucontext_t *uc = (ucontext_t*) secret; + REDIS_NOTUSED(info); + + redisLog(REDIS_WARNING, + "======= Ooops! Redis %s got signal: -%d- =======", REDIS_VERSION, sig); + redisLog(REDIS_WARNING, "%s", sdscatprintf(sdsempty(), + "redis_version:%s; " + "uptime_in_seconds:%d; " + "connected_clients:%d; " + "connected_slaves:%d; " + "used_memory:%zu; " + "changes_since_last_save:%lld; " + "bgsave_in_progress:%d; " + "last_save_time:%d; " + "total_connections_received:%lld; " + "total_commands_processed:%lld; " + "role:%s;" + ,REDIS_VERSION, + uptime, + listLength(server.clients)-listLength(server.slaves), + listLength(server.slaves), + server.usedmemory, + server.dirty, + server.bgsaveinprogress, + server.lastsave, + server.stat_numconnections, + server.stat_numcommands, + server.masterhost == NULL ? "master" : "slave" + )); + + trace_size = backtrace(trace, 100); + /* overwrite sigaction with caller's address */ + if (getMcontextEip(uc) != NULL) { + trace[1] = getMcontextEip(uc); + } + messages = backtrace_symbols(trace, trace_size); + + for (i=1; i /proc/sys/vm/overcommit_memory' in your init scripts."); + redisLog(REDIS_WARNING,"WARNING overcommit_memory is set to 0! Background save may fail under low condition memory. To fix this issue add 'vm.overcommit_memory = 1' to /etc/sysctl.conf and then reboot or run the command 'sysctl vm.overcommit_memory=1' for this to take effect."); } } #endif /* __linux__ */ @@ -4126,10 +5050,6 @@ static void daemonize(void) { } int main(int argc, char **argv) { -#ifdef __linux__ - linuxOvercommitMemoryWarning(); -#endif - initServerConfig(); if (argc == 2) { ResetServerSaveParams(); @@ -4143,6 +5063,9 @@ int main(int argc, char **argv) { initServer(); if (server.daemonize) daemonize(); redisLog(REDIS_NOTICE,"Server started, Redis version " REDIS_VERSION); +#ifdef __linux__ + linuxOvercommitMemoryWarning(); +#endif if (rdbLoad(server.dbfilename) == REDIS_OK) redisLog(REDIS_NOTICE,"DB loaded from disk"); if (aeCreateFileEvent(server.el, server.fd, AE_READABLE,