use std::mem::{swap, take};
use super::{
super::{cmd_strings as cs, resp_server_session::RespServerSession},
session_parse_state,
};
use crate::types::RespCommand;
pub const MAX_RESP_ARRAY_LENGTH: usize = 1 << 20;
type PrimaryEntry = (&'static str, RespCommand, bool);
static PRIMARY_TABLE: &[PrimaryEntry] = &[
("ACL", RespCommand::Acl, true),
("APPEND", RespCommand::Append, false),
("ASKING", RespCommand::Asking, false),
("ASYNC", RespCommand::Async, false),
("AUTH", RespCommand::Auth, false),
("BGSAVE", RespCommand::Bgsave, false),
("BITCOUNT", RespCommand::Bitcount, false),
("BITFIELD", RespCommand::Bitfield, false),
("BITFIELD_RO", RespCommand::BitfieldRo, false),
("BITOP", RespCommand::Bitop, true),
("BITPOS", RespCommand::Bitpos, false),
("BLMOVE", RespCommand::Blmove, false),
("BLMPOP", RespCommand::Blmpop, false),
("BLPOP", RespCommand::Blpop, false),
("BRPOP", RespCommand::Brpop, false),
("BRPOPLPUSH", RespCommand::Brpoplpush, false),
("BZMPOP", RespCommand::Bzmpop, false),
("BZPOPMAX", RespCommand::Bzpopmax, false),
("BZPOPMIN", RespCommand::Bzpopmin, false),
("CLIENT", RespCommand::Client, true),
("CLUSTER", RespCommand::Cluster, true),
("COMMAND", RespCommand::Command, true),
("COMMITAOF", RespCommand::Commitaof, false),
("CONFIG", RespCommand::Config, true),
("CUSTOMOBJECTSCAN", RespCommand::Coscan, false),
("DBSIZE", RespCommand::Dbsize, false),
("DEBUG", RespCommand::Debug, false),
("DECR", RespCommand::Decr, false),
("DECRBY", RespCommand::Decrby, false),
("DEL", RespCommand::Del, false),
("DELIFGREATER", RespCommand::Delifgreater, false),
("DISCARD", RespCommand::Discard, false),
("DUMP", RespCommand::Dump, false),
("ECHO", RespCommand::Echo, false),
("EVAL", RespCommand::Eval, false),
("EVALSHA", RespCommand::Evalsha, false),
("EXEC", RespCommand::Exec, false),
("EXISTS", RespCommand::Exists, false),
("EXPDELSCAN", RespCommand::Expdelscan, false),
("EXPIRE", RespCommand::Expire, false),
("EXPIREAT", RespCommand::Expireat, false),
("EXPIRETIME", RespCommand::Expiretime, false),
("FAILOVER", RespCommand::Failover, false),
("FLUSHALL", RespCommand::Flushall, false),
("FLUSHDB", RespCommand::Flushdb, false),
("GEOADD", RespCommand::Geoadd, false),
("GEODIST", RespCommand::Geodist, false),
("GEOHASH", RespCommand::Geohash, false),
("GEOPOS", RespCommand::Geopos, false),
("GEORADIUS", RespCommand::Georadius, false),
("GEORADIUSBYMEMBER", RespCommand::Georadiusbymember, false),
(
"GEORADIUSBYMEMBER_RO",
RespCommand::GeoradiusbymemberRo,
false,
),
("GEORADIUS_RO", RespCommand::GeoradiusRo, false),
("GEOSEARCH", RespCommand::Geosearch, false),
("GEOSEARCHSTORE", RespCommand::Geosearchstore, false),
("GET", RespCommand::Get, false),
("GETBIT", RespCommand::Getbit, false),
("GETDEL", RespCommand::Getdel, false),
("GETEX", RespCommand::Getex, false),
("GETIFNOTMATCH", RespCommand::Getifnotmatch, false),
("GETRANGE", RespCommand::Getrange, false),
("GETSET", RespCommand::Getset, false),
("GETWITHETAG", RespCommand::Getwithetag, false),
("HCOLLECT", RespCommand::Hcollect, false),
("HDEL", RespCommand::Hdel, false),
("HELLO", RespCommand::Hello, false),
("HEXISTS", RespCommand::Hexists, false),
("HEXPIRE", RespCommand::Hexpire, false),
("HEXPIREAT", RespCommand::Hexpireat, false),
("HEXPIRETIME", RespCommand::Hexpiretime, false),
("HGET", RespCommand::Hget, false),
("HGETALL", RespCommand::Hgetall, false),
("HINCRBY", RespCommand::Hincrby, false),
("HINCRBYFLOAT", RespCommand::Hincrbyfloat, false),
("HKEYS", RespCommand::Hkeys, false),
("HLEN", RespCommand::Hlen, false),
("HMGET", RespCommand::Hmget, false),
("HMSET", RespCommand::Hmset, false),
("HPERSIST", RespCommand::Hpersist, false),
("HPEXPIRE", RespCommand::Hpexpire, false),
("HPEXPIREAT", RespCommand::Hpexpireat, false),
("HPEXPIRETIME", RespCommand::Hpexpiretime, false),
("HPTTL", RespCommand::Hpttl, false),
("HRANDFIELD", RespCommand::Hrandfield, false),
("HSCAN", RespCommand::Hscan, false),
("HSET", RespCommand::Hset, false),
("HSETNX", RespCommand::Hsetnx, false),
("HSTRLEN", RespCommand::Hstrlen, false),
("HTTL", RespCommand::Httl, false),
("HVALS", RespCommand::Hvals, false),
("INCR", RespCommand::Incr, false),
("INCRBY", RespCommand::Incrby, false),
("INCRBYFLOAT", RespCommand::Incrbyfloat, false),
("INFO", RespCommand::Info, false),
("KEYS", RespCommand::Keys, false),
("LASTSAVE", RespCommand::Lastsave, false),
("LATENCY", RespCommand::Latency, true),
("LCS", RespCommand::Lcs, false),
("LINDEX", RespCommand::Lindex, false),
("LINSERT", RespCommand::Linsert, false),
("LLEN", RespCommand::Llen, false),
("LMOVE", RespCommand::Lmove, false),
("LMPOP", RespCommand::Lmpop, false),
("LPOP", RespCommand::Lpop, false),
("LPOS", RespCommand::Lpos, false),
("LPUSH", RespCommand::Lpush, false),
("LPUSHX", RespCommand::Lpushx, false),
("LRANGE", RespCommand::Lrange, false),
("LREM", RespCommand::Lrem, false),
("LSET", RespCommand::Lset, false),
("LTRIM", RespCommand::Ltrim, false),
("MEMORY", RespCommand::Memory, true),
("MGET", RespCommand::Mget, false),
("MIGRATE", RespCommand::Migrate, false),
("MODULE", RespCommand::Module, true),
("MONITOR", RespCommand::Monitor, false),
("MSET", RespCommand::Mset, false),
("MSETNX", RespCommand::Msetnx, false),
("MULTI", RespCommand::Multi, false),
("OBJECT", RespCommand::Object, true),
("PERSIST", RespCommand::Persist, false),
("PEXPIRE", RespCommand::Pexpire, false),
("PEXPIREAT", RespCommand::Pexpireat, false),
("PEXPIRETIME", RespCommand::Pexpiretime, false),
("PFADD", RespCommand::Pfadd, false),
("PFCOUNT", RespCommand::Pfcount, false),
("PFMERGE", RespCommand::Pfmerge, false),
("PING", RespCommand::Ping, false),
("PSETEX", RespCommand::Psetex, false),
("PSUBSCRIBE", RespCommand::Psubscribe, false),
("PTTL", RespCommand::Pttl, false),
("PUBLISH", RespCommand::Publish, false),
("PUBSUB", RespCommand::Pubsub, true),
("PUNSUBSCRIBE", RespCommand::Punsubscribe, false),
("QUIT", RespCommand::Quit, false),
("READONLY", RespCommand::Readonly, false),
("READWRITE", RespCommand::Readwrite, false),
("REGISTERCS", RespCommand::Registercs, false),
("RENAME", RespCommand::Rename, false),
("RENAMENX", RespCommand::Renamenx, false),
("REPLICAOF", RespCommand::Replicaof, false),
("RESTORE", RespCommand::Restore, false),
("RI.CONFIG", RespCommand::Riconfig, false),
("RI.CREATE", RespCommand::Ricreate, false),
("RI.DEL", RespCommand::Ridel, false),
("RI.EXISTS", RespCommand::Riexists, false),
("RI.GET", RespCommand::Riget, false),
("RI.METRICS", RespCommand::Rimetrics, false),
("RI.RANGE", RespCommand::Rirange, false),
("RI.SCAN", RespCommand::Riscan, false),
("RI.SET", RespCommand::Riset, false),
("ROLE", RespCommand::Role, false),
("RPOP", RespCommand::Rpop, false),
("RPOPLPUSH", RespCommand::Rpoplpush, false),
("RPUSH", RespCommand::Rpush, false),
("RPUSHX", RespCommand::Rpushx, false),
("RUNTXP", RespCommand::Runtxp, false),
("SADD", RespCommand::Sadd, false),
("SAVE", RespCommand::Save, false),
("SCAN", RespCommand::Scan, false),
("SCARD", RespCommand::Scard, false),
("SCRIPT", RespCommand::Script, true),
("SDIFF", RespCommand::Sdiff, false),
("SDIFFSTORE", RespCommand::Sdiffstore, false),
("SECONDARYOF", RespCommand::Secondaryof, false),
("SELECT", RespCommand::Select, false),
("SET", RespCommand::Set, false),
("SETBIT", RespCommand::Setbit, false),
("SETEX", RespCommand::Setex, false),
("SETIFGREATER", RespCommand::Setifgreater, false),
("SETIFMATCH", RespCommand::Setifmatch, false),
("SETNX", RespCommand::Setnx, false),
("SETRANGE", RespCommand::Setrange, false),
("SETWITHETAG", RespCommand::Setwithetag, false),
("SINTER", RespCommand::Sinter, false),
("SINTERCARD", RespCommand::Sintercard, false),
("SINTERSTORE", RespCommand::Sinterstore, false),
("SISMEMBER", RespCommand::Sismember, false),
("SLAVEOF", RespCommand::Secondaryof, false),
("SLOWLOG", RespCommand::Slowlog, true),
("SMEMBERS", RespCommand::Smembers, false),
("SMISMEMBER", RespCommand::Smismember, false),
("SMOVE", RespCommand::Smove, false),
("SPOP", RespCommand::Spop, false),
("SPUBLISH", RespCommand::Spublish, false),
("SRANDMEMBER", RespCommand::Srandmember, false),
("SREM", RespCommand::Srem, false),
("SSCAN", RespCommand::Sscan, false),
("SSUBSCRIBE", RespCommand::Ssubscribe, false),
("STRLEN", RespCommand::Strlen, false),
("SUBSCRIBE", RespCommand::Subscribe, false),
("SUBSTR", RespCommand::Substr, false),
("SUNION", RespCommand::Sunion, false),
("SUNIONSTORE", RespCommand::Sunionstore, false),
("SWAPDB", RespCommand::Swapdb, false),
("TIME", RespCommand::Time, false),
("TTL", RespCommand::Ttl, false),
("TYPE", RespCommand::Type, false),
("UNLINK", RespCommand::Unlink, false),
("UNSUBSCRIBE", RespCommand::Unsubscribe, false),
("UNWATCH", RespCommand::Unwatch, false),
("VADD", RespCommand::Vadd, false),
("VCARD", RespCommand::Vcard, false),
("VDIM", RespCommand::Vdim, false),
("VEMB", RespCommand::Vemb, false),
("VGETATTR", RespCommand::Vgetattr, false),
("VINFO", RespCommand::Vinfo, false),
("VISMEMBER", RespCommand::Vismember, false),
("VLINKS", RespCommand::Vlinks, false),
("VRANDMEMBER", RespCommand::Vrandmember, false),
("VREM", RespCommand::Vrem, false),
("VSETATTR", RespCommand::Vsetattr, false),
("VSIM", RespCommand::Vsim, false),
("WATCH", RespCommand::Watch, false),
("WATCHMS", RespCommand::Watchms, false),
("WATCHOS", RespCommand::Watchos, false),
("ZADD", RespCommand::Zadd, false),
("ZCARD", RespCommand::Zcard, false),
("ZCOLLECT", RespCommand::Zcollect, false),
("ZCOUNT", RespCommand::Zcount, false),
("ZDIFF", RespCommand::Zdiff, false),
("ZDIFFSTORE", RespCommand::Zdiffstore, false),
("ZEXPIRE", RespCommand::Zexpire, false),
("ZEXPIREAT", RespCommand::Zexpireat, false),
("ZEXPIRETIME", RespCommand::Zexpiretime, false),
("ZINCRBY", RespCommand::Zincrby, false),
("ZINTER", RespCommand::Zinter, false),
("ZINTERCARD", RespCommand::Zintercard, false),
("ZINTERSTORE", RespCommand::Zinterstore, false),
("ZLEXCOUNT", RespCommand::Zlexcount, false),
("ZMPOP", RespCommand::Zmpop, false),
("ZMSCORE", RespCommand::Zmscore, false),
("ZPERSIST", RespCommand::Zpersist, false),
("ZPEXPIRE", RespCommand::Zpexpire, false),
("ZPEXPIREAT", RespCommand::Zpexpireat, false),
("ZPEXPIRETIME", RespCommand::Zpexpiretime, false),
("ZPOPMAX", RespCommand::Zpopmax, false),
("ZPOPMIN", RespCommand::Zpopmin, false),
("ZPTTL", RespCommand::Zpttl, false),
("ZRANDMEMBER", RespCommand::Zrandmember, false),
("ZRANGE", RespCommand::Zrange, false),
("ZRANGEBYLEX", RespCommand::Zrangebylex, false),
("ZRANGEBYSCORE", RespCommand::Zrangebyscore, false),
("ZRANGESTORE", RespCommand::Zrangestore, false),
("ZRANK", RespCommand::Zrank, false),
("ZREM", RespCommand::Zrem, false),
("ZREMRANGEBYLEX", RespCommand::Zremrangebylex, false),
("ZREMRANGEBYRANK", RespCommand::Zremrangebyrank, false),
("ZREMRANGEBYSCORE", RespCommand::Zremrangebyscore, false),
("ZREVRANGE", RespCommand::Zrevrange, false),
("ZREVRANGEBYLEX", RespCommand::Zrevrangebylex, false),
("ZREVRANGEBYSCORE", RespCommand::Zrevrangebyscore, false),
("ZREVRANK", RespCommand::Zrevrank, false),
("ZSCAN", RespCommand::Zscan, false),
("ZSCORE", RespCommand::Zscore, false),
("ZTTL", RespCommand::Zttl, false),
("ZUNION", RespCommand::Zunion, false),
("ZUNIONSTORE", RespCommand::Zunionstore, false),
];
static CLIENT_SUBTABLE: &[(&str, RespCommand)] = &[
("GETNAME", RespCommand::ClientGetname),
("ID", RespCommand::ClientId),
("INFO", RespCommand::ClientInfo),
("KILL", RespCommand::ClientKill),
("LIST", RespCommand::ClientList),
("SETINFO", RespCommand::ClientSetinfo),
("SETNAME", RespCommand::ClientSetname),
("UNBLOCK", RespCommand::ClientUnblock),
];
static CONFIG_SUBTABLE: &[(&str, RespCommand)] = &[
("GET", RespCommand::ConfigGet),
("REWRITE", RespCommand::ConfigRewrite),
("SET", RespCommand::ConfigSet),
];
static COMMAND_SUBTABLE: &[(&str, RespCommand)] = &[
("COUNT", RespCommand::CommandCount),
("DOCS", RespCommand::CommandDocs),
("GETKEYS", RespCommand::CommandGetkeys),
("GETKEYSANDFLAGS", RespCommand::CommandGetkeysandflags),
("INFO", RespCommand::CommandInfo),
];
static ACL_SUBTABLE: &[(&str, RespCommand)] = &[
("CAT", RespCommand::AclCat),
("DELUSER", RespCommand::AclDeluser),
("GENPASS", RespCommand::AclGenpass),
("GETUSER", RespCommand::AclGetuser),
("LIST", RespCommand::AclList),
("LOAD", RespCommand::AclLoad),
("SAVE", RespCommand::AclSave),
("SETUSER", RespCommand::AclSetuser),
("USERS", RespCommand::AclUsers),
("WHOAMI", RespCommand::AclWhoami),
];
static SCRIPT_SUBTABLE: &[(&str, RespCommand)] = &[
("EXISTS", RespCommand::ScriptExists),
("FLUSH", RespCommand::ScriptFlush),
("LOAD", RespCommand::ScriptLoad),
];
static PUBSUB_SUBTABLE: &[(&str, RespCommand)] = &[
("CHANNELS", RespCommand::PubsubChannels),
("NUMPAT", RespCommand::PubsubNumpat),
("NUMSUB", RespCommand::PubsubNumsub),
];
static LATENCY_SUBTABLE: &[(&str, RespCommand)] = &[
("HELP", RespCommand::LatencyHelp),
("HISTOGRAM", RespCommand::LatencyHistogram),
("RESET", RespCommand::LatencyReset),
];
static SLOWLOG_SUBTABLE: &[(&str, RespCommand)] = &[
("GET", RespCommand::SlowlogGet),
("HELP", RespCommand::SlowlogHelp),
("LEN", RespCommand::SlowlogLen),
("RESET", RespCommand::SlowlogReset),
];
static MODULE_SUBTABLE: &[(&str, RespCommand)] = &[("LOADCS", RespCommand::ModuleLoadcs)];
static MEMORY_SUBTABLE: &[(&str, RespCommand)] = &[("USAGE", RespCommand::MemoryUsage)];
static OBJECT_SUBTABLE: &[(&str, RespCommand)] = &[
("ENCODING", RespCommand::ObjectEncoding),
("FREQ", RespCommand::ObjectFreq),
("HELP", RespCommand::ObjectHelp),
("IDLETIME", RespCommand::ObjectIdletime),
("REFCOUNT", RespCommand::ObjectRefcount),
];
static CLUSTER_SUBTABLE: &[(&str, RespCommand)] = &[
("ADDSLOTS", RespCommand::ClusterAddslots),
("ADDSLOTSRANGE", RespCommand::ClusterAddslotsrange),
("ADVANCE_TIME", RespCommand::ClusterAdvanceTime),
("APPENDLOG", RespCommand::ClusterAppendlog),
("ATTACH_SYNC", RespCommand::ClusterAttachSync),
("BANLIST", RespCommand::ClusterBanlist),
(
"BEGIN_REPLICA_RECOVER",
RespCommand::ClusterBeginReplicaRecover,
),
("BUMPEPOCH", RespCommand::ClusterBumpepoch),
("COUNTKEYSINSLOT", RespCommand::ClusterCountkeysinslot),
("DELKEYSINSLOT", RespCommand::ClusterDelkeysinslot),
("DELKEYSINSLOTRANGE", RespCommand::ClusterDelkeysinslotrange),
("DELSLOTS", RespCommand::ClusterDelslots),
("DELSLOTSRANGE", RespCommand::ClusterDelslotsrange),
("ENDPOINT", RespCommand::ClusterEndpoint),
("FAILOVER", RespCommand::ClusterFailover),
(
"FAILREPLICATIONOFFSET",
RespCommand::ClusterFailreplicationoffset,
),
("FAILSTOPWRITES", RespCommand::ClusterFailstopwrites),
("FLUSHALL", RespCommand::ClusterFlushall),
("FORGET", RespCommand::ClusterForget),
("GETKEYSINSLOT", RespCommand::ClusterGetkeysinslot),
("GOSSIP", RespCommand::ClusterGossip),
("HELP", RespCommand::ClusterHelp),
("INFO", RespCommand::ClusterInfo),
(
"INITIATE_REPLICA_SYNC",
RespCommand::ClusterInitiateReplicaSync,
),
("KEYSLOT", RespCommand::ClusterKeyslot),
("MEET", RespCommand::ClusterMeet),
("MIGRATE", RespCommand::ClusterMigrate),
("MLOG_KEY_TIME", RespCommand::ClusterMlogKeyTime),
("MTASKS", RespCommand::ClusterMtasks),
("MYID", RespCommand::ClusterMyid),
("MYPARENTID", RespCommand::ClusterMyparentid),
("NODES", RespCommand::ClusterNodes),
("PUBLISH", RespCommand::ClusterPublish),
("REPLICAS", RespCommand::ClusterReplicas),
("REPLICATE", RespCommand::ClusterReplicate),
("RESERVE", RespCommand::ClusterReserve),
("RESET", RespCommand::ClusterReset),
(
"SEND_CKPT_FILE_SEGMENT",
RespCommand::ClusterSendCkptFileSegment,
),
("SEND_CKPT_METADATA", RespCommand::ClusterSendCkptMetadata),
("SET-CONFIG-EPOCH", RespCommand::ClusterSetconfigepoch),
("SETSLOT", RespCommand::ClusterSetslot),
("SETSLOTSRANGE", RespCommand::ClusterSetslotsrange),
("SHARDS", RespCommand::ClusterShards),
("SLOTS", RespCommand::ClusterSlots),
("SLOTSTATE", RespCommand::ClusterSlotstate),
("SNAPSHOT_DATA", RespCommand::ClusterSnapshotData),
("SPUBLISH", RespCommand::ClusterSpublish),
("SYNC", RespCommand::ClusterSync),
];
static BITOP_SUBTABLE: &[(&str, RespCommand)] = &[
("AND", RespCommand::BitopAnd),
("DIFF", RespCommand::BitopDiff),
("NOT", RespCommand::BitopNot),
("OR", RespCommand::BitopOr),
("XOR", RespCommand::BitopXor),
];
fn lookup_in_table(table: &[(&str, RespCommand)], name: &[u8]) -> Option<RespCommand> {
table
.binary_search_by(|(entry_name, _)| entry_name.as_bytes().cmp(name))
.ok()
.map(|idx| table[idx].1)
}
fn lookup_primary(name: &[u8]) -> Option<(RespCommand, bool)> {
let idx = PRIMARY_TABLE
.binary_search_by(|(entry_name, ..)| entry_name.as_bytes().cmp(name))
.ok()?;
let (_, cmd, has_subcommands) = PRIMARY_TABLE[idx];
Some((cmd, has_subcommands))
}
fn lookup_subcommand(parent: RespCommand, name: &[u8]) -> Option<RespCommand> {
let table = match parent {
RespCommand::Client => CLIENT_SUBTABLE,
RespCommand::Config => CONFIG_SUBTABLE,
RespCommand::Command => COMMAND_SUBTABLE,
RespCommand::Acl => ACL_SUBTABLE,
RespCommand::Script => SCRIPT_SUBTABLE,
RespCommand::Pubsub => PUBSUB_SUBTABLE,
RespCommand::Latency => LATENCY_SUBTABLE,
RespCommand::Slowlog => SLOWLOG_SUBTABLE,
RespCommand::Module => MODULE_SUBTABLE,
RespCommand::Memory => MEMORY_SUBTABLE,
RespCommand::Object => OBJECT_SUBTABLE,
RespCommand::Cluster => CLUSTER_SUBTABLE,
RespCommand::Bitop => BITOP_SUBTABLE,
_ => &[],
};
lookup_in_table(table, name)
}
static AOF_INDEPENDENT_COMMANDS: &[RespCommand] = &[
RespCommand::Async,
RespCommand::Ping,
RespCommand::Select,
RespCommand::Swapdb,
RespCommand::Echo,
RespCommand::Monitor,
RespCommand::ModuleLoadcs,
RespCommand::Registercs,
RespCommand::Info,
RespCommand::Time,
RespCommand::Lastsave,
RespCommand::AclCat,
RespCommand::AclDeluser,
RespCommand::AclGenpass,
RespCommand::AclGetuser,
RespCommand::AclList,
RespCommand::AclLoad,
RespCommand::AclSave,
RespCommand::AclSetuser,
RespCommand::AclUsers,
RespCommand::AclWhoami,
RespCommand::ClientId,
RespCommand::ClientInfo,
RespCommand::ClientList,
RespCommand::ClientKill,
RespCommand::ClientGetname,
RespCommand::ClientSetname,
RespCommand::ClientSetinfo,
RespCommand::ClientUnblock,
RespCommand::Command,
RespCommand::CommandCount,
RespCommand::CommandDocs,
RespCommand::CommandInfo,
RespCommand::CommandGetkeys,
RespCommand::CommandGetkeysandflags,
RespCommand::MemoryUsage,
RespCommand::ConfigGet,
RespCommand::ConfigRewrite,
RespCommand::ConfigSet,
RespCommand::LatencyHelp,
RespCommand::LatencyHistogram,
RespCommand::LatencyReset,
RespCommand::SlowlogHelp,
RespCommand::SlowlogLen,
RespCommand::SlowlogGet,
RespCommand::SlowlogReset,
RespCommand::Multi,
];
#[inline]
pub fn is_aof_independent(cmd: RespCommand) -> bool {
AOF_INDEPENDENT_COMMANDS.contains(&cmd)
}
static FAST_PATTERN_TABLE: &[(&[u8], RespCommand, u8)] = &[
(b"*2\r\n$3\r\nGET\r\n", RespCommand::Get, 1),
(b"*3\r\n$3\r\nSET\r\n", RespCommand::Set, 2),
(b"*2\r\n$3\r\nDEL\r\n", RespCommand::Del, 1),
(b"*2\r\n$3\r\nTTL\r\n", RespCommand::Ttl, 1),
(b"*1\r\n$4\r\nPING\r\n", RespCommand::Ping, 0),
(b"*2\r\n$4\r\nINCR\r\n", RespCommand::Incr, 1),
(b"*2\r\n$4\r\nDECR\r\n", RespCommand::Decr, 1),
(b"*1\r\n$4\r\nEXEC\r\n", RespCommand::Exec, 0),
(b"*2\r\n$4\r\nPTTL\r\n", RespCommand::Pttl, 1),
(b"*1\r\n$5\r\nMULTI\r\n", RespCommand::Multi, 0),
(b"*3\r\n$5\r\nSETNX\r\n", RespCommand::Setnx, 2),
(b"*4\r\n$5\r\nSETEX\r\n", RespCommand::Setex, 3),
(b"*2\r\n$6\r\nEXISTS\r\n", RespCommand::Exists, 1),
(b"*2\r\n$6\r\nGETDEL\r\n", RespCommand::Getdel, 1),
(b"*3\r\n$6\r\nAPPEND\r\n", RespCommand::Append, 2),
(b"*3\r\n$6\r\nINCRBY\r\n", RespCommand::Incrby, 2),
(b"*3\r\n$6\r\nDECRBY\r\n", RespCommand::Decrby, 2),
(b"*4\r\n$6\r\nPSETEX\r\n", RespCommand::Psetex, 3),
];
#[inline]
fn pattern_matches(buffer: &[u8], start: usize, pattern: &[u8]) -> bool {
buffer.len() >= start + pattern.len() && &buffer[start..start + pattern.len()] == pattern
}
#[derive(Default)]
pub(crate) struct MruCommandCache {
slot0: Option<MruEntry>,
slot1: Option<MruEntry>,
}
#[derive(Clone, Copy)]
struct MruEntry {
pattern: [u8; 16],
len: u8,
cmd: RespCommand,
count: u8,
}
impl MruCommandCache {
fn update(&mut self, frame: &[u8], consumed: usize, cmd: RespCommand, count: u8) {
debug_assert!((13..=16).contains(&consumed) && frame.len() >= consumed);
let mut entry = MruEntry {
pattern: [0; 16],
len: consumed as u8,
cmd,
count,
};
entry.pattern[..consumed].copy_from_slice(&frame[..consumed]);
self.slot1 = self.slot0;
self.slot0 = Some(entry);
}
fn lookup(&self, buffer: &[u8], start: usize) -> Option<(MruEntry, usize)> {
for (slot_idx, slot) in [self.slot0, self.slot1].into_iter().enumerate() {
let Some(entry) = slot else { continue };
let len = entry.len as usize;
if pattern_matches(buffer, start, &entry.pattern[..len]) {
return Some((entry, slot_idx));
}
}
None
}
fn promote(&mut self, slot_idx: usize) {
if slot_idx == 1 {
swap(&mut self.slot0, &mut self.slot1);
}
}
}
impl RespServerSession {
pub fn parse_command(&mut self) -> Option<RespCommand> {
self.parse_command_with(true)
}
fn parse_command_with(&mut self, write_error_on_failure: bool) -> Option<RespCommand> {
let mut count: isize = -1;
self.end_read_head = self.read_head;
let mut cmd = self.fast_parse_command(&mut count);
if cmd == RespCommand::None {
let cmd_start_offset = self.read_head;
cmd = self.array_parse_command(&mut count, write_error_on_failure)?;
if cmd != RespCommand::Invalid
&& cmd != RespCommand::None
&& !matches!(
cmd,
RespCommand::Customtxn
| RespCommand::Customprocedure
| RespCommand::Customrawstringcmd
| RespCommand::Customobjcmd
)
{
self.update_command_cache(cmd_start_offset, cmd, count);
}
}
if count > MAX_RESP_ARRAY_LENGTH as isize {
return None;
}
let count = count.max(0) as usize;
self.parse_state.initialize(count);
let mut ptr = self.read_head;
for i in 0..count {
if !session_parse_state::read(
&mut self.parse_state,
i,
&self.recv_buffer,
&mut ptr,
self.bytes_read,
) {
return None;
}
}
self.end_read_head = ptr;
self.handle_aof_commit_mode(cmd);
Some(cmd)
}
fn update_command_cache(&mut self, cmd_start_offset: usize, cmd: RespCommand, arg_count: isize) {
let arg_count = arg_count.max(0) as usize;
if arg_count > usize::from(u8::MAX) {
return;
}
if self.bytes_read - cmd_start_offset < 16 {
return;
}
let consumed = self.read_head - cmd_start_offset;
if !(13..=16).contains(&consumed) {
return;
}
let frame = &self.recv_buffer[cmd_start_offset..cmd_start_offset + 16];
self.mru_cache.update(frame, consumed, cmd, arg_count as u8);
}
pub fn fast_parse_command(&mut self, count: &mut isize) -> RespCommand {
let start = self.read_head;
let remaining = self.bytes_read.saturating_sub(start);
if remaining >= 16 && self.recv_buffer[start] == b'*' {
for (pattern, cmd, arg_count) in FAST_PATTERN_TABLE {
if pattern_matches(&self.recv_buffer, start, pattern) {
self.read_head += pattern.len();
*count = isize::from(*arg_count);
return *cmd;
}
}
if let Some((entry, slot_idx)) = self.mru_cache.lookup(&self.recv_buffer, start) {
self.mru_cache.promote(slot_idx);
self.read_head += entry.len as usize;
*count = isize::from(entry.count);
return entry.cmd;
}
}
if remaining >= 8
&& self.recv_buffer[start] == b'*'
&& self.recv_buffer[start + 2] == b'\r'
&& self.recv_buffer[start + 3] == b'\n'
&& self.recv_buffer[start + 4] == b'$'
&& self.recv_buffer[start + 6] == b'\r'
&& self.recv_buffer[start + 7] == b'\n'
{
*count = i64::from(self.recv_buffer[start + 1]) as isize - isize::from(b'1');
let length = i64::from(self.recv_buffer[start + 5]) - i64::from(b'0');
if (1..=9).contains(&length) && remaining >= length as usize + 10 {
let frame_len = length as usize + 10;
let frame_end = start + frame_len;
self.read_head += frame_len;
for (pattern, cmd, _) in FAST_PATTERN_TABLE {
if pattern.len() == frame_len && pattern_matches(&self.recv_buffer, start, pattern) {
return *cmd;
}
}
let buf = &self.recv_buffer;
let last_word = &buf[start + length as usize + 2..frame_end];
let prefix = &buf[start + 8..start + 10];
return match (*count, length) {
(2, 7) if last_word == b"UBLISH\r\n" && buf[start + 8] == b'P' => RespCommand::Publish,
(2, 8) if last_word == b"UBLISH\r\n" && prefix == b"SP" => RespCommand::Spublish,
(3, 8) if last_word == b"TRANGE\r\n" && prefix == b"SE" => RespCommand::Setrange,
(3, 8) if last_word == b"TRANGE\r\n" && prefix == b"GE" => RespCommand::Getrange,
(3..=7, 3) if last_word == b"3\r\nSET\r\n" => RespCommand::Setexnx,
(1..=3, 5) if last_word == b"\nGETEX\r\n" => RespCommand::Getex,
(2..=3, 6) if last_word == b"EXPIRE\r\n" => RespCommand::Expire,
(2..=3, 7) if last_word == b"EXPIRE\r\n" && buf[start + 8] == b'P' => {
RespCommand::Pexpire
}
_ => self.matched_none(start, count),
};
}
*count = -1;
return RespCommand::None;
}
self.fast_parse_inline_command(count)
}
pub fn fast_parse_inline_command(&mut self, count: &mut isize) -> RespCommand {
let start = self.read_head;
*count = 0;
if self.bytes_read - start >= 6 {
let word = &self.recv_buffer[start..start + 4];
if &self.recv_buffer[start + 4..start + 6] == b"\r\n" {
self.read_head += 6;
if word == b"PING" {
return RespCommand::Ping;
}
if word == b"QUIT" {
return RespCommand::Quit;
}
self.read_head -= 6;
}
}
RespCommand::None
}
fn matched_none(&mut self, old_read_head: usize, count: &mut isize) -> RespCommand {
self.read_head = old_read_head;
*count = -1;
RespCommand::None
}
pub fn try_parse_custom_command(&mut self, command: &[u8]) -> Option<RespCommand> {
let _ = command;
None
}
pub fn attempt_skip_line(&mut self) -> bool {
let mut string_end = self.read_head;
while string_end + 1 < self.bytes_read {
if self.recv_buffer[string_end] == b'\r' && self.recv_buffer[string_end + 1] == b'\n' {
self.read_head = string_end + 2;
self.end_read_head = self.read_head;
return true;
}
string_end += 1;
}
false
}
pub fn parse_resp_command_buffer(&mut self, buffer: &[u8]) -> Option<RespCommand> {
let saved = self.take_receive_state();
self.recv_buffer.clear();
self.recv_buffer.extend_from_slice(buffer);
self.bytes_read = self.recv_buffer.len();
self.read_head = 0;
let parsed = self.parse_command_with(false);
self.restore_receive_state(saved);
parsed.filter(|cmd| *cmd != RespCommand::Invalid)
}
pub fn fuzz_parse_command_buffer(&mut self, buffer: &[u8]) -> (bool, RespCommand) {
let saved = self.take_receive_state();
self.recv_buffer.clear();
self.recv_buffer.extend_from_slice(buffer);
self.bytes_read = self.recv_buffer.len();
self.read_head = 0;
let cmd = if self.bytes_read >= 4 {
self
.parse_command_with(false)
.unwrap_or(RespCommand::Invalid)
} else {
RespCommand::Invalid
};
self.restore_receive_state(saved);
(cmd != RespCommand::Invalid, cmd)
}
pub fn handle_aof_commit_mode(&mut self, cmd: RespCommand) {
if self.pending_output_len() == 0 {
self.wait_for_aof_blocking = false;
}
if self.txn_state == TxnState::Started {
return;
}
self.wait_for_aof_blocking = self.wait_for_aof_blocking || !is_aof_independent(cmd);
}
pub fn array_parse_command(
&mut self,
count: &mut isize,
write_error_on_failure: bool,
) -> Option<RespCommand> {
self.end_read_head = self.read_head;
let start = self.read_head;
if start < self.bytes_read && self.make_upper_case(start, self.bytes_read - start) {
let cmd = self.fast_parse_command(count);
if cmd != RespCommand::None {
return Some(cmd);
}
}
if start >= self.bytes_read || self.recv_buffer[start] != b'*' {
if !self.attempt_skip_line() {
return None;
}
return Some(RespCommand::Invalid);
}
let mut ptr = start + 1;
let mut array_len = 0usize;
while ptr < self.bytes_read && self.recv_buffer[ptr].is_ascii_digit() {
array_len = array_len
.saturating_mul(10)
.saturating_add(usize::from(self.recv_buffer[ptr] - b'0'));
ptr += 1;
}
if ptr + 2 > self.bytes_read || &self.recv_buffer[ptr..ptr + 2] != b"\r\n" {
return None;
}
ptr += 2;
self.read_head = ptr;
*count = array_len as isize;
let mut specific_error: Option<Vec<u8>> = None;
let cmd = self.hash_lookup_command(count, &mut specific_error)?;
if write_error_on_failure && cmd == RespCommand::Invalid {
if let Some(error) = specific_error {
self.output.push(b'-');
self.output.extend_from_slice(&error);
self.output.extend_from_slice(b"\r\n");
self.command_error_written = true;
} else {
self.abort_error_message(cs::RESP_ERR_GENERIC_UNK_CMD);
}
}
Some(cmd)
}
pub fn hash_lookup_command(
&mut self,
count: &mut isize,
specific_error: &mut Option<Vec<u8>>,
) -> Option<RespCommand> {
let command = self.get_command()?;
*count -= 1;
match lookup_primary(&command) {
None => {
if let Some(custom_cmd) = self.try_parse_custom_command(&command) {
return Some(custom_cmd);
}
Some(RespCommand::Invalid)
}
Some((cmd, has_subcommands)) => {
if has_subcommands {
self.handle_subcommand_lookup(cmd, count, specific_error)
} else {
Some(cmd)
}
}
}
}
pub fn handle_subcommand_lookup(
&mut self,
parent_cmd: RespCommand,
count: &mut isize,
specific_error: &mut Option<Vec<u8>>,
) -> Option<RespCommand> {
if parent_cmd == RespCommand::Command && *count == 0 {
return Some(RespCommand::Command);
}
if *count == 0 {
let parent = format!("{parent_cmd:?}").to_uppercase();
*specific_error = Some(if parent_cmd == RespCommand::Bitop {
cs::RESP_ERR_GENERIC_SYNTAX_ERROR.as_bytes().to_vec()
} else {
cs::GENERIC_ERR_WRONG_NUM_ARGS
.replace("{0}", &parent)
.into_bytes()
});
return Some(RespCommand::Invalid);
}
let sub_command = self.get_upper_case_command()?;
*count -= 1;
if let Some(sub_cmd) = lookup_subcommand(parent_cmd, &sub_command) {
return Some(sub_cmd);
}
let sub_text = String::from_utf8_lossy(&sub_command);
let parent = format!("{parent_cmd:?}").to_uppercase();
*specific_error = Some(if parent_cmd == RespCommand::Bitop {
cs::RESP_ERR_GENERIC_SYNTAX_ERROR.as_bytes().to_vec()
} else if matches!(parent_cmd, RespCommand::Cluster | RespCommand::Latency) {
cs::GENERIC_ERR_UNKNOWN_SUB_COMMAND
.replace("{0}", &sub_text)
.replace("{1}", &parent)
.into_bytes()
} else {
cs::GENERIC_ERR_UNKNOWN_SUB_COMMAND_NO_HELP
.replace("{0}", &sub_text)
.into_bytes()
});
Some(RespCommand::Invalid)
}
fn take_receive_state(&mut self) -> (Vec<u8>, usize, usize, usize) {
(
take(&mut self.recv_buffer),
self.bytes_read,
self.read_head,
self.end_read_head,
)
}
fn restore_receive_state(&mut self, saved: (Vec<u8>, usize, usize, usize)) {
let (buffer, bytes_read, read_head, end_read_head) = saved;
self.recv_buffer = buffer;
self.bytes_read = bytes_read;
self.read_head = read_head;
self.end_read_head = end_read_head;
}
}
use crate::resp::resp_server_session::TxnState;
#[cfg(test)]
mod tests {
use super::*;
type ParentSubtable = (
&'static str,
RespCommand,
&'static [(&'static str, RespCommand)],
);
fn parse_one(
session: &mut RespServerSession,
buffer: &[u8],
) -> (Option<RespCommand>, Vec<Vec<u8>>) {
session.recv_buffer.clear();
session.recv_buffer.extend_from_slice(buffer);
session.bytes_read = session.recv_buffer.len();
session.read_head = 0;
let cmd = session.parse_command();
let args = (0..session.parse_state.count)
.map(|i| {
session
.parse_state
.get_arg_slice_by_ref(i)
.as_slice()
.to_vec()
})
.collect();
(cmd, args)
}
#[test]
fn fast_paths_parse_hot_commands() {
let mut s = RespServerSession::default();
let (cmd, args) = parse_one(&mut s, b"*2\r\n$3\r\nGET\r\n$3\r\nfoo\r\n");
assert_eq!(cmd, Some(RespCommand::Get));
assert_eq!(args, vec![b"foo".to_vec()]);
let (cmd, args) = parse_one(&mut s, b"*3\r\n$3\r\nSET\r\n$1\r\nk\r\n$2\r\nv1\r\n");
assert_eq!(cmd, Some(RespCommand::Set));
assert_eq!(args, vec![b"k".to_vec(), b"v1".to_vec()]);
let (cmd, args) = parse_one(&mut s, b"*1\r\n$4\r\nPING\r\n");
assert_eq!(cmd, Some(RespCommand::Ping));
assert!(args.is_empty());
let (cmd, args) = parse_one(&mut s, b"*2\r\n$4\r\nPING\r\n$3\r\nhey\r\n");
assert_eq!(cmd, Some(RespCommand::Ping));
assert_eq!(args, vec![b"hey".to_vec()]);
let (cmd, args) = parse_one(&mut s, b"*2\r\n$3\r\nget\r\n$3\r\nfoo\r\n");
assert_eq!(cmd, Some(RespCommand::Get));
assert_eq!(args, vec![b"foo".to_vec()]);
let (cmd, _) = parse_one(&mut s, b"PING\r\n");
assert_eq!(cmd, Some(RespCommand::Ping));
let (cmd, args) = parse_one(&mut s, b"*2\r\n$6\r\nEXISTS\r\n$1\r\nk\r\n");
assert_eq!(cmd, Some(RespCommand::Exists));
assert_eq!(args, vec![b"k".to_vec()]);
}
#[test]
fn scalar_fast_path_varargs_table() {
let mut s = RespServerSession::default();
let (cmd, args) = parse_one(&mut s, b"*3\r\n$7\r\nPUBLISH\r\n$3\r\nfoo\r\n$1\r\nb\r\n");
assert_eq!(cmd, Some(RespCommand::Publish));
assert_eq!(args, vec![b"foo".to_vec(), b"b".to_vec()]);
let (cmd, _) = parse_one(&mut s, b"*3\r\n$8\r\nSPUBLISH\r\n$1\r\na\r\n$1\r\nb\r\n");
assert_eq!(cmd, Some(RespCommand::Spublish));
let (cmd, _) = parse_one(
&mut s,
b"*4\r\n$8\r\nSETRANGE\r\n$1\r\nk\r\n$1\r\n0\r\n$1\r\nv\r\n",
);
assert_eq!(cmd, Some(RespCommand::Setrange));
let (cmd, _) = parse_one(
&mut s,
b"*4\r\n$8\r\nGETRANGE\r\n$1\r\nk\r\n$1\r\n0\r\n$1\r\n9\r\n",
);
assert_eq!(cmd, Some(RespCommand::Getrange));
let (cmd, args) = parse_one(
&mut s,
b"*4\r\n$3\r\nSET\r\n$1\r\nk\r\n$1\r\nv\r\n$2\r\nNX\r\n",
);
assert_eq!(cmd, Some(RespCommand::Setexnx));
assert_eq!(args, vec![b"k".to_vec(), b"v".to_vec(), b"NX".to_vec()]);
let (cmd, args) = parse_one(&mut s, b"*3\r\n$5\r\nGETEX\r\n$1\r\nk\r\n$7\r\nPERSIST\r\n");
assert_eq!(cmd, Some(RespCommand::Getex));
assert_eq!(args, vec![b"k".to_vec(), b"PERSIST".to_vec()]);
let (cmd, _) = parse_one(&mut s, b"*3\r\n$6\r\nEXPIRE\r\n$1\r\nk\r\n$3\r\n100\r\n");
assert_eq!(cmd, Some(RespCommand::Expire));
let (cmd, _) = parse_one(&mut s, b"*3\r\n$7\r\nPEXPIRE\r\n$1\r\nk\r\n$3\r\n100\r\n");
assert_eq!(cmd, Some(RespCommand::Pexpire));
}
#[test]
fn mru_cache_captures_hash_lookup_commands() {
let mut s = RespServerSession::default();
let (cmd, _) = parse_one(&mut s, b"*2\r\n$5\r\nLPUSH\r\n$1\r\nk\r\n");
assert_eq!(cmd, Some(RespCommand::Lpush));
let (cmd, args) = parse_one(&mut s, b"*2\r\n$5\r\nLPUSH\r\n$1\r\nz\r\n");
assert_eq!(cmd, Some(RespCommand::Lpush));
assert_eq!(args, vec![b"z".to_vec()]);
let (cmd, _) = parse_one(
&mut s,
b"*4\r\n$4\r\nHSET\r\n$1\r\nk\r\n$1\r\nf\r\n$1\r\nv\r\n",
);
assert_eq!(cmd, Some(RespCommand::Hset));
let (cmd, _) = parse_one(&mut s, b"*2\r\n$5\r\nLPUSH\r\n$1\r\nq\r\n");
assert_eq!(cmd, Some(RespCommand::Lpush));
}
#[test]
fn subcommand_dispatch_and_unknown() {
let mut s = RespServerSession::default();
let (cmd, args) = parse_one(&mut s, b"*2\r\n$6\r\nclient\r\n$2\r\nid\r\n");
assert_eq!(cmd, Some(RespCommand::ClientId));
assert!(args.is_empty());
let (cmd, _) = parse_one(&mut s, b"*2\r\n$6\r\nCONFIG\r\n$3\r\nGET\r\n");
assert_eq!(cmd, Some(RespCommand::ConfigGet));
let (cmd, _) = parse_one(&mut s, b"*2\r\n$6\r\nCLIENT\r\n$4\r\nNOPE\r\n");
assert_eq!(cmd, Some(RespCommand::Invalid));
let out = s.take_output();
assert_eq!(out, b"-ERR unknown subcommand 'NOPE'.\r\n");
let (cmd, _) = parse_one(&mut s, b"*2\r\n$7\r\nCLUSTER\r\n$4\r\nNOPE\r\n");
assert_eq!(cmd, Some(RespCommand::Invalid));
let out = s.take_output();
assert_eq!(out, b"-ERR unknown subcommand 'NOPE'. Try CLUSTER HELP\r\n");
let (cmd, _) = parse_one(&mut s, b"*1\r\n$6\r\nCLIENT\r\n");
assert_eq!(cmd, Some(RespCommand::Invalid));
let out = s.take_output();
assert_eq!(
out,
b"-ERR wrong number of arguments for 'CLIENT' command\r\n"
);
let (cmd, _) = parse_one(&mut s, b"*1\r\n$4\r\nnope\r\n");
assert_eq!(cmd, Some(RespCommand::Invalid));
let out = s.take_output();
assert_eq!(out, b"-ERR unknown command\r\n");
}
#[test]
fn every_table_entry_round_trips_via_parse() {
let mut s = RespServerSession::default();
for (name, cmd, has_subcommands) in PRIMARY_TABLE {
if *has_subcommands {
continue;
}
let frame = format!("*2\r\n${}\r\n{name}\r\n$1\r\nk\r\n", name.len());
let (parsed, _) = parse_one(&mut s, frame.as_bytes());
assert_eq!(parsed, Some(*cmd), "主表 {name} 解析不可达");
}
let subtables: [ParentSubtable; 13] = [
("CLIENT", RespCommand::Client, CLIENT_SUBTABLE),
("CONFIG", RespCommand::Config, CONFIG_SUBTABLE),
("COMMAND", RespCommand::Command, COMMAND_SUBTABLE),
("ACL", RespCommand::Acl, ACL_SUBTABLE),
("SCRIPT", RespCommand::Script, SCRIPT_SUBTABLE),
("PUBSUB", RespCommand::Pubsub, PUBSUB_SUBTABLE),
("LATENCY", RespCommand::Latency, LATENCY_SUBTABLE),
("SLOWLOG", RespCommand::Slowlog, SLOWLOG_SUBTABLE),
("MODULE", RespCommand::Module, MODULE_SUBTABLE),
("MEMORY", RespCommand::Memory, MEMORY_SUBTABLE),
("OBJECT", RespCommand::Object, OBJECT_SUBTABLE),
("CLUSTER", RespCommand::Cluster, CLUSTER_SUBTABLE),
("BITOP", RespCommand::Bitop, BITOP_SUBTABLE),
];
for (parent_name, parent_cmd, table) in subtables {
let entry = lookup_primary(parent_name.as_bytes());
assert_eq!(
entry,
Some((parent_cmd, true)),
"主表缺父命令 {parent_name}"
);
for (sub_name, sub_cmd) in table {
let frame = format!(
"*2\r\n${}\r\n{parent_name}\r\n${}\r\n{sub_name}\r\n",
parent_name.len(),
sub_name.len()
);
let (parsed, _) = parse_one(&mut s, frame.as_bytes());
assert_eq!(
parsed,
Some(*sub_cmd),
"{parent_name} {sub_name} 解析不可达"
);
}
}
}
#[test]
fn primary_table_order_and_full_dispatch() {
let mut s = RespServerSession::default();
let (cmd, args) = parse_one(&mut s, b"*3\r\n$4\r\nHDEL\r\n$1\r\nk\r\n$1\r\nf\r\n");
assert_eq!(cmd, Some(RespCommand::Hdel));
assert_eq!(args, vec![b"k".to_vec(), b"f".to_vec()]);
for (frame, expect) in [
(
&b"*3\r\n$8\r\nEXPIREAT\r\n$1\r\nk\r\n$1\r\n1\r\n"[..],
RespCommand::Expireat,
),
(
&b"*3\r\n$8\r\nZREVRANK\r\n$1\r\nk\r\n$1\r\nm\r\n"[..],
RespCommand::Zrevrank,
),
(
&b"*2\r\n$6\r\nSUBSTR\r\n$1\r\nk\r\n"[..],
RespCommand::Substr,
),
(
&b"*5\r\n$5\r\nLMOVE\r\n$1\r\na\r\n$1\r\nb\r\n$4\r\nLEFT\r\n$5\r\nRIGHT\r\n"[..],
RespCommand::Lmove,
),
] {
let (cmd, _) = parse_one(&mut s, frame);
assert_eq!(cmd, Some(expect), "{frame:?}");
}
for (frame, expect) in [
(
&b"*2\r\n$3\r\nACL\r\n$3\r\nCAT\r\n"[..],
RespCommand::AclCat,
),
(
&b"*2\r\n$6\r\nMODULE\r\n$6\r\nLOADCS\r\n"[..],
RespCommand::ModuleLoadcs,
),
(
&b"*2\r\n$6\r\nMEMORY\r\n$5\r\nUSAGE\r\n$1\r\nk\r\n"[..],
RespCommand::MemoryUsage,
),
(
&b"*3\r\n$6\r\nOBJECT\r\n$8\r\nENCODING\r\n$1\r\nk\r\n"[..],
RespCommand::ObjectEncoding,
),
(
&b"*2\r\n$7\r\nCLUSTER\r\n$4\r\nMYID\r\n"[..],
RespCommand::ClusterMyid,
),
(
&b"*3\r\n$7\r\nCLUSTER\r\n$16\r\nSET-CONFIG-EPOCH\r\n$1\r\n0\r\n"[..],
RespCommand::ClusterSetconfigepoch,
),
] {
let (cmd, _) = parse_one(&mut s, frame);
assert_eq!(cmd, Some(expect), "{frame:?}");
}
let (cmd, args) = parse_one(
&mut s,
b"*4\r\n$5\r\nBITOP\r\n$3\r\nNOT\r\n$3\r\ndst\r\n$3\r\nsrc\r\n",
);
assert_eq!(cmd, Some(RespCommand::BitopNot));
assert_eq!(args, vec![b"dst".to_vec(), b"src".to_vec()]);
}
#[test]
fn incomplete_command_returns_none() {
let mut s = RespServerSession::default();
let (cmd, _) = parse_one(&mut s, b"*2\r\n$3\r\nSET\r\n$1\r\nk");
assert_eq!(cmd, None);
let (cmd, _) = parse_one(&mut s, b"*2\r");
assert_eq!(cmd, None);
}
#[test]
fn parse_resp_command_buffer_restores_state() {
let mut s = RespServerSession::default();
let read_head_before = s.read_head;
let cmd = s.parse_resp_command_buffer(b"*2\r\n$6\r\nclient\r\n$4\r\ninfo\r\n");
assert_eq!(cmd, Some(RespCommand::ClientInfo));
assert_eq!(s.read_head, read_head_before);
let (ok, cmd) = s.fuzz_parse_command_buffer(b"*1\r\n$3\r\nGE");
assert!(!ok);
assert_eq!(cmd, RespCommand::Invalid);
let (ok, cmd) = s.fuzz_parse_command_buffer(b"*1\r\n$3\r\nGET\r\n");
assert!(ok);
assert_eq!(cmd, RespCommand::Get);
}
#[test]
fn aof_commit_mode_marks_dependent_commands() {
let mut s = RespServerSession::default();
s.handle_aof_commit_mode(RespCommand::Set);
assert!(s.wait_for_aof_blocking);
s.handle_aof_commit_mode(RespCommand::Ping);
assert!(!s.wait_for_aof_blocking);
s.handle_aof_commit_mode(RespCommand::Set);
s.write_direct_large(b"+OK\r\n");
s.handle_aof_commit_mode(RespCommand::Ping);
assert!(s.wait_for_aof_blocking);
s.output.clear();
s.handle_aof_commit_mode(RespCommand::Ping);
assert!(!s.wait_for_aof_blocking);
s.txn_state = TxnState::Started;
s.handle_aof_commit_mode(RespCommand::Set);
assert!(!s.wait_for_aof_blocking);
}
#[test]
fn inline_malformed_skips_line() {
let mut s = RespServerSession::default();
s.recv_buffer
.extend_from_slice(b"garbage\r\n*1\r\n$4\r\nPING\r\n");
s.bytes_read = s.recv_buffer.len();
s.read_head = 0;
let consumed = s.try_consume_messages(b"garbage\r\n*1\r\n$4\r\nPING\r\n");
assert!(consumed.is_some());
}
#[test]
fn is_aof_independent_matches_csharp_set() {
let csharp_set: &[RespCommand] = &[
RespCommand::Async,
RespCommand::Ping,
RespCommand::Select,
RespCommand::Swapdb,
RespCommand::Echo,
RespCommand::Monitor,
RespCommand::ModuleLoadcs,
RespCommand::Registercs,
RespCommand::Info,
RespCommand::Time,
RespCommand::Lastsave,
RespCommand::AclCat,
RespCommand::AclDeluser,
RespCommand::AclGenpass,
RespCommand::AclGetuser,
RespCommand::AclList,
RespCommand::AclLoad,
RespCommand::AclSave,
RespCommand::AclSetuser,
RespCommand::AclUsers,
RespCommand::AclWhoami,
RespCommand::ClientId,
RespCommand::ClientInfo,
RespCommand::ClientList,
RespCommand::ClientKill,
RespCommand::ClientGetname,
RespCommand::ClientSetname,
RespCommand::ClientSetinfo,
RespCommand::ClientUnblock,
RespCommand::Command,
RespCommand::CommandCount,
RespCommand::CommandDocs,
RespCommand::CommandInfo,
RespCommand::CommandGetkeys,
RespCommand::CommandGetkeysandflags,
RespCommand::MemoryUsage,
RespCommand::ConfigGet,
RespCommand::ConfigRewrite,
RespCommand::ConfigSet,
RespCommand::LatencyHelp,
RespCommand::LatencyHistogram,
RespCommand::LatencyReset,
RespCommand::SlowlogHelp,
RespCommand::SlowlogLen,
RespCommand::SlowlogGet,
RespCommand::SlowlogReset,
RespCommand::Multi,
];
assert_eq!(AOF_INDEPENDENT_COMMANDS.len(), csharp_set.len());
for cmd in csharp_set {
assert!(is_aof_independent(*cmd), "{cmd:?} 应为 AOF 无关");
}
for cmd in AOF_INDEPENDENT_COMMANDS {
assert!(csharp_set.contains(cmd), "{cmd:?} 不在 C# 集内");
}
assert!(!is_aof_independent(RespCommand::Set));
assert!(!is_aof_independent(RespCommand::Get));
assert!(!is_aof_independent(RespCommand::Invalid));
assert!(!is_aof_independent(RespCommand::None));
assert!(!is_aof_independent(RespCommand::Latency));
assert!(!is_aof_independent(RespCommand::Slowlog));
}
#[test]
fn ordered_tables_strictly_ascending() {
for pair in PRIMARY_TABLE.windows(2) {
assert!(
pair[0].0.as_bytes() < pair[1].0.as_bytes(),
"PRIMARY_TABLE 乱序: {} >= {}",
pair[0].0,
pair[1].0
);
}
for table in [
CLIENT_SUBTABLE,
CONFIG_SUBTABLE,
COMMAND_SUBTABLE,
ACL_SUBTABLE,
SCRIPT_SUBTABLE,
PUBSUB_SUBTABLE,
LATENCY_SUBTABLE,
SLOWLOG_SUBTABLE,
MODULE_SUBTABLE,
MEMORY_SUBTABLE,
OBJECT_SUBTABLE,
CLUSTER_SUBTABLE,
BITOP_SUBTABLE,
] {
for pair in table.windows(2) {
assert!(
pair[0].0.as_bytes() < pair[1].0.as_bytes(),
"子命令表乱序: {} >= {}",
pair[0].0,
pair[1].0
);
}
}
}
#[test]
fn parent_commands_have_subtable_dispatch() {
let parents: Vec<(&str, RespCommand)> = PRIMARY_TABLE
.iter()
.filter(|(.., has_sub)| *has_sub)
.map(|(name, cmd, _)| (*name, *cmd))
.collect();
let expected: [&str; 13] = [
"ACL", "BITOP", "CLIENT", "CLUSTER", "COMMAND", "CONFIG", "LATENCY", "MEMORY", "MODULE",
"OBJECT", "PUBSUB", "SCRIPT", "SLOWLOG",
];
let mut names: Vec<&str> = parents.iter().map(|(name, _)| *name).collect();
names.sort_unstable();
assert_eq!(names, expected);
for (name, cmd) in parents {
let table: &[(&str, RespCommand)] = match cmd {
RespCommand::Client => CLIENT_SUBTABLE,
RespCommand::Config => CONFIG_SUBTABLE,
RespCommand::Command => COMMAND_SUBTABLE,
RespCommand::Acl => ACL_SUBTABLE,
RespCommand::Script => SCRIPT_SUBTABLE,
RespCommand::Pubsub => PUBSUB_SUBTABLE,
RespCommand::Latency => LATENCY_SUBTABLE,
RespCommand::Slowlog => SLOWLOG_SUBTABLE,
RespCommand::Module => MODULE_SUBTABLE,
RespCommand::Memory => MEMORY_SUBTABLE,
RespCommand::Object => OBJECT_SUBTABLE,
RespCommand::Cluster => CLUSTER_SUBTABLE,
RespCommand::Bitop => BITOP_SUBTABLE,
_ => &[],
};
assert!(!table.is_empty(), "父命令 {name} 缺子表分派");
assert_eq!(
lookup_subcommand(cmd, table[0].0.as_bytes()),
Some(table[0].1)
);
let last = table[table.len() - 1];
assert_eq!(lookup_subcommand(cmd, last.0.as_bytes()), Some(last.1));
}
}
}