use crate::cmd::{
OOM_ERR, cmd_expire, cmd_expireat, cmd_hello, cmd_set, cmd_spop_rand, cmd_ttl, emit_bulk_array,
emit_int_result, is_growing_write_verb, rest_borrowed, store_err, upper_verb, wrong_args,
};
use crate::state::Ctx;
use kevy_resp::{
ArgvView, encode_bulk, encode_error, encode_integer, encode_null_bulk, encode_simple_string,
};
use kevy_store::Store;
pub(crate) fn dispatch<A: ArgvView + ?Sized>(
ctx: &Ctx<'_>,
store: &mut Store,
args: &A,
) -> Vec<u8> {
let mut out = Vec::new();
dispatch_into(ctx, store, args, &mut out);
out
}
pub(crate) fn dispatch_into<A: ArgvView + ?Sized>(
ctx: &Ctx<'_>,
store: &mut Store,
args: &A,
out: &mut Vec<u8>,
) {
dispatch_with_proto(ctx, store, args, out, false);
}
pub(crate) fn dispatch_into_resp3<A: ArgvView + ?Sized>(
ctx: &Ctx<'_>,
store: &mut Store,
args: &A,
out: &mut Vec<u8>,
) {
dispatch_with_proto(ctx, store, args, out, true);
}
fn dispatch_with_proto<A: ArgvView + ?Sized>(
ctx: &Ctx<'_>,
store: &mut Store,
args: &A,
out: &mut Vec<u8>,
proto_v3: bool,
) {
let Some(name) = args.first() else {
encode_error(out, "ERR empty command");
return;
};
let mut buf = [0u8; 32];
let cmd = upper_verb(name, &mut buf);
if crate::cmd::is_write_verb(cmd)
&& ctx.shard.gate_bits(ctx.state) & crate::state::SCOPE_ACTIVE != 0
&& let Some(key) = args.get(1)
&& let Some(redirect) = ctx.state.route_write(key, ctx.shard)
{
match redirect {
crate::state::WriteRedirect::Misdirected(addr) => {
crate::state::encode_misdirected(out, &addr);
}
crate::state::WriteRedirect::Quiesced { to_addr } => {
crate::state::encode_quiesced(out, &to_addr);
}
}
return;
}
match cmd {
b"GET" => {
if args.len() == 2 {
match store.get(&args[1]) {
Ok(Some(v)) => encode_bulk(out, &v),
Ok(None) => encode_null_bulk(out),
Err(e) => store_err(out, e),
}
} else {
wrong_args(out, "get");
}
return;
}
b"SET" => {
if store.maxmemory() > 0 {
if store.precheck_for_write().is_err() {
encode_error(out, OOM_ERR);
return;
}
cmd_set(store, args, out);
store.try_evict_after_write();
} else {
cmd_set(store, args, out);
}
store.try_demote_after_write();
return;
}
_ => {}
}
let is_grow = is_growing_write_verb(cmd);
if store.maxmemory() > 0 && is_grow && store.precheck_for_write().is_err() {
encode_error(out, OOM_ERR);
return;
}
let handled = (proto_v3
&& crate::dispatch_resp3::try_resp3_overrides(ctx, cmd, store, args, out))
|| dispatch_conn(ctx, cmd, store, args, out)
|| crate::ops::dispatch_ops(ctx, cmd, store, args, out)
|| crate::dispatch_strings::dispatch_string(cmd, store, args, out)
|| crate::dispatch_bitmap::dispatch_bitmap(cmd, store, args, out)
|| crate::dispatch_collections::dispatch_hash(cmd, store, args, out)
|| crate::dispatch_collections::dispatch_list(cmd, store, args, out)
|| dispatch_set(cmd, store, args, out)
|| crate::dispatch_collections::dispatch_zset(cmd, store, args, out)
|| crate::dispatch_geo::dispatch_geo(cmd, store, args, out)
|| crate::dispatch_stream::dispatch_stream(cmd, store, args, out)
|| crate::cmd_lua::dispatch_lua(ctx, cmd, store, args, out)
|| dispatch_generic(cmd, store, args, out)
|| crate::dispatch_replay::dispatch_multikey_stub(cmd, store, args, out);
if !handled {
crate::cmd::unhandled_verb(out, name, args.len());
return;
}
if is_grow && store.maxmemory() > 0 {
store.try_evict_after_write();
}
if is_grow {
store.try_demote_after_write();
}
}
fn cmd_time(out: &mut Vec<u8>) {
let now =
std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap_or_default();
kevy_resp::encode_array_len(out, 2);
encode_bulk(out, now.as_secs().to_string().as_bytes());
encode_bulk(out, now.subsec_micros().to_string().as_bytes());
}
fn dispatch_conn<A: ArgvView + ?Sized>(
ctx: &Ctx<'_>,
cmd: &[u8],
store: &Store,
args: &A,
out: &mut Vec<u8>,
) -> bool {
match cmd {
b"PING" => match args.len() {
1 => encode_simple_string(out, "PONG"),
2 => encode_bulk(out, &args[1]),
_ => wrong_args(out, "ping"),
},
b"TIME" => cmd_time(out),
b"IDX.CREATE" => crate::cmd_index::cmd_idx_create(ctx, store, args, out),
b"VIEW.CREATE" => crate::cmd_view::cmd_view_create(ctx, args, out),
b"VIEW.DROP" => crate::cmd_view::cmd_view_drop(ctx, args, out),
b"IDX.DROP" => crate::cmd_index::cmd_idx_drop(ctx, args, out),
b"IDX.ADVISE" => crate::cmd_index_advise::cmd_idx_advise(ctx, args, out),
b"TABLE.DECLARE" => crate::cmd_table::cmd_table_declare(ctx, store, args, out),
b"TABLE.ENSURE" => crate::cmd_table::cmd_table_ensure(ctx, store, args, out),
b"TABLE.REPLACE" => crate::cmd_table::cmd_table_replace(ctx, store, args, out),
b"TABLE.DROP" => crate::cmd_table::cmd_table_drop(ctx, args, out),
b"TABLE.LIST" => encode_error(out, "ERR usage: TABLE.LIST"),
b"TABLE.VERIFY" => encode_error(out, "ERR usage: TABLE.VERIFY name"),
b"ECHO" => {
if args.len() == 2 {
encode_bulk(out, &args[1]);
} else {
wrong_args(out, "echo");
}
}
b"COMMAND" => crate::cmd_command::cmd_command(args, out),
b"FAILOVER" => crate::cmd_failover::cmd_failover(ctx, args, out),
b"HELLO" => cmd_hello(out),
b"QUIT" => encode_simple_string(out, "OK"),
b"SELECT" => cmd_select(args, out),
_ => return false,
}
true
}
fn cmd_select<A: ArgvView + ?Sized>(args: &A, out: &mut Vec<u8>) {
if args.len() != 2 {
wrong_args(out, "select");
return;
}
let idx_bytes = &args[1];
let Ok(s) = std::str::from_utf8(idx_bytes) else {
encode_error(out, "ERR value is not an integer or out of range");
return;
};
let parsed: Result<i64, _> = s.parse();
match parsed {
Ok(0) => encode_simple_string(out, "OK"),
Ok(_) => encode_error(
out,
"ERR kevy only supports DB 0 (multi-database support is on the v1.1.0 backlog)",
),
Err(_) => encode_error(out, "ERR value is not an integer or out of range"),
}
}
fn dispatch_set<A: ArgvView + ?Sized>(
cmd: &[u8],
store: &mut Store,
args: &A,
out: &mut Vec<u8>,
) -> bool {
match cmd {
b"SADD" => {
if args.len() < 3 {
wrong_args(out, "sadd");
} else {
emit_int_result(
store.sadd(&args[1], &rest_borrowed(args, 2)).map(|n| n as i64),
out,
);
}
}
b"SREM" => {
if args.len() < 3 {
wrong_args(out, "srem");
} else {
emit_int_result(
store.srem(&args[1], &rest_borrowed(args, 2)).map(|n| n as i64),
out,
);
}
}
b"SCARD" => {
if args.len() == 2 {
emit_int_result(store.scard(&args[1]).map(|n| n as i64), out);
} else {
wrong_args(out, "scard");
}
}
b"SISMEMBER" => {
if args.len() == 3 {
emit_int_result(store.sismember(&args[1], &args[2]).map(i64::from), out);
} else {
wrong_args(out, "sismember");
}
}
b"SMEMBERS" => {
if args.len() == 2 {
emit_bulk_array(store.smembers(&args[1]), out);
} else {
wrong_args(out, "smembers");
}
}
b"SPOP" => cmd_spop_rand(store, args, true, out),
b"SRANDMEMBER" => cmd_spop_rand(store, args, false, out),
_ => return false,
}
true
}
fn dispatch_generic<A: ArgvView + ?Sized>(
cmd: &[u8],
store: &mut Store,
args: &A,
out: &mut Vec<u8>,
) -> bool {
match cmd {
b"DEL" | b"UNLINK" => {
let verb = if cmd == b"DEL" { "del" } else { "unlink" };
if args.len() < 2 {
wrong_args(out, verb);
} else {
encode_integer(out, store.del(&rest_borrowed(args, 1)) as i64);
}
}
b"EXISTS" | b"TOUCH" => {
let verb = if cmd == b"EXISTS" { "exists" } else { "touch" };
if args.len() < 2 {
wrong_args(out, verb);
} else {
encode_integer(out, store.exists(&rest_borrowed(args, 1)) as i64);
}
}
b"EXPIRE" => cmd_expire(store, args, 1000, "expire", out),
b"PEXPIRE" => cmd_expire(store, args, 1, "pexpire", out),
b"EXPIREAT" => cmd_expireat(store, args, 1000, "expireat", out),
b"PEXPIREAT" => cmd_expireat(store, args, 1, "pexpireat", out),
b"TTL" => cmd_ttl(store, args, true, "ttl", out),
b"PTTL" => cmd_ttl(store, args, false, "pttl", out),
b"PERSIST" => {
if args.len() == 2 {
encode_integer(out, i64::from(store.persist(&args[1])));
} else {
wrong_args(out, "persist");
}
}
b"TYPE" => {
if args.len() == 2 {
encode_simple_string(out, store.type_of(&args[1]));
} else {
wrong_args(out, "type");
}
}
b"DBSIZE" => encode_integer(out, store.dbsize() as i64),
b"FLUSHDB" | b"FLUSHALL" => {
store.flushall();
encode_simple_string(out, "OK");
}
_ => return false,
}
true
}