use kevy_resp::{
ArgvView, RespVersion, encode_array_len, encode_bulk, encode_double, encode_error,
encode_integer,
};
use kevy_rt::NotifyClass;
use kevy_store::{ScoreBound, Store, StoreError};
pub(crate) fn upper_verb<'a>(name: &[u8], buf: &'a mut [u8; 32]) -> &'a [u8] {
let n = name.len();
if n <= buf.len() {
buf[..n].copy_from_slice(name);
buf[..n].make_ascii_uppercase();
&buf[..n]
} else {
&buf[..0]
}
}
pub(crate) fn wrong_args(out: &mut Vec<u8>, cmd: &str) {
encode_error(
out,
&format!("ERR wrong number of arguments for '{cmd}' command"),
);
}
pub(crate) fn cmd_hello(out: &mut Vec<u8>) {
encode_array_len(out, 14);
encode_bulk(out, b"server");
encode_bulk(out, b"kevy");
encode_bulk(out, b"version");
encode_bulk(out, env!("CARGO_PKG_VERSION").as_bytes());
encode_bulk(out, b"proto");
encode_integer(out, 2);
encode_bulk(out, b"id");
encode_integer(out, 0);
encode_bulk(out, b"mode");
encode_bulk(out, b"standalone");
encode_bulk(out, b"role");
encode_bulk(out, b"master");
encode_bulk(out, b"modules");
encode_array_len(out, 0);
}
pub(crate) const ERR_NOT_INT: &str = "ERR value is not an integer or out of range";
pub(crate) const WRONGTYPE: &str =
"WRONGTYPE Operation against a key holding the wrong kind of value";
pub(crate) const OOM_ERR: &str =
"OOM command not allowed when used memory > 'maxmemory'.";
pub(crate) fn is_write_verb(cmd: &[u8]) -> bool {
matches!(
cmd,
b"SET"
| b"SETNX"
| b"SETEX"
| b"PSETEX"
| b"GETSET"
| b"GETDEL"
| b"INCRBYFLOAT"
| b"DEL"
| b"UNLINK"
| b"INCR"
| b"DECR"
| b"INCRBY"
| b"DECRBY"
| b"APPEND"
| b"EXPIRE"
| b"PEXPIRE"
| b"EXPIREAT"
| b"PEXPIREAT"
| b"PERSIST"
| b"FLUSHDB"
| b"FLUSHALL"
| b"HSET"
| b"HSETNX"
| b"HMSET"
| b"HDEL"
| b"HINCRBY"
| b"LPUSH"
| b"RPUSH"
| b"LPOP"
| b"RPOP"
| b"LSET"
| b"LREM"
| b"LTRIM"
| b"RPOPLPUSH"
| b"BRPOPLPUSH"
| b"LMOVE"
| b"SADD"
| b"SREM"
| b"SPOP"
| b"ZADD"
| b"ZREM"
| b"ZINCRBY"
| b"ZPOPMIN"
| b"BZPOPMIN"
| b"ZREMRANGEBYRANK"
| b"ZREMRANGEBYSCORE"
| b"GEOADD"
| b"GEOSEARCHSTORE"
| b"GEORADIUS"
| b"GEORADIUSBYMEMBER"
| b"XADD"
| b"XDEL"
| b"XTRIM"
| b"XSETID"
| b"XGROUP"
| b"XREADGROUP"
| b"XACK"
| b"XCLAIM"
| b"XAUTOCLAIM"
| b"MSET"
| b"EVAL" | b"EVALSHA"
)
}
pub(crate) fn notify_class_for_verb(cmd: &[u8]) -> Option<NotifyClass> {
Some(match cmd {
b"SET" | b"SETNX" | b"SETEX" | b"PSETEX" | b"GETSET" | b"GETDEL"
| b"APPEND" | b"INCR" | b"DECR" | b"INCRBY" | b"DECRBY" | b"INCRBYFLOAT" => {
NotifyClass::String
}
b"HSET" | b"HSETNX" | b"HMSET" | b"HDEL" | b"HINCRBY" => NotifyClass::Hash,
b"LPUSH" | b"RPUSH" | b"LPOP" | b"RPOP" | b"LSET" | b"LREM" | b"LTRIM"
| b"RPOPLPUSH" | b"LMOVE" => NotifyClass::List,
b"SADD" | b"SREM" | b"SPOP" => NotifyClass::Set,
b"ZADD" | b"ZREM" | b"ZINCRBY" | b"ZPOPMIN" | b"ZREMRANGEBYRANK"
| b"ZREMRANGEBYSCORE" | b"GEOADD" => NotifyClass::Zset,
b"XADD" | b"XDEL" | b"XTRIM" | b"XSETID" | b"XGROUP" | b"XACK" | b"XCLAIM"
| b"XAUTOCLAIM" | b"XREADGROUP" => NotifyClass::Stream,
b"DEL" | b"UNLINK" | b"EXPIRE" | b"PEXPIRE" | b"PERSIST" => NotifyClass::Generic,
_ => return None,
})
}
pub(crate) fn is_growing_write_verb(cmd: &[u8]) -> bool {
matches!(
cmd,
b"SET"
| b"SETNX"
| b"SETEX"
| b"PSETEX"
| b"GETSET"
| b"INCRBYFLOAT"
| b"INCR"
| b"DECR"
| b"INCRBY"
| b"DECRBY"
| b"APPEND"
| b"HSET"
| b"HSETNX"
| b"HMSET"
| b"HINCRBY"
| b"LPUSH"
| b"RPUSH"
| b"RPOPLPUSH"
| b"BRPOPLPUSH"
| b"LMOVE"
| b"LSET"
| b"SADD"
| b"ZADD"
| b"ZINCRBY"
| b"GEOADD"
| b"GEOSEARCHSTORE"
| b"GEORADIUS"
| b"GEORADIUSBYMEMBER"
| b"XADD"
| b"XGROUP"
| b"XREADGROUP"
| b"XCLAIM"
| b"XAUTOCLAIM"
| b"MSET"
)
}
pub(crate) fn store_err(out: &mut Vec<u8>, e: StoreError) {
let msg = match e {
StoreError::WrongType => WRONGTYPE,
StoreError::NotInteger => ERR_NOT_INT,
StoreError::Overflow => "ERR increment or decrement would overflow",
StoreError::OutOfRange => "ERR index out of range",
StoreError::NoSuchKey => "ERR no such key",
StoreError::NotFloat => "ERR value is not a valid float",
StoreError::OutOfMemory => OOM_ERR,
};
encode_error(out, msg);
}
pub(crate) fn emit_int_result(res: Result<i64, StoreError>, out: &mut Vec<u8>) {
match res {
Ok(n) => encode_integer(out, n),
Err(e) => store_err(out, e),
}
}
pub(crate) fn emit_bulk_array(res: Result<Vec<Vec<u8>>, StoreError>, out: &mut Vec<u8>) {
match res {
Ok(items) => {
encode_array_len(out, items.len() as i64);
for it in &items {
encode_bulk(out, it);
}
}
Err(e) => store_err(out, e),
}
}
pub(crate) fn cmd_hset<A: ArgvView + ?Sized>(store: &mut Store, args: &A, out: &mut Vec<u8>) {
if args.len() < 4 || !args.len().is_multiple_of(2) {
return wrong_args(out, "hset");
}
let pairs: Vec<(&[u8], &[u8])> = (2..args.len())
.step_by(2)
.map(|i| (&args[i], &args[i + 1]))
.collect();
emit_int_result(store.hset_borrowed(&args[1], &pairs).map(|n| n as i64), out);
}
pub(crate) fn cmd_zadd<A: ArgvView + ?Sized>(store: &mut Store, args: &A, out: &mut Vec<u8>) {
if args.len() < 4 || !(args.len() - 2).is_multiple_of(2) {
return wrong_args(out, "zadd");
}
let mut pairs: Vec<(f64, &[u8])> = Vec::with_capacity((args.len() - 2) / 2);
let mut i = 2;
while i < args.len() {
let Some(score) = arg_f64(&args[i]) else {
return encode_error(out, "ERR value is not a valid float");
};
pairs.push((score, &args[i + 1]));
i += 2;
}
emit_int_result(store.zadd_borrowed(&args[1], &pairs).map(|n| n as i64), out);
}
pub(crate) fn cmd_zrange<A: ArgvView + ?Sized>(
store: &mut Store,
args: &A,
out: &mut Vec<u8>,
proto: RespVersion,
) {
if args.len() < 4 || args.len() > 5 {
return wrong_args(out, "zrange");
}
let withscores = args.len() == 5;
if withscores && !args[4].eq_ignore_ascii_case(b"WITHSCORES") {
return encode_error(out, "ERR syntax error");
}
let (Some(s), Some(e)) = (arg_i64(&args[2]), arg_i64(&args[3])) else {
return encode_error(out, ERR_NOT_INT);
};
emit_zrange(store.zrange(&args[1], s, e), withscores, proto, out);
}
pub(crate) fn cmd_zrangebyscore<A: ArgvView + ?Sized>(
store: &mut Store,
args: &A,
out: &mut Vec<u8>,
proto: RespVersion,
) {
if args.len() < 4 {
return wrong_args(out, "zrangebyscore");
}
let (Some(min), Some(max)) = (parse_score_bound(&args[2]), parse_score_bound(&args[3])) else {
return encode_error(out, "ERR min or max is not a float");
};
let mut withscores = false;
let mut limit: Option<(i64, i64)> = None;
let mut i = 4;
while i < args.len() {
let tok = &args[i];
if tok.eq_ignore_ascii_case(b"WITHSCORES") {
if withscores {
return encode_error(out, "ERR syntax error");
}
withscores = true;
i += 1;
} else if tok.eq_ignore_ascii_case(b"LIMIT") {
if limit.is_some() || i + 2 >= args.len() {
return encode_error(out, "ERR syntax error");
}
let Some(off) = std::str::from_utf8(&args[i + 1])
.ok()
.and_then(|s| s.parse::<i64>().ok())
else {
return encode_error(out, ERR_NOT_INT);
};
let Some(cnt) = std::str::from_utf8(&args[i + 2])
.ok()
.and_then(|s| s.parse::<i64>().ok())
else {
return encode_error(out, ERR_NOT_INT);
};
limit = Some((off, cnt));
i += 3;
} else {
return encode_error(out, "ERR syntax error");
}
}
let res = store.zrange_by_score(&args[1], min, max);
match res {
Err(e) => store_err(out, e),
Ok(mut items) => {
if let Some((off, cnt)) = limit {
let start = off.max(0) as usize;
if start >= items.len() {
items.clear();
} else if cnt < 0 {
items.drain(..start);
} else {
let end = (start + cnt as usize).min(items.len());
items = items[start..end].to_vec();
}
}
emit_zrange(Ok(items), withscores, proto, out);
}
}
}
pub(crate) fn emit_zrange(
res: Result<Vec<(Vec<u8>, f64)>, StoreError>,
withscores: bool,
proto: RespVersion,
out: &mut Vec<u8>,
) {
match res {
Err(e) => store_err(out, e),
Ok(items) => match (withscores, proto) {
(false, _) => {
encode_array_len(out, items.len() as i64);
for (m, _) in &items {
encode_bulk(out, m);
}
}
(true, RespVersion::V2) => {
encode_array_len(out, (items.len() * 2) as i64);
for (m, sc) in &items {
encode_bulk(out, m);
encode_bulk(out, &fmt_score(*sc));
}
}
(true, RespVersion::V3) => {
encode_array_len(out, items.len() as i64);
for (m, sc) in &items {
encode_array_len(out, 2);
encode_bulk(out, m);
encode_double(out, *sc);
}
}
},
}
}
pub(crate) fn arg_f64(b: &[u8]) -> Option<f64> {
let s = std::str::from_utf8(b).ok()?.trim();
let f: f64 = match s.to_ascii_lowercase().as_str() {
"inf" | "+inf" | "infinity" | "+infinity" => f64::INFINITY,
"-inf" | "-infinity" => f64::NEG_INFINITY,
_ => s.parse().ok()?,
};
if f.is_nan() { None } else { Some(f) }
}
pub(crate) fn parse_score_bound(b: &[u8]) -> Option<ScoreBound> {
match b.strip_prefix(b"(") {
Some(rest) => Some(ScoreBound {
value: arg_f64(rest)?,
exclusive: true,
}),
None => Some(ScoreBound {
value: arg_f64(b)?,
exclusive: false,
}),
}
}
pub(crate) fn fmt_score(s: f64) -> Vec<u8> {
if s.is_infinite() {
return if s > 0.0 {
b"inf".to_vec()
} else {
b"-inf".to_vec()
};
}
#[allow(clippy::float_cmp)]
let is_integer_valued = s == s.trunc();
if is_integer_valued && s.abs() < 1e17 {
return (s as i64).to_string().into_bytes();
}
format!("{s}").into_bytes()
}
pub(crate) fn rest_borrowed<'a, A: ArgvView + ?Sized>(
args: &'a A,
from: usize,
) -> Vec<&'a [u8]> {
(from..args.len()).map(|i| &args[i]).collect()
}
pub(crate) fn arg_i64(b: &[u8]) -> Option<i64> {
std::str::from_utf8(b).ok()?.parse::<i64>().ok()
}
pub(crate) fn scan_pattern<A: ArgvView + ?Sized>(args: &A) -> Option<Vec<u8>> {
let mut i = 2;
while i + 1 < args.len() {
if args[i].eq_ignore_ascii_case(b"MATCH") {
return Some(args[i + 1].to_vec());
}
i += 2;
}
None
}
pub(crate) use crate::cmd_data::*;