[6.0] src/db.c
This commit is contained in:
@@ -63,10 +63,7 @@ robj *lookupKey(redisDb *db, robj *key, int flags) {
|
||||
/* Update the access time for the ageing algorithm.
|
||||
* Don't do it if we have a saving child, as this will trigger
|
||||
* a copy on write madness. */
|
||||
if (server.rdb_child_pid == -1 &&
|
||||
server.aof_child_pid == -1 &&
|
||||
!(flags & LOOKUP_NOTOUCH))
|
||||
{
|
||||
if (!hasActiveChildProcess() && !(flags & LOOKUP_NOTOUCH)){
|
||||
if (server.maxmemory_policy & MAXMEMORY_FLAG_LFU) {
|
||||
updateLFU(val);
|
||||
} else {
|
||||
@@ -86,6 +83,7 @@ robj *lookupKey(redisDb *db, robj *key, int flags) {
|
||||
* 1. A key gets expired if it reached it's TTL.
|
||||
* 2. The key last access time is updated.
|
||||
* 3. The global keys hits/misses stats are updated (reported in INFO).
|
||||
* 4. If keyspace notifications are enabled, a "keymiss" notification is fired.
|
||||
*
|
||||
* This API should not be used when we write to the key after obtaining
|
||||
* the object linked to the key, but only for read only operations.
|
||||
@@ -109,6 +107,7 @@ robj *lookupKeyReadWithFlags(redisDb *db, robj *key, int flags) {
|
||||
* to return NULL ASAP. */
|
||||
if (server.masterhost == NULL) {
|
||||
server.stat_keyspace_misses++;
|
||||
notifyKeyspaceEvent(NOTIFY_KEY_MISS, "keymiss", key, db->id);
|
||||
return NULL;
|
||||
}
|
||||
|
||||
@@ -120,7 +119,7 @@ robj *lookupKeyReadWithFlags(redisDb *db, robj *key, int flags) {
|
||||
* However, if the command caller is not the master, and as additional
|
||||
* safety measure, the command invoked is a read-only command, we can
|
||||
* safely return NULL here, and provide a more consistent behavior
|
||||
* to clients accessign expired values in a read-only fashion, that
|
||||
* to clients accessing expired values in a read-only fashion, that
|
||||
* will say the key as non existing.
|
||||
*
|
||||
* Notably this covers GETs when slaves are used to scale reads. */
|
||||
@@ -130,12 +129,15 @@ robj *lookupKeyReadWithFlags(redisDb *db, robj *key, int flags) {
|
||||
server.current_client->cmd->flags & CMD_READONLY)
|
||||
{
|
||||
server.stat_keyspace_misses++;
|
||||
notifyKeyspaceEvent(NOTIFY_KEY_MISS, "keymiss", key, db->id);
|
||||
return NULL;
|
||||
}
|
||||
}
|
||||
val = lookupKey(db,key,flags);
|
||||
if (val == NULL)
|
||||
if (val == NULL) {
|
||||
server.stat_keyspace_misses++;
|
||||
notifyKeyspaceEvent(NOTIFY_KEY_MISS, "keymiss", key, db->id);
|
||||
}
|
||||
else
|
||||
server.stat_keyspace_hits++;
|
||||
return val;
|
||||
@@ -152,9 +154,13 @@ robj *lookupKeyRead(redisDb *db, robj *key) {
|
||||
*
|
||||
* Returns the linked value object if the key exists or NULL if the key
|
||||
* does not exist in the specified DB. */
|
||||
robj *lookupKeyWrite(redisDb *db, robj *key) {
|
||||
robj *lookupKeyWriteWithFlags(redisDb *db, robj *key, int flags) {
|
||||
expireIfNeeded(db,key);
|
||||
return lookupKey(db,key,LOOKUP_NONE);
|
||||
return lookupKey(db,key,flags);
|
||||
}
|
||||
|
||||
robj *lookupKeyWrite(redisDb *db, robj *key) {
|
||||
return lookupKeyWriteWithFlags(db, key, LOOKUP_NONE);
|
||||
}
|
||||
|
||||
robj *lookupKeyReadOrReply(client *c, robj *key, robj *reply) {
|
||||
@@ -179,9 +185,28 @@ void dbAdd(redisDb *db, robj *key, robj *val) {
|
||||
|
||||
serverAssertWithInfo(NULL,key,retval == DICT_OK);
|
||||
if (val->type == OBJ_LIST ||
|
||||
val->type == OBJ_ZSET)
|
||||
val->type == OBJ_ZSET ||
|
||||
val->type == OBJ_STREAM)
|
||||
signalKeyAsReady(db, key);
|
||||
if (server.cluster_enabled) slotToKeyAdd(key->ptr);
|
||||
}
|
||||
|
||||
/* This is a special version of dbAdd() that is used only when loading
|
||||
* keys from the RDB file: the key is passed as an SDS string that is
|
||||
* retained by the function (and not freed by the caller).
|
||||
*
|
||||
* Moreover this function will not abort if the key is already busy, to
|
||||
* give more control to the caller, nor will signal the key as ready
|
||||
* since it is not useful in this context.
|
||||
*
|
||||
* The function returns 1 if the key was added to the database, taking
|
||||
* ownership of the SDS string, otherwise 0 is returned, and is up to the
|
||||
* caller to free the SDS string. */
|
||||
int dbAddRDBLoad(redisDb *db, sds key, robj *val) {
|
||||
int retval = dictAdd(db->dict, key, val);
|
||||
if (retval != DICT_OK) return 0;
|
||||
if (server.cluster_enabled) slotToKeyAdd(key);
|
||||
return 1;
|
||||
}
|
||||
|
||||
/* Overwrite an existing key with a new value. Incrementing the reference
|
||||
@@ -213,20 +238,30 @@ void dbOverwrite(redisDb *db, robj *key, robj *val) {
|
||||
*
|
||||
* 1) The ref count of the value object is incremented.
|
||||
* 2) clients WATCHing for the destination key notified.
|
||||
* 3) The expire time of the key is reset (the key is made persistent).
|
||||
* 3) The expire time of the key is reset (the key is made persistent),
|
||||
* unless 'keepttl' is true.
|
||||
*
|
||||
* All the new keys in the database should be created via this interface. */
|
||||
void setKey(redisDb *db, robj *key, robj *val) {
|
||||
* All the new keys in the database should be created via this interface.
|
||||
* The client 'c' argument may be set to NULL if the operation is performed
|
||||
* in a context where there is no clear client performing the operation. */
|
||||
void genericSetKey(client *c, redisDb *db, robj *key, robj *val, int keepttl, int signal) {
|
||||
if (lookupKeyWrite(db,key) == NULL) {
|
||||
dbAdd(db,key,val);
|
||||
} else {
|
||||
dbOverwrite(db,key,val);
|
||||
}
|
||||
incrRefCount(val);
|
||||
removeExpire(db,key);
|
||||
signalModifiedKey(db,key);
|
||||
if (!keepttl) removeExpire(db,key);
|
||||
if (signal) signalModifiedKey(c,db,key);
|
||||
}
|
||||
|
||||
/* Common case for genericSetKey() where the TTL is not retained. */
|
||||
void setKey(client *c, redisDb *db, robj *key, robj *val) {
|
||||
genericSetKey(c,db,key,val,0,1);
|
||||
}
|
||||
|
||||
/* Return true if the specified key exists in the specified database.
|
||||
* LRU/LFU info is not updated in any way. */
|
||||
int dbExists(redisDb *db, robj *key) {
|
||||
return dictFind(db->dict,key->ptr) != NULL;
|
||||
}
|
||||
@@ -244,7 +279,7 @@ robj *dbRandomKey(redisDb *db) {
|
||||
sds key;
|
||||
robj *keyobj;
|
||||
|
||||
de = dictGetRandomKey(db->dict);
|
||||
de = dictGetFairRandomKey(db->dict);
|
||||
if (de == NULL) return NULL;
|
||||
|
||||
key = dictGetKey(de);
|
||||
@@ -276,7 +311,7 @@ int dbSyncDelete(redisDb *db, robj *key) {
|
||||
* the key, because it is shared with the main dictionary. */
|
||||
if (dictSize(db->expires) > 0) dictDelete(db->expires,key->ptr);
|
||||
if (dictDelete(db->dict,key->ptr) == DICT_OK) {
|
||||
if (server.cluster_enabled) slotToKeyDel(key);
|
||||
if (server.cluster_enabled) slotToKeyDel(key->ptr);
|
||||
return 1;
|
||||
} else {
|
||||
return 0;
|
||||
@@ -336,14 +371,19 @@ robj *dbUnshareStringValue(redisDb *db, robj *key, robj *o) {
|
||||
* DB number if we want to flush only a single Redis database number.
|
||||
*
|
||||
* Flags are be EMPTYDB_NO_FLAGS if no special flags are specified or
|
||||
* EMPTYDB_ASYNC if we want the memory to be freed in a different thread
|
||||
* 1. EMPTYDB_ASYNC if we want the memory to be freed in a different thread.
|
||||
* 2. EMPTYDB_BACKUP if we want to empty the backup dictionaries created by
|
||||
* disklessLoadMakeBackups. In that case we only free memory and avoid
|
||||
* firing module events.
|
||||
* and the function to return ASAP.
|
||||
*
|
||||
* On success the fuction returns the number of keys removed from the
|
||||
* On success the function returns the number of keys removed from the
|
||||
* database(s). Otherwise -1 is returned in the specific case the
|
||||
* DB number is out of range, and errno is set to EINVAL. */
|
||||
PORT_LONGLONG emptyDb(int dbnum, int flags, void(callback)(void*)) {
|
||||
PORT_LONGLONG emptyDb(redisDb *dbarray, int dbnum, int flags, void(callback)(void*)) {
|
||||
int async = (flags & EMPTYDB_ASYNC);
|
||||
int backup = (flags & EMPTYDB_BACKUP); /* Just free the memory, nothing else */
|
||||
RedisModuleFlushInfoV1 fi = {REDISMODULE_FLUSHINFO_VERSION,!async,dbnum};
|
||||
PORT_LONGLONG removed = 0;
|
||||
|
||||
if (dbnum < -1 || dbnum >= server.dbnum) {
|
||||
@@ -351,6 +391,19 @@ PORT_LONGLONG emptyDb(int dbnum, int flags, void(callback)(void*)) {
|
||||
return -1;
|
||||
}
|
||||
|
||||
/* Pre-flush actions */
|
||||
if (!backup) {
|
||||
/* Fire the flushdb modules event. */
|
||||
moduleFireServerEvent(REDISMODULE_EVENT_FLUSHDB,
|
||||
REDISMODULE_SUBEVENT_FLUSHDB_START,
|
||||
&fi);
|
||||
|
||||
/* Make sure the WATCHed keys are affected by the FLUSH* commands.
|
||||
* Note that we need to call the function while the keys are still
|
||||
* there. */
|
||||
signalFlushedDb(dbnum);
|
||||
}
|
||||
|
||||
int startdb, enddb;
|
||||
if (dbnum == -1) {
|
||||
startdb = 0;
|
||||
@@ -360,25 +413,40 @@ PORT_LONGLONG emptyDb(int dbnum, int flags, void(callback)(void*)) {
|
||||
}
|
||||
|
||||
for (int j = startdb; j <= enddb; j++) {
|
||||
removed += dictSize(server.db[j].dict);
|
||||
removed += dictSize(dbarray[j].dict);
|
||||
if (async) {
|
||||
emptyDbAsync(&server.db[j]);
|
||||
emptyDbAsync(&dbarray[j]);
|
||||
} else {
|
||||
dictEmpty(server.db[j].dict,callback);
|
||||
dictEmpty(server.db[j].expires,callback);
|
||||
dictEmpty(dbarray[j].dict,callback);
|
||||
dictEmpty(dbarray[j].expires,callback);
|
||||
}
|
||||
}
|
||||
if (server.cluster_enabled) {
|
||||
if (async) {
|
||||
slotToKeyFlushAsync();
|
||||
} else {
|
||||
slotToKeyFlush();
|
||||
|
||||
/* Post-flush actions */
|
||||
if (!backup) {
|
||||
if (server.cluster_enabled) {
|
||||
if (async) {
|
||||
slotToKeyFlushAsync();
|
||||
} else {
|
||||
slotToKeyFlush();
|
||||
}
|
||||
}
|
||||
if (dbnum == -1) flushSlaveKeysWithExpireList();
|
||||
|
||||
/* Also fire the end event. Note that this event will fire almost
|
||||
* immediately after the start event if the flush is asynchronous. */
|
||||
moduleFireServerEvent(REDISMODULE_EVENT_FLUSHDB,
|
||||
REDISMODULE_SUBEVENT_FLUSHDB_END,
|
||||
&fi);
|
||||
}
|
||||
if (dbnum == -1) flushSlaveKeysWithExpireList();
|
||||
|
||||
return removed;
|
||||
}
|
||||
|
||||
PORT_LONGLONG emptyDb(int dbnum, int flags, void(callback)(void*)) {
|
||||
return emptyDbGeneric(server.db, dbnum, flags, callback);
|
||||
}
|
||||
|
||||
int selectDb(client *c, int id) {
|
||||
if (id < 0 || id >= server.dbnum)
|
||||
return C_ERR;
|
||||
@@ -386,6 +454,15 @@ int selectDb(client *c, int id) {
|
||||
return C_OK;
|
||||
}
|
||||
|
||||
PORT_LONGLONG dbTotalServerKeyCount() {
|
||||
PORT_LONGLONG total = 0;
|
||||
int j;
|
||||
for (j = 0; j < server.dbnum; j++) {
|
||||
total += dictSize(server.db[j].dict);
|
||||
}
|
||||
return total;
|
||||
}
|
||||
|
||||
/*-----------------------------------------------------------------------------
|
||||
* Hooks for key space changes.
|
||||
*
|
||||
@@ -395,12 +472,16 @@ int selectDb(client *c, int id) {
|
||||
* Every time a DB is flushed the function signalFlushDb() is called.
|
||||
*----------------------------------------------------------------------------*/
|
||||
|
||||
void signalModifiedKey(redisDb *db, robj *key) {
|
||||
/* Note that the 'c' argument may be NULL if the key was modified out of
|
||||
* a context of a client. */
|
||||
void signalModifiedKey(client *c, redisDb *db, robj *key) {
|
||||
touchWatchedKey(db,key);
|
||||
trackingInvalidateKey(c,key);
|
||||
}
|
||||
|
||||
void signalFlushedDb(int dbid) {
|
||||
touchWatchedKeysOnFlush(dbid);
|
||||
trackingInvalidateKeysOnFlush(dbid);
|
||||
}
|
||||
|
||||
/*-----------------------------------------------------------------------------
|
||||
@@ -429,32 +510,10 @@ int getFlushCommandFlags(client *c, int *flags) {
|
||||
return C_OK;
|
||||
}
|
||||
|
||||
/* FLUSHDB [ASYNC]
|
||||
*
|
||||
* Flushes the currently SELECTed Redis DB. */
|
||||
void flushdbCommand(client *c) {
|
||||
int flags;
|
||||
|
||||
if (getFlushCommandFlags(c,&flags) == C_ERR) return;
|
||||
signalFlushedDb(c->db->id);
|
||||
server.dirty += emptyDb(c->db->id,flags,NULL);
|
||||
addReply(c,shared.ok);
|
||||
}
|
||||
|
||||
/* FLUSHALL [ASYNC]
|
||||
*
|
||||
* Flushes the whole server data set. */
|
||||
void flushallCommand(client *c) {
|
||||
int flags;
|
||||
|
||||
if (getFlushCommandFlags(c,&flags) == C_ERR) return;
|
||||
signalFlushedDb(-1);
|
||||
/* Flushes the whole server data set. */
|
||||
void flushAllDataAndResetRDB(int flags) {
|
||||
server.dirty += emptyDb(-1,flags,NULL);
|
||||
addReply(c,shared.ok);
|
||||
if (server.rdb_child_pid != -1) {
|
||||
IF_WIN32(AbortForkOperation(), kill(server.rdb_child_pid,SIGUSR1));
|
||||
rdbRemoveTempFile(server.rdb_child_pid);
|
||||
}
|
||||
if (server.rdb_child_pid != -1) killRDBChild();
|
||||
if (server.saveparamslen > 0) {
|
||||
/* Normally rdbSave() will reset dirty, but we don't want this here
|
||||
* as otherwise FLUSHALL will not be replicated nor put into the AOF. */
|
||||
@@ -465,6 +524,41 @@ void flushallCommand(client *c) {
|
||||
server.dirty = saved_dirty;
|
||||
}
|
||||
server.dirty++;
|
||||
#if defined(USE_JEMALLOC)
|
||||
/* jemalloc 5 doesn't release pages back to the OS when there's no traffic.
|
||||
* for large databases, flushdb blocks for long anyway, so a bit more won't
|
||||
* harm and this way the flush and purge will be synchroneus. */
|
||||
if (!(flags & EMPTYDB_ASYNC))
|
||||
jemalloc_purge();
|
||||
#endif
|
||||
}
|
||||
|
||||
/* FLUSHDB [ASYNC]
|
||||
*
|
||||
* Flushes the currently SELECTed Redis DB. */
|
||||
void flushdbCommand(client *c) {
|
||||
int flags;
|
||||
|
||||
if (getFlushCommandFlags(c,&flags) == C_ERR) return;
|
||||
server.dirty += emptyDb(c->db->id,flags,NULL);
|
||||
addReply(c,shared.ok);
|
||||
#if defined(USE_JEMALLOC)
|
||||
/* jemalloc 5 doesn't release pages back to the OS when there's no traffic.
|
||||
* for large databases, flushdb blocks for long anyway, so a bit more won't
|
||||
* harm and this way the flush and purge will be synchroneus. */
|
||||
if (!(flags & EMPTYDB_ASYNC))
|
||||
jemalloc_purge();
|
||||
#endif
|
||||
}
|
||||
|
||||
/* FLUSHALL [ASYNC]
|
||||
*
|
||||
* Flushes the whole server data set. */
|
||||
void flushallCommand(client *c) {
|
||||
int flags;
|
||||
if (getFlushCommandFlags(c,&flags) == C_ERR) return;
|
||||
flushAllDataAndResetRDB(flags);
|
||||
addReply(c,shared.ok);
|
||||
}
|
||||
|
||||
/* This command implements DEL and LAZYDEL. */
|
||||
@@ -476,7 +570,7 @@ void delGenericCommand(client *c, int lazy) {
|
||||
int deleted = lazy ? dbAsyncDelete(c->db,c->argv[j]) :
|
||||
dbSyncDelete(c->db,c->argv[j]);
|
||||
if (deleted) {
|
||||
signalModifiedKey(c->db,c->argv[j]);
|
||||
signalModifiedKey(c,c->db,c->argv[j]);
|
||||
notifyKeyspaceEvent(NOTIFY_GENERIC,
|
||||
"del",c->argv[j],c->db->id);
|
||||
server.dirty++;
|
||||
@@ -487,7 +581,7 @@ void delGenericCommand(client *c, int lazy) {
|
||||
}
|
||||
|
||||
void delCommand(client *c) {
|
||||
delGenericCommand(c,0);
|
||||
delGenericCommand(c,server.lazyfree_lazy_user_del);
|
||||
}
|
||||
|
||||
void unlinkCommand(client *c) {
|
||||
@@ -528,7 +622,7 @@ void randomkeyCommand(client *c) {
|
||||
robj *key;
|
||||
|
||||
if ((key = dbRandomKey(c->db)) == NULL) {
|
||||
addReply(c,shared.nullbulk);
|
||||
addReplyNull(c);
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -542,7 +636,7 @@ void keysCommand(client *c) {
|
||||
sds pattern = c->argv[1]->ptr;
|
||||
int plen = (int)sdslen(pattern), allkeys;
|
||||
PORT_ULONG numkeys = 0;
|
||||
void *replylen = addDeferredMultiBulkLength(c);
|
||||
void *replylen = addReplyDeferredLen(c);
|
||||
|
||||
di = dictGetSafeIterator(c->db->dict);
|
||||
allkeys = (pattern[0] == '*' && plen == 1);
|
||||
@@ -560,7 +654,7 @@ void keysCommand(client *c) {
|
||||
}
|
||||
}
|
||||
dictReleaseIterator(di);
|
||||
setDeferredMultiBulkLength(c,replylen,numkeys);
|
||||
setDeferredArrayLen(c,replylen,numkeys);
|
||||
}
|
||||
|
||||
/* This callback is used by scanGenericCommand in order to collect elements
|
||||
@@ -614,7 +708,7 @@ int parseScanCursorOrReply(client *c, robj *o, PORT_ULONG *cursor) {
|
||||
}
|
||||
|
||||
/* This command implements SCAN, HSCAN and SSCAN commands.
|
||||
* If object 'o' is passed, then it must be a Hash or Set object, otherwise
|
||||
* If object 'o' is passed, then it must be a Hash, Set or Zset object, otherwise
|
||||
* if 'o' is NULL the command will operate on the dictionary associated with
|
||||
* the current database.
|
||||
*
|
||||
@@ -630,6 +724,7 @@ void scanGenericCommand(client *c, robj *o, PORT_ULONG cursor) {
|
||||
listNode *node, *nextnode;
|
||||
PORT_LONG count = 10;
|
||||
sds pat = NULL;
|
||||
sds typename = NULL;
|
||||
int patlen = 0, use_pattern = 0;
|
||||
dict *ht;
|
||||
|
||||
@@ -666,6 +761,10 @@ void scanGenericCommand(client *c, robj *o, PORT_ULONG cursor) {
|
||||
use_pattern = !(pat[0] == '*' && patlen == 1);
|
||||
|
||||
i += 2;
|
||||
} else if (!strcasecmp(c->argv[i]->ptr, "type") && o == NULL && j >= 2) {
|
||||
/* SCAN for a particular type only applies to the db dict */
|
||||
typename = c->argv[i+1]->ptr;
|
||||
i+= 2;
|
||||
} else {
|
||||
addReply(c,shared.syntaxerr);
|
||||
goto cleanup;
|
||||
@@ -760,10 +859,17 @@ void scanGenericCommand(client *c, robj *o, PORT_ULONG cursor) {
|
||||
}
|
||||
}
|
||||
|
||||
/* Filter an element if it isn't the type we want. */
|
||||
if (!filter && o == NULL && typename){
|
||||
robj* typecheck = lookupKeyReadWithFlags(c->db, kobj, LOOKUP_NOTOUCH);
|
||||
char* type = getObjectTypeName(typecheck);
|
||||
if (strcasecmp((char*) typename, type)) filter = 1;
|
||||
}
|
||||
|
||||
/* Filter element if it is an expired key. */
|
||||
if (!filter && o == NULL && expireIfNeeded(c->db, kobj)) filter = 1;
|
||||
|
||||
/* Remove the element and its associted value if needed. */
|
||||
/* Remove the element and its associated value if needed. */
|
||||
if (filter) {
|
||||
decrRefCount(kobj);
|
||||
listDelNode(keys, node);
|
||||
@@ -785,10 +891,10 @@ void scanGenericCommand(client *c, robj *o, PORT_ULONG cursor) {
|
||||
}
|
||||
|
||||
/* Step 4: Reply to the client. */
|
||||
addReplyMultiBulkLen(c, 2);
|
||||
addReplyArrayLen(c, 2);
|
||||
addReplyBulkLongLong(c,cursor);
|
||||
|
||||
addReplyMultiBulkLen(c, listLength(keys));
|
||||
addReplyArrayLen(c, listLength(keys));
|
||||
while ((node = listFirst(keys)) != NULL) {
|
||||
robj *kobj = listNodeValue(node);
|
||||
addReplyBulk(c, kobj);
|
||||
@@ -816,11 +922,8 @@ void lastsaveCommand(client *c) {
|
||||
addReplyLongLong(c,server.lastsave);
|
||||
}
|
||||
|
||||
void typeCommand(client *c) {
|
||||
robj *o;
|
||||
char *type;
|
||||
|
||||
o = lookupKeyReadWithFlags(c->db,c->argv[1],LOOKUP_NOTOUCH);
|
||||
char* getObjectTypeName(robj *o) {
|
||||
char* type;
|
||||
if (o == NULL) {
|
||||
type = "none";
|
||||
} else {
|
||||
@@ -838,7 +941,13 @@ void typeCommand(client *c) {
|
||||
default: type = "unknown"; break;
|
||||
}
|
||||
}
|
||||
addReplyStatus(c,type);
|
||||
return type;
|
||||
}
|
||||
|
||||
void typeCommand(client *c) {
|
||||
robj *o;
|
||||
o = lookupKeyReadWithFlags(c->db,c->argv[1],LOOKUP_NOTOUCH);
|
||||
addReplyStatus(c, getObjectTypeName(o));
|
||||
}
|
||||
|
||||
void shutdownCommand(client *c) {
|
||||
@@ -893,8 +1002,8 @@ void renameGenericCommand(client *c, int nx) {
|
||||
dbAdd(c->db,c->argv[2],o);
|
||||
if (expire != -1) setExpire(c,c->db,c->argv[2],expire);
|
||||
dbDelete(c->db,c->argv[1]);
|
||||
signalModifiedKey(c->db,c->argv[1]);
|
||||
signalModifiedKey(c->db,c->argv[2]);
|
||||
signalModifiedKey(c,c->db,c->argv[1]);
|
||||
signalModifiedKey(c,c->db,c->argv[2]);
|
||||
notifyKeyspaceEvent(NOTIFY_GENERIC,"rename_from",
|
||||
c->argv[1],c->db->id);
|
||||
notifyKeyspaceEvent(NOTIFY_GENERIC,"rename_to",
|
||||
@@ -962,6 +1071,13 @@ void moveCommand(client *c) {
|
||||
|
||||
/* OK! key moved, free the entry in the source DB */
|
||||
dbDelete(src,c->argv[1]);
|
||||
signalModifiedKey(c,src,c->argv[1]);
|
||||
signalModifiedKey(c,dst,c->argv[1]);
|
||||
notifyKeyspaceEvent(NOTIFY_GENERIC,
|
||||
"move_from",c->argv[1],src->id);
|
||||
notifyKeyspaceEvent(NOTIFY_GENERIC,
|
||||
"move_to",c->argv[1],dst->id);
|
||||
|
||||
server.dirty++;
|
||||
addReply(c,shared.cone);
|
||||
}
|
||||
@@ -992,7 +1108,7 @@ void scanDatabaseForReadyLists(redisDb *db) {
|
||||
*
|
||||
* Returns C_ERR if at least one of the DB ids are out of range, otherwise
|
||||
* C_OK is returned. */
|
||||
int dbSwapDatabases(int id1, int id2) {
|
||||
int dbSwapDatabases(PORT_LONG id1, PORT_LONG id2) {
|
||||
if (id1 < 0 || id1 >= server.dbnum ||
|
||||
id2 < 0 || id2 >= server.dbnum) return C_ERR;
|
||||
if (id1 == id2) return C_OK;
|
||||
@@ -1005,10 +1121,12 @@ int dbSwapDatabases(int id1, int id2) {
|
||||
db1->dict = db2->dict;
|
||||
db1->expires = db2->expires;
|
||||
db1->avg_ttl = db2->avg_ttl;
|
||||
db1->expires_cursor = db2->expires_cursor;
|
||||
|
||||
db2->dict = aux.dict;
|
||||
db2->expires = aux.expires;
|
||||
db2->avg_ttl = aux.avg_ttl;
|
||||
db2->expires_cursor = aux.expires_cursor;
|
||||
|
||||
/* Now we need to handle clients blocked on lists: as an effect
|
||||
* of swapping the two DBs, a client that was waiting for list
|
||||
@@ -1048,6 +1166,8 @@ void swapdbCommand(client *c) {
|
||||
addReplyError(c,"DB index is out of range");
|
||||
return;
|
||||
} else {
|
||||
RedisModuleSwapDbInfo si = {REDISMODULE_SWAPDBINFO_VERSION,id1,id2};
|
||||
moduleFireServerEvent(REDISMODULE_EVENT_SWAPDB,0,&si);
|
||||
server.dirty++;
|
||||
addReply(c,shared.ok);
|
||||
}
|
||||
@@ -1196,28 +1316,64 @@ int expireIfNeeded(redisDb *db, robj *key) {
|
||||
propagateExpire(db,key,server.lazyfree_lazy_expire);
|
||||
notifyKeyspaceEvent(NOTIFY_EXPIRED,
|
||||
"expired",key,db->id);
|
||||
return server.lazyfree_lazy_expire ? dbAsyncDelete(db,key) :
|
||||
dbSyncDelete(db,key);
|
||||
int retval = server.lazyfree_lazy_expire ? dbAsyncDelete(db,key) :
|
||||
dbSyncDelete(db,key);
|
||||
if (retval) signalModifiedKey(NULL,db,key);
|
||||
return retval;
|
||||
}
|
||||
|
||||
/* -----------------------------------------------------------------------------
|
||||
* API to get key arguments from commands
|
||||
* ---------------------------------------------------------------------------*/
|
||||
|
||||
/* Prepare the getKeysResult struct to hold numkeys, either by using the
|
||||
* pre-allocated keysbuf or by allocating a new array on the heap.
|
||||
*
|
||||
* This function must be called at least once before starting to populate
|
||||
* the result, and can be called repeatedly to enlarge the result array.
|
||||
*/
|
||||
int *getKeysPrepareResult(getKeysResult *result, int numkeys) {
|
||||
/* GETKEYS_RESULT_INIT initializes keys to NULL, point it to the pre-allocated stack
|
||||
* buffer here. */
|
||||
if (!result->keys) {
|
||||
serverAssert(!result->numkeys);
|
||||
result->keys = result->keysbuf;
|
||||
}
|
||||
|
||||
/* Resize if necessary */
|
||||
if (numkeys > result->size) {
|
||||
if (result->keys != result->keysbuf) {
|
||||
/* We're not using a static buffer, just (re)alloc */
|
||||
result->keys = zrealloc(result->keys, numkeys * sizeof(int));
|
||||
} else {
|
||||
/* We are using a static buffer, copy its contents */
|
||||
result->keys = zmalloc(numkeys * sizeof(int));
|
||||
if (result->numkeys)
|
||||
memcpy(result->keys, result->keysbuf, result->numkeys * sizeof(int));
|
||||
}
|
||||
result->size = numkeys;
|
||||
}
|
||||
|
||||
return result->keys;
|
||||
}
|
||||
|
||||
/* The base case is to use the keys position as given in the command table
|
||||
* (firstkey, lastkey, step). */
|
||||
int *getKeysUsingCommandTable(struct redisCommand *cmd,robj **argv, int argc, int *numkeys) {
|
||||
int getKeysUsingCommandTable(struct redisCommand *cmd,robj **argv, int argc, getKeysResult *result) {
|
||||
int j, i = 0, last, *keys;
|
||||
UNUSED(argv);
|
||||
|
||||
if (cmd->firstkey == 0) {
|
||||
*numkeys = 0;
|
||||
return NULL;
|
||||
result->numkeys = 0;
|
||||
return 0;
|
||||
}
|
||||
|
||||
last = cmd->lastkey;
|
||||
if (last < 0) last = argc+last;
|
||||
keys = zmalloc(sizeof(int)*(((PORT_ULONG) last - cmd->firstkey)+1)); WIN_PORT_FIX /* cat (PORT_ULONG) */
|
||||
|
||||
int count = ((last - cmd->firstkey)+1);
|
||||
keys = getKeysPrepareResult(result, count);
|
||||
|
||||
for (j = cmd->firstkey; j <= last; j += cmd->keystep) {
|
||||
if (j >= argc) {
|
||||
/* Modules commands, and standard commands with a not fixed number
|
||||
@@ -1227,23 +1383,23 @@ int *getKeysUsingCommandTable(struct redisCommand *cmd,robj **argv, int argc, in
|
||||
* return no keys and expect the command implementation to report
|
||||
* an arity or syntax error. */
|
||||
if (cmd->flags & CMD_MODULE || cmd->arity < 0) {
|
||||
zfree(keys);
|
||||
*numkeys = 0;
|
||||
return NULL;
|
||||
getKeysFreeResult(result);
|
||||
result->numkeys = 0;
|
||||
return 0;
|
||||
} else {
|
||||
serverPanic("Redis built-in command declared keys positions not matching the arity requirements.");
|
||||
}
|
||||
}
|
||||
keys[i++] = j;
|
||||
}
|
||||
*numkeys = i;
|
||||
return keys;
|
||||
result->numkeys = i;
|
||||
return i;
|
||||
}
|
||||
|
||||
/* Return all the arguments that are keys in the command passed via argc / argv.
|
||||
*
|
||||
* The command returns the positions of all the key arguments inside the array,
|
||||
* so the actual return value is an heap allocated array of integers. The
|
||||
* so the actual return value is a heap allocated array of integers. The
|
||||
* length of the array is returned by reference into *numkeys.
|
||||
*
|
||||
* 'cmd' must be point to the corresponding entry into the redisCommand
|
||||
@@ -1251,25 +1407,26 @@ int *getKeysUsingCommandTable(struct redisCommand *cmd,robj **argv, int argc, in
|
||||
*
|
||||
* This function uses the command table if a command-specific helper function
|
||||
* is not required, otherwise it calls the command-specific function. */
|
||||
int *getKeysFromCommand(struct redisCommand *cmd, robj **argv, int argc, int *numkeys) {
|
||||
int getKeysFromCommand(struct redisCommand *cmd, robj **argv, int argc, getKeysResult *result) {
|
||||
if (cmd->flags & CMD_MODULE_GETKEYS) {
|
||||
return moduleGetCommandKeysViaAPI(cmd,argv,argc,numkeys);
|
||||
return moduleGetCommandKeysViaAPI(cmd,argv,argc,result);
|
||||
} else if (!(cmd->flags & CMD_MODULE) && cmd->getkeys_proc) {
|
||||
return cmd->getkeys_proc(cmd,argv,argc,numkeys);
|
||||
return cmd->getkeys_proc(cmd,argv,argc,result);
|
||||
} else {
|
||||
return getKeysUsingCommandTable(cmd,argv,argc,numkeys);
|
||||
return getKeysUsingCommandTable(cmd,argv,argc,result);
|
||||
}
|
||||
}
|
||||
|
||||
/* Free the result of getKeysFromCommand. */
|
||||
void getKeysFreeResult(int *result) {
|
||||
zfree(result);
|
||||
void getKeysFreeResult(getKeysResult *result) {
|
||||
if (result && result->keys != result->keysbuf)
|
||||
zfree(result->keys);
|
||||
}
|
||||
|
||||
/* Helper function to extract keys from following commands:
|
||||
* ZUNIONSTORE <destkey> <num-keys> <key> <key> ... <key> <options>
|
||||
* ZINTERSTORE <destkey> <num-keys> <key> <key> ... <key> <options> */
|
||||
int *zunionInterGetKeys(struct redisCommand *cmd, robj **argv, int argc, int *numkeys) {
|
||||
int zunionInterGetKeys(struct redisCommand *cmd, robj **argv, int argc, getKeysResult *result) {
|
||||
int i, num, *keys;
|
||||
UNUSED(cmd);
|
||||
|
||||
@@ -1277,28 +1434,30 @@ int *zunionInterGetKeys(struct redisCommand *cmd, robj **argv, int argc, int *nu
|
||||
/* Sanity check. Don't return any key if the command is going to
|
||||
* reply with syntax error. */
|
||||
if (num < 1 || num > (argc-3)) {
|
||||
*numkeys = 0;
|
||||
return NULL;
|
||||
result->numkeys = 0;
|
||||
return 0;
|
||||
}
|
||||
|
||||
/* Keys in z{union,inter}store come from two places:
|
||||
* argv[1] = storage key,
|
||||
* argv[3...n] = keys to intersect */
|
||||
keys = zmalloc(sizeof(int)*((PORT_ULONG)num+1)); WIN_PORT_FIX /* cast (PORT_ULONG) */
|
||||
/* Total keys = {union,inter} keys + storage key */
|
||||
keys = getKeysPrepareResult(result, num+1);
|
||||
result->numkeys = num+1;
|
||||
|
||||
/* Add all key positions for argv[3...n] to keys[] */
|
||||
for (i = 0; i < num; i++) keys[i] = 3+i;
|
||||
|
||||
/* Finally add the argv[1] key position (the storage key target). */
|
||||
keys[num] = 1;
|
||||
*numkeys = num+1; /* Total keys = {union,inter} keys + storage key */
|
||||
return keys;
|
||||
|
||||
return result->numkeys;
|
||||
}
|
||||
|
||||
/* Helper function to extract keys from the following commands:
|
||||
* EVAL <script> <num-keys> <key> <key> ... <key> [more stuff]
|
||||
* EVALSHA <script> <num-keys> <key> <key> ... <key> [more stuff] */
|
||||
int *evalGetKeys(struct redisCommand *cmd, robj **argv, int argc, int *numkeys) {
|
||||
int evalGetKeys(struct redisCommand *cmd, robj **argv, int argc, getKeysResult *result) {
|
||||
int i, num, *keys;
|
||||
UNUSED(cmd);
|
||||
|
||||
@@ -1306,17 +1465,17 @@ int *evalGetKeys(struct redisCommand *cmd, robj **argv, int argc, int *numkeys)
|
||||
/* Sanity check. Don't return any key if the command is going to
|
||||
* reply with syntax error. */
|
||||
if (num <= 0 || num > (argc-3)) {
|
||||
*numkeys = 0;
|
||||
return NULL;
|
||||
result->numkeys = 0;
|
||||
return 0;
|
||||
}
|
||||
|
||||
keys = zmalloc(sizeof(int)*num);
|
||||
*numkeys = num;
|
||||
keys = getKeysPrepareResult(result, num);
|
||||
result->numkeys = num;
|
||||
|
||||
/* Add all key positions for argv[3...n] to keys[] */
|
||||
for (i = 0; i < num; i++) keys[i] = 3+i;
|
||||
|
||||
return keys;
|
||||
return result->numkeys;
|
||||
}
|
||||
|
||||
/* Helper function to extract keys from the SORT command.
|
||||
@@ -1326,13 +1485,12 @@ int *evalGetKeys(struct redisCommand *cmd, robj **argv, int argc, int *numkeys)
|
||||
* The first argument of SORT is always a key, however a list of options
|
||||
* follow in SQL-alike style. Here we parse just the minimum in order to
|
||||
* correctly identify keys in the "STORE" option. */
|
||||
int *sortGetKeys(struct redisCommand *cmd, robj **argv, int argc, int *numkeys) {
|
||||
int sortGetKeys(struct redisCommand *cmd, robj **argv, int argc, getKeysResult *result) {
|
||||
int i, j, num, *keys, found_store = 0;
|
||||
UNUSED(cmd);
|
||||
|
||||
num = 0;
|
||||
keys = zmalloc(sizeof(int)*2); /* Alloc 2 places for the worst case. */
|
||||
|
||||
keys = getKeysPrepareResult(result, 2); /* Alloc 2 places for the worst case. */
|
||||
keys[num++] = 1; /* <sort-key> is always present. */
|
||||
|
||||
/* Search for STORE option. By default we consider options to don't
|
||||
@@ -1364,11 +1522,11 @@ int *sortGetKeys(struct redisCommand *cmd, robj **argv, int argc, int *numkeys)
|
||||
}
|
||||
}
|
||||
}
|
||||
*numkeys = num + found_store;
|
||||
return keys;
|
||||
result->numkeys = num + found_store;
|
||||
return result->numkeys;
|
||||
}
|
||||
|
||||
int *migrateGetKeys(struct redisCommand *cmd, robj **argv, int argc, int *numkeys) {
|
||||
int migrateGetKeys(struct redisCommand *cmd, robj **argv, int argc, getKeysResult *result) {
|
||||
int i, num, first, *keys;
|
||||
UNUSED(cmd);
|
||||
|
||||
@@ -1389,17 +1547,17 @@ int *migrateGetKeys(struct redisCommand *cmd, robj **argv, int argc, int *numkey
|
||||
}
|
||||
}
|
||||
|
||||
keys = zmalloc(sizeof(int)*num);
|
||||
keys = getKeysPrepareResult(result, num);
|
||||
for (i = 0; i < num; i++) keys[i] = first+i;
|
||||
*numkeys = num;
|
||||
return keys;
|
||||
result->numkeys = num;
|
||||
return num;
|
||||
}
|
||||
|
||||
/* Helper function to extract keys from following commands:
|
||||
* GEORADIUS key x y radius unit [WITHDIST] [WITHHASH] [WITHCOORD] [ASC|DESC]
|
||||
* [COUNT count] [STORE key] [STOREDIST key]
|
||||
* GEORADIUSBYMEMBER key member radius unit ... options ... */
|
||||
int *georadiusGetKeys(struct redisCommand *cmd, robj **argv, int argc, int *numkeys) {
|
||||
int georadiusGetKeys(struct redisCommand *cmd, robj **argv, int argc, getKeysResult *result) {
|
||||
int i, num, *keys;
|
||||
UNUSED(cmd);
|
||||
|
||||
@@ -1422,20 +1580,60 @@ int *georadiusGetKeys(struct redisCommand *cmd, robj **argv, int argc, int *numk
|
||||
* argv[1] = key,
|
||||
* argv[5...n] = stored key if present
|
||||
*/
|
||||
keys = zmalloc(sizeof(int) * num);
|
||||
keys = getKeysPrepareResult(result, num);
|
||||
|
||||
/* Add all key positions to keys[] */
|
||||
keys[0] = 1;
|
||||
if(num > 1) {
|
||||
keys[1] = stored_key;
|
||||
}
|
||||
*numkeys = num;
|
||||
return keys;
|
||||
result->numkeys = num;
|
||||
return num;
|
||||
}
|
||||
|
||||
/* LCS ... [KEYS <key1> <key2>] ... */
|
||||
int lcsGetKeys(struct redisCommand *cmd, robj **argv, int argc, getKeysResult *result) {
|
||||
int i;
|
||||
int *keys = getKeysPrepareResult(result, 2);
|
||||
UNUSED(cmd);
|
||||
|
||||
/* We need to parse the options of the command in order to check for the
|
||||
* "KEYS" argument before the "STRINGS" argument. */
|
||||
for (i = 1; i < argc; i++) {
|
||||
char *arg = argv[i]->ptr;
|
||||
int moreargs = (argc-1) - i;
|
||||
|
||||
if (!strcasecmp(arg, "strings")) {
|
||||
break;
|
||||
} else if (!strcasecmp(arg, "keys") && moreargs >= 2) {
|
||||
keys[0] = i+1;
|
||||
keys[1] = i+2;
|
||||
result->numkeys = 2;
|
||||
return result->numkeys;
|
||||
}
|
||||
}
|
||||
result->numkeys = 0;
|
||||
return result->numkeys;
|
||||
}
|
||||
|
||||
/* Helper function to extract keys from memory command.
|
||||
* MEMORY USAGE <key> */
|
||||
int memoryGetKeys(struct redisCommand *cmd, robj **argv, int argc, getKeysResult *result) {
|
||||
UNUSED(cmd);
|
||||
|
||||
getKeysPrepareResult(result, 1);
|
||||
if (argc >= 3 && !strcasecmp(argv[1]->ptr,"usage")) {
|
||||
result->keys[0] = 2;
|
||||
result->numkeys = 1;
|
||||
return result->numkeys;
|
||||
}
|
||||
result->numkeys = 0;
|
||||
return 0;
|
||||
}
|
||||
|
||||
/* XREAD [BLOCK <milliseconds>] [COUNT <count>] [GROUP <groupname> <ttl>]
|
||||
* STREAMS key_1 key_2 ... key_N ID_1 ID_2 ... ID_N */
|
||||
int *xreadGetKeys(struct redisCommand *cmd, robj **argv, int argc, int *numkeys) {
|
||||
int xreadGetKeys(struct redisCommand *cmd, robj **argv, int argc, getKeysResult *result) {
|
||||
int i, num = 0, *keys;
|
||||
UNUSED(cmd);
|
||||
|
||||
@@ -1465,33 +1663,33 @@ int *xreadGetKeys(struct redisCommand *cmd, robj **argv, int argc, int *numkeys)
|
||||
|
||||
/* Syntax error. */
|
||||
if (streams_pos == -1 || num == 0 || num % 2 != 0) {
|
||||
*numkeys = 0;
|
||||
return NULL;
|
||||
result->numkeys = 0;
|
||||
return 0;
|
||||
}
|
||||
num /= 2; /* We have half the keys as there are arguments because
|
||||
there are also the IDs, one per key. */
|
||||
|
||||
keys = zmalloc(sizeof(int) * num);
|
||||
keys = getKeysPrepareResult(result, num);
|
||||
for (i = streams_pos+1; i < argc-num; i++) keys[i-streams_pos-1] = i;
|
||||
*numkeys = num;
|
||||
return keys;
|
||||
result->numkeys = num;
|
||||
return num;
|
||||
}
|
||||
|
||||
/* Slot to Key API. This is used by Redis Cluster in order to obtain in
|
||||
* a fast way a key that belongs to a specified hash slot. This is useful
|
||||
* while rehashing the cluster and in other conditions when we need to
|
||||
* understand if we have keys for a given hash slot. */
|
||||
void slotToKeyUpdateKey(robj *key, int add) {
|
||||
unsigned int hashslot = keyHashSlot(key->ptr,sdslen(key->ptr));
|
||||
void slotToKeyUpdateKey(sds key, int add) {
|
||||
size_t keylen = sdslen(key);
|
||||
unsigned int hashslot = keyHashSlot(key,keylen);
|
||||
unsigned char buf[64];
|
||||
unsigned char *indexed = buf;
|
||||
size_t keylen = sdslen(key->ptr);
|
||||
|
||||
server.cluster->slots_keys_count[hashslot] += add ? 1 : -1;
|
||||
if (keylen+2 > 64) indexed = zmalloc(keylen+2);
|
||||
indexed[0] = (hashslot >> 8) & 0xff;
|
||||
indexed[1] = hashslot & 0xff;
|
||||
memcpy(indexed+2,key->ptr,keylen);
|
||||
memcpy(indexed+2,key,keylen);
|
||||
if (add) {
|
||||
raxInsert(server.cluster->slots_to_keys,indexed,keylen+2,NULL,NULL);
|
||||
} else {
|
||||
@@ -1500,11 +1698,11 @@ void slotToKeyUpdateKey(robj *key, int add) {
|
||||
if (indexed != buf) zfree(indexed);
|
||||
}
|
||||
|
||||
void slotToKeyAdd(robj *key) {
|
||||
void slotToKeyAdd(sds key) {
|
||||
slotToKeyUpdateKey(key,1);
|
||||
}
|
||||
|
||||
void slotToKeyDel(robj *key) {
|
||||
void slotToKeyDel(sds key) {
|
||||
slotToKeyUpdateKey(key,0);
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user