Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
25 commits
Select commit Hold shift + click to select a range
827a5ee
Implement XACKDEL & Address Code Reviews
Sep 8, 2026
f93c54c
Implement XDELEX & Address Code Reviews
Sep 8, 2026
9557388
Shared helpers for XDEL/XDELEX/XACKDEL
Sep 9, 2026
45effc1
Address CodeRabbit comments
Sep 9, 2026
71e9dbd
Address CodeRabbit Comments on XDELEX ACKED mode gap
Sep 9, 2026
295a4e1
Propagate XACKDEL/XDELEX manually with pre-9.2 commands
Sep 9, 2026
5d2bd88
Create common parser for `XDEL`-like commands
Sep 9, 2026
9ea19ed
Align XDEL-like parse helper to others
Sep 9, 2026
9fdb69a
Revert unrelated testing change
Sep 9, 2026
98c4c30
Remove clunky XDELEX/XACKDEL helpers
Sep 9, 2026
71669bb
Clang Format
Sep 9, 2026
b814bcf
Add streamDeletePELEntry helper across XACKDEL, XDELEX, XACK
Sep 9, 2026
37e6690
Mirror fix for XDELEX ACKED mode gap in XACKDEL as well
Sep 9, 2026
c029281
Use size_t for id count to clarify constraints
Sep 9, 2026
96128f8
Unify XDEL, XDELEX and XACKDEL into a single xdelGenericCommand
zuiderkwast Sep 10, 2026
3e67a03
Inline single-caller helper streamParseStrictIDsOrReply
zuiderkwast Sep 10, 2026
03ff02d
Clang format
zuiderkwast Sep 10, 2026
3882fdc
Only allocate working arrays each xdelGenericCommand variant needs
zuiderkwast Sep 10, 2026
11ed5b7
Clarify Complexity in XDELEX JSON Defn
nickiaq Sep 10, 2026
0766138
Clarify Complexity in XACKDEL JSON Defn
nickiaq Sep 10, 2026
75b1be9
Revert formatting change to unrelated line
Sep 10, 2026
7a0addb
Revert formatting change to unrelated line
Sep 10, 2026
51d6e79
Regenerate commands.def after accepting JSON changes in GitHub
Sep 10, 2026
704fca7
Add Cross Version (pre-9.2.0) Replication Tests
Sep 10, 2026
346b74f
Merge branch 'unstable' into commands-xdelex-xackdel
nickiaq Sep 10, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
73 changes: 73 additions & 0 deletions src/commands.def
Original file line number Diff line number Diff line change
Expand Up @@ -10452,6 +10452,42 @@ struct COMMAND_ARG XACK_Args[] = {
{MAKE_ARG("id",ARG_TYPE_STRING,-1,NULL,NULL,NULL,CMD_ARG_MULTIPLE,0,NULL)},
};

/********** XACKDEL ********************/

#ifndef SKIP_CMD_HISTORY_TABLE
/* XACKDEL history */
#define XACKDEL_History NULL
#endif

#ifndef SKIP_CMD_TIPS_TABLE
/* XACKDEL tips */
#define XACKDEL_Tips NULL
#endif

#ifndef SKIP_CMD_KEY_SPECS_TABLE
/* XACKDEL key specs */
keySpec XACKDEL_Keyspecs[1] = {
{NULL,CMD_KEY_RW|CMD_KEY_UPDATE,KSPEC_BS_INDEX,.bs.index={1},KSPEC_FK_RANGE,.fk.range={0,1,0}}
};
#endif

/* XACKDEL mode argument table */
struct COMMAND_ARG XACKDEL_mode_Subargs[] = {
{MAKE_ARG("keepref",ARG_TYPE_PURE_TOKEN,-1,"KEEPREF",NULL,NULL,CMD_ARG_NONE,0,NULL)},
{MAKE_ARG("delref",ARG_TYPE_PURE_TOKEN,-1,"DELREF",NULL,NULL,CMD_ARG_NONE,0,NULL)},
{MAKE_ARG("acked",ARG_TYPE_PURE_TOKEN,-1,"ACKED",NULL,NULL,CMD_ARG_NONE,0,NULL)},
};

/* XACKDEL argument table */
struct COMMAND_ARG XACKDEL_Args[] = {
{MAKE_ARG("key",ARG_TYPE_KEY,0,NULL,NULL,NULL,CMD_ARG_NONE,0,NULL)},
{MAKE_ARG("group",ARG_TYPE_STRING,-1,NULL,NULL,NULL,CMD_ARG_NONE,0,NULL)},
{MAKE_ARG("mode",ARG_TYPE_ONEOF,-1,NULL,NULL,NULL,CMD_ARG_OPTIONAL,3,NULL),.subargs=XACKDEL_mode_Subargs},
{MAKE_ARG("ids",ARG_TYPE_PURE_TOKEN,-1,"IDS",NULL,NULL,CMD_ARG_NONE,0,NULL)},
{MAKE_ARG("numids",ARG_TYPE_INTEGER,-1,NULL,NULL,NULL,CMD_ARG_NONE,0,NULL)},
{MAKE_ARG("id",ARG_TYPE_STRING,-1,NULL,NULL,NULL,CMD_ARG_MULTIPLE,0,NULL)},
};

/********** XADD ********************/

#ifndef SKIP_CMD_HISTORY_TABLE
Expand Down Expand Up @@ -10612,6 +10648,41 @@ struct COMMAND_ARG XDEL_Args[] = {
{MAKE_ARG("id",ARG_TYPE_STRING,-1,NULL,NULL,NULL,CMD_ARG_MULTIPLE,0,NULL)},
};

/********** XDELEX ********************/

#ifndef SKIP_CMD_HISTORY_TABLE
/* XDELEX history */
#define XDELEX_History NULL
#endif

#ifndef SKIP_CMD_TIPS_TABLE
/* XDELEX tips */
#define XDELEX_Tips NULL
#endif

#ifndef SKIP_CMD_KEY_SPECS_TABLE
/* XDELEX key specs */
keySpec XDELEX_Keyspecs[1] = {
{NULL,CMD_KEY_RW|CMD_KEY_UPDATE,KSPEC_BS_INDEX,.bs.index={1},KSPEC_FK_RANGE,.fk.range={0,1,0}}
};
#endif

/* XDELEX mode argument table */
struct COMMAND_ARG XDELEX_mode_Subargs[] = {
{MAKE_ARG("keepref",ARG_TYPE_PURE_TOKEN,-1,"KEEPREF",NULL,NULL,CMD_ARG_NONE,0,NULL)},
{MAKE_ARG("delref",ARG_TYPE_PURE_TOKEN,-1,"DELREF",NULL,NULL,CMD_ARG_NONE,0,NULL)},
{MAKE_ARG("acked",ARG_TYPE_PURE_TOKEN,-1,"ACKED",NULL,NULL,CMD_ARG_NONE,0,NULL)},
};

/* XDELEX argument table */
struct COMMAND_ARG XDELEX_Args[] = {
{MAKE_ARG("key",ARG_TYPE_KEY,0,NULL,NULL,NULL,CMD_ARG_NONE,0,NULL)},
{MAKE_ARG("mode",ARG_TYPE_ONEOF,-1,NULL,NULL,NULL,CMD_ARG_OPTIONAL,3,NULL),.subargs=XDELEX_mode_Subargs},
{MAKE_ARG("ids",ARG_TYPE_PURE_TOKEN,-1,"IDS",NULL,NULL,CMD_ARG_NONE,0,NULL)},
{MAKE_ARG("numids",ARG_TYPE_INTEGER,-1,NULL,NULL,NULL,CMD_ARG_NONE,0,NULL)},
{MAKE_ARG("id",ARG_TYPE_STRING,-1,NULL,NULL,NULL,CMD_ARG_MULTIPLE,0,NULL)},
};

/********** XGROUP CREATE ********************/

#ifndef SKIP_CMD_HISTORY_TABLE
Expand Down Expand Up @@ -12214,10 +12285,12 @@ struct COMMAND_STRUCT serverCommandTable[] = {
{MAKE_CMD("zunionstore","Stores the union of multiple sorted sets in a key.","O(N)+O(M log(M)) with N being the sum of the sizes of the input sorted sets, and M being the number of elements in the resulting sorted set.","2.0.0",CMD_DOC_NONE,NULL,NULL,"sorted_set",COMMAND_GROUP_SORTED_SET,ZUNIONSTORE_History,0,ZUNIONSTORE_Tips,0,zunionstoreCommand,-4,CMD_WRITE|CMD_DENYOOM,ACL_CATEGORY_SLOW|ACL_CATEGORY_SORTEDSET|ACL_CATEGORY_WRITE,NULL,ZUNIONSTORE_Keyspecs,2,zunionInterDiffStoreGetKeys,5),.args=ZUNIONSTORE_Args},
/* stream */
{MAKE_CMD("xack","Returns the number of messages that were successfully acknowledged by the consumer group member of a stream.","O(1) for each message ID processed.","5.0.0",CMD_DOC_NONE,NULL,NULL,"stream",COMMAND_GROUP_STREAM,XACK_History,0,XACK_Tips,0,xackCommand,-4,CMD_WRITE|CMD_FAST,ACL_CATEGORY_FAST|ACL_CATEGORY_STREAM|ACL_CATEGORY_WRITE,NULL,XACK_Keyspecs,1,NULL,3),.args=XACK_Args},
{MAKE_CMD("xackdel","Acknowledge and (if possible) delete stream message(s).","O(1) for each single item to delete in the stream, regardless of the stream size","9.2.0",CMD_DOC_NONE,NULL,NULL,"stream",COMMAND_GROUP_STREAM,XACKDEL_History,0,XACKDEL_Tips,0,xackdelCommand,-6,CMD_WRITE|CMD_FAST,ACL_CATEGORY_FAST|ACL_CATEGORY_WRITE|ACL_CATEGORY_STREAM,NULL,XACKDEL_Keyspecs,1,NULL,6),.args=XACKDEL_Args},
{MAKE_CMD("xadd","Appends a new message to a stream. Creates the key if it doesn't exist.","O(1) when adding a new entry, O(N) when trimming where N being the number of entries evicted.","5.0.0",CMD_DOC_NONE,NULL,NULL,"stream",COMMAND_GROUP_STREAM,XADD_History,2,XADD_Tips,1,xaddCommand,-5,CMD_WRITE|CMD_DENYOOM|CMD_FAST,ACL_CATEGORY_FAST|ACL_CATEGORY_STREAM|ACL_CATEGORY_WRITE,NULL,XADD_Keyspecs,1,NULL,5),.args=XADD_Args},
{MAKE_CMD("xautoclaim","Changes, or acquires, ownership of messages in a consumer group, as if the messages were delivered to a consumer group member.","O(1) if COUNT is small.","6.2.0",CMD_DOC_NONE,NULL,NULL,"stream",COMMAND_GROUP_STREAM,XAUTOCLAIM_History,1,XAUTOCLAIM_Tips,1,xautoclaimCommand,-6,CMD_WRITE|CMD_FAST,ACL_CATEGORY_FAST|ACL_CATEGORY_STREAM|ACL_CATEGORY_WRITE,NULL,XAUTOCLAIM_Keyspecs,1,NULL,7),.args=XAUTOCLAIM_Args},
{MAKE_CMD("xclaim","Changes, or acquires, ownership of a message in a consumer group, as if the message was delivered to a consumer group member.","O(log N) with N being the number of messages in the PEL of the consumer group.","5.0.0",CMD_DOC_NONE,NULL,NULL,"stream",COMMAND_GROUP_STREAM,XCLAIM_History,0,XCLAIM_Tips,1,xclaimCommand,-6,CMD_WRITE|CMD_FAST,ACL_CATEGORY_FAST|ACL_CATEGORY_STREAM|ACL_CATEGORY_WRITE,NULL,XCLAIM_Keyspecs,1,NULL,11),.args=XCLAIM_Args},
{MAKE_CMD("xdel","Returns the number of messages after removing them from a stream.","O(1) for each single item to delete in the stream, regardless of the stream size.","5.0.0",CMD_DOC_NONE,NULL,NULL,"stream",COMMAND_GROUP_STREAM,XDEL_History,0,XDEL_Tips,0,xdelCommand,-3,CMD_WRITE|CMD_FAST,ACL_CATEGORY_FAST|ACL_CATEGORY_STREAM|ACL_CATEGORY_WRITE,NULL,XDEL_Keyspecs,1,NULL,2),.args=XDEL_Args},
{MAKE_CMD("xdelex","Delete stream message(s) with extended options","O(1) for each single item to delete in the stream, regardless of the stream size","9.2.0",CMD_DOC_NONE,NULL,NULL,"stream",COMMAND_GROUP_STREAM,XDELEX_History,0,XDELEX_Tips,0,xdelexCommand,-5,CMD_WRITE|CMD_FAST,ACL_CATEGORY_FAST|ACL_CATEGORY_WRITE|ACL_CATEGORY_STREAM,NULL,XDELEX_Keyspecs,1,NULL,5),.args=XDELEX_Args},
{MAKE_CMD("xgroup","A container for consumer groups commands.","Depends on subcommand.","5.0.0",CMD_DOC_NONE,NULL,NULL,"stream",COMMAND_GROUP_STREAM,XGROUP_History,0,XGROUP_Tips,0,NULL,-2,0,ACL_CATEGORY_SLOW,NULL,XGROUP_Keyspecs,0,NULL,0),.subcommands=XGROUP_Subcommands},
{MAKE_CMD("xinfo","A container for stream introspection commands.","Depends on subcommand.","5.0.0",CMD_DOC_NONE,NULL,NULL,"stream",COMMAND_GROUP_STREAM,XINFO_History,0,XINFO_Tips,0,NULL,-2,0,ACL_CATEGORY_SLOW,NULL,XINFO_Keyspecs,0,NULL,0),.subcommands=XINFO_Subcommands},
{MAKE_CMD("xlen","Returns the number of messages in a stream.","O(1)","5.0.0",CMD_DOC_NONE,NULL,NULL,"stream",COMMAND_GROUP_STREAM,XLEN_History,0,XLEN_Tips,0,xlenCommand,2,CMD_READONLY|CMD_FAST,ACL_CATEGORY_FAST|ACL_CATEGORY_READ|ACL_CATEGORY_STREAM,NULL,XLEN_Keyspecs,1,NULL,1),.args=XLEN_Args},
Expand Down
97 changes: 97 additions & 0 deletions src/commands/xackdel.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,97 @@
{
"XACKDEL": {
"summary": "Acknowledge and (if possible) delete stream message(s).",
"complexity": "O(1) for each single item to delete in the stream, regardless of the stream size",
"group": "stream",
"since": "9.2.0",
"arity": -6,
"function": "xackdelCommand",
"command_flags": [
"WRITE",
"FAST"
],
"acl_categories": [
"FAST",
"WRITE",
"STREAM"
],
"key_specs": [
{
"flags": [
"RW",
"UPDATE"
],
"begin_search": {
"index": {
"pos": 1
}
},
"find_keys": {
"range": {
"lastkey": 0,
"step": 1,
"limit": 0
}
}
}
],
"arguments": [
{
"name": "key",
"type": "key",
"key_spec_index": 0
},
{
"name": "group",
"type": "string"
},
{
"name": "mode",
"type": "oneof",
"optional": true,
"arguments": [
{
"name": "keepref",
"type": "pure-token",
"token": "KEEPREF"
},
{
"name": "delref",
"type": "pure-token",
"token": "DELREF"
},
{
"name": "acked",
"type": "pure-token",
"token": "ACKED"
}
]
},
{
"name": "ids",
"type": "pure-token",
"token": "IDS"
},
{
"name": "numids",
"type": "integer"
},
{
"name": "id",
"type": "string",
"multiple": true
}
],
"reply_schema": {
"description": "The command returns an integer for each stream message: -1=message not found, 1=acked and deleted, 2=acked but not deleted.",
"type": "array",
"minItems": 1,
"items": {
"description": "Status of the stream message, -1=message not found, 1=acked and deleted, 2=acked but not deleted.",
"type": "integer",
"minimum": -1,
"maximum": 2
}
}
}
}
93 changes: 93 additions & 0 deletions src/commands/xdelex.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,93 @@
{
"XDELEX": {
"summary": "Delete stream message(s) with extended options",
"complexity": "O(1) for each single item to delete in the stream, regardless of the stream size",
"group": "stream",
"since": "9.2.0",
"arity": -5,
"function": "xdelexCommand",
"command_flags": [
"WRITE",
"FAST"
],
"acl_categories": [
"FAST",
"WRITE",
"STREAM"
],
"key_specs": [
{
"flags": [
"RW",
"UPDATE"
],
"begin_search": {
"index": {
"pos": 1
}
},
"find_keys": {
"range": {
"lastkey": 0,
"step": 1,
"limit": 0
}
}
}
],
"arguments": [
{
"name": "key",
"type": "key",
"key_spec_index": 0
},
{
"name": "mode",
"type": "oneof",
"optional": true,
"arguments": [
{
"name": "keepref",
"type": "pure-token",
"token": "KEEPREF"
},
{
"name": "delref",
"type": "pure-token",
"token": "DELREF"
},
{
"name": "acked",
"type": "pure-token",
"token": "ACKED"
}
]
},
{
"name": "ids",
"token": "IDS",
"type": "pure-token"
},
{
"name": "numids",
"type": "integer"
},
{
"name": "id",
"type": "string",
"multiple": true
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
],
"reply_schema": {
"description": "The command returns an integer for each stream message: -1=message not found, 1=message was deleted, 2=message was not deleted due to existing references (ACKED mode).",
"type": "array",
"minItems": 0,
"items": {
"description": "Status of the stream message, -1=message not found, 1=message was deleted, 2=message was not deleted due to existing references (ACKED mode).",
"type": "integer",
"minimum": -1,
"maximum": 2
}
}
}
}
2 changes: 2 additions & 0 deletions src/server.c
Original file line number Diff line number Diff line change
Expand Up @@ -2301,6 +2301,8 @@ void createSharedObjects(void) {
shared.srem = createSharedString("SREM");
shared.xgroup = createSharedString("XGROUP");
shared.xclaim = createSharedString("XCLAIM");
shared.xdel = createSharedString("XDEL");
shared.xack = createSharedString("XACK");
shared.script = createSharedString("SCRIPT");
shared.replconf = createSharedString("REPLCONF");
shared.pexpireat = createSharedString("PEXPIREAT");
Expand Down
4 changes: 3 additions & 1 deletion src/server.h
Original file line number Diff line number Diff line change
Expand Up @@ -1529,7 +1529,7 @@ struct sharedObjectsStruct {
*execaborterr, *noautherr, *noreplicaserr, *busykeyerr, *oomerr, *plus, *messagebulk, *pmessagebulk,
*subscribebulk, *unsubscribebulk, *psubscribebulk, *punsubscribebulk, *del, *unlink, *rpop, *lpop, *lpush, *zadd,
*rpoplpush, *lmove, *blmove, *zpopmin, *zpopmax, *emptyscan, *multi, *exec, *left, *right, *hset, *hsetex, *hdel, *hpexpireat, *hpersist, *srem,
*xgroup, *xclaim, *script, *replconf, *eval, *cluster, *syncslots, *persist, *set, *pexpireat, *pexpire, *time, *pxat, *absttl,
*xgroup, *xclaim, *xdel, *xack, *script, *replconf, *eval, *cluster, *syncslots, *persist, *set, *pexpireat, *pexpire, *time, *pxat, *absttl,
*retrycount, *force, *justid, *entriesread, *lastid, *ping, *setid, *keepttl, *load, *createconsumer, *getack,
*special_asterisk, *special_equals, *default_username, *redacted, *ssubscribebulk, *sunsubscribebulk, *fields,
*finish, *state, *success, *failed, *name, *message,
Expand Down Expand Up @@ -4345,6 +4345,8 @@ void xclaimCommand(client *c);
void xautoclaimCommand(client *c);
void xinfoCommand(client *c);
void xdelCommand(client *c);
void xackdelCommand(client *c);
void xdelexCommand(client *c);
void xtrimCommand(client *c);
void lolwutCommand(client *c);
void aclCommand(client *c);
Expand Down
Loading
Loading