use kevy_rt::{
ArgvView, BlockKind, Commands, NotifyClass, ResolvedCmd, RespVersion, Route, TxnKind,
parse_slowlog_sub,
};
use kevy_store::Store;
use crate::cmd::{self, scan_pattern, upper_verb};
use crate::{
Argv, KevyCommands, cmd_block, cmd_block_serve, cmd_hello, cmd_resolve, config_global,
dispatch, map_appendfsync, map_eviction_policy, ops,
};
impl Commands for KevyCommands {
fn route<A: ArgvView + ?Sized>(&self, args: &A) -> Route {
let Some(name) = args.first() else {
return Route::Local;
};
let mut buf = [0u8; 32];
match upper_verb(name, &mut buf) {
b"HELLO" => Route::Hello,
b"PING" | b"ECHO" | b"QUIT" | b"COMMAND" | b"CONFIG"
| b"INFO" | b"CLUSTER" | b"DEBUG" | b"SHUTDOWN"
| b"CLIENT" | b"SELECT" | b"ROLE"
| b"REPLICAOF" | b"SLAVEOF" => Route::Local,
b"WAIT" => crate::cmd_repl::wait_route(args),
b"REPL.TOKEN" => crate::cmd_repl::token_route(args),
b"REPL.WAIT" => crate::cmd_repl::repl_wait_route(args),
b"DBSIZE" => Route::Dbsize,
b"FLUSHDB" | b"FLUSHALL" => Route::Flush,
b"SAVE" => Route::Save,
b"BGSAVE" => Route::BgSave,
b"BGREWRITEAOF" => Route::RewriteAof,
b"MSET" if args.len() >= 3 && !args.len().is_multiple_of(2) => Route::MSet,
b"MGET" if args.len() >= 2 => Route::MGet,
b"ZINTERSTORE" if args.len() >= 4 => Route::ZAlgebraStore(kevy_rt::ZCombine::ZInter),
b"ZUNIONSTORE" if args.len() >= 4 => Route::ZAlgebraStore(kevy_rt::ZCombine::ZUnion),
b"ZDIFFSTORE" if args.len() >= 4 => Route::ZAlgebraStore(kevy_rt::ZCombine::ZDiff),
b"SINTERSTORE" if args.len() >= 3 => Route::ZAlgebraStore(kevy_rt::ZCombine::SInter),
b"SUNIONSTORE" if args.len() >= 3 => Route::ZAlgebraStore(kevy_rt::ZCombine::SUnion),
b"SDIFFSTORE" if args.len() >= 3 => Route::ZAlgebraStore(kevy_rt::ZCombine::SDiff),
b"ZINTERCARD" if args.len() >= 3 => Route::ZInterCard,
b"IDX.QUERY" if args.len() >= 4 => Route::Extension,
b"IDX.EXPLAIN" if args.len() >= 2 => Route::Extension,
b"IDX.REBUILD" if args.len() == 2 => Route::Extension,
b"IDX.COUNT" if args.len() >= 4 => Route::Extension,
b"IDX.VERIFY" if args.len() == 2 => Route::Extension,
b"IDX.LIST" if args.len() == 1 => Route::Extension,
b"VIEW.QUERY" if args.len() >= 2 => Route::Extension,
b"VIEW.LIST" if args.len() == 1 => Route::Extension,
b"VIEW.VERIFY" if args.len() == 2 => Route::Extension,
b"VIEW.REBUILD" if args.len() == 2 => Route::Extension,
b"VIEW.EXPLAIN" if args.len() == 2 => Route::Extension,
b"PREFIX.STATS" if args.len() == 2 => Route::PrefixStats,
b"PREFIX.DIGEST" if args.len() == 2 => Route::Extension,
b"FEED.READ" if args.len() >= 4 => Route::FeedRead,
b"FEED.TAIL" if args.len() == 2 => Route::FeedTail,
b"FEED.SHARDS" if args.len() == 1 => Route::FeedShards,
b"SINTER" if args.len() >= 2 => Route::SInter,
b"SUNION" if args.len() >= 2 => Route::SUnion,
b"SDIFF" if args.len() >= 2 => Route::SDiff,
b"KEYS" if args.len() == 2 => Route::Keys(Some(args[1].to_vec())),
b"SCAN" if args.len() >= 2 => Route::Scan(scan_pattern(args)),
b"RANDOMKEY" if args.len() == 1 => Route::RandomKey,
b"SUBSCRIBE" if args.len() >= 2 => Route::Subscribe,
b"UNSUBSCRIBE" => Route::Unsubscribe, b"PSUBSCRIBE" if args.len() >= 2 => Route::Psubscribe,
b"PUNSUBSCRIBE" => Route::Punsubscribe, b"PUBLISH" if args.len() == 3 => Route::Publish,
b"WATCH" if args.len() >= 2 => Route::Watch,
b"UNWATCH" => Route::Unwatch,
b"RENAME" => Route::Rename { nx: false },
b"RENAMENX" => Route::Rename { nx: true },
b"EVAL" | b"EVALSHA" | b"EVAL_RO" | b"EVALSHA_RO" => {
if args.len() >= 4 {
let nk = std::str::from_utf8(&args[2])
.ok()
.and_then(|s| s.parse::<i64>().ok())
.unwrap_or(0);
if nk >= 1 && (args.len() as i64) >= 3 + nk {
Route::Single(3)
} else {
Route::Local
}
} else {
Route::Local
}
}
b"SCRIPT" => Route::Local,
b"XREAD" => cmd_block::xread_route(args),
b"XREADGROUP" => cmd_block::xreadgroup_route(args),
b"XGROUP" | b"XINFO" if args.len() >= 3 => Route::Single(2),
b"SLOWLOG" => Route::Slowlog(parse_slowlog_sub(args)),
b"DEL" | b"UNLINK" => {
if args.len() == 2 {
Route::Single(1)
} else {
Route::DelKeys
}
}
b"EXISTS" => {
if args.len() == 2 {
Route::Single(1)
} else {
Route::ExistsKeys
}
}
_ => {
if args.len() >= 2 {
Route::Single(1)
} else {
Route::Local }
}
}
}
fn dispatch<A: ArgvView + ?Sized>(&self, store: &mut Store, args: &A) -> Vec<u8> {
dispatch::dispatch(store, args)
}
fn dispatch_into<A: ArgvView + ?Sized>(&self, store: &mut Store, args: &A, out: &mut Vec<u8>) {
dispatch::dispatch_into(store, args, out);
}
fn dispatch_resp3<A: ArgvView + ?Sized>(&self, store: &mut Store, args: &A) -> Vec<u8> {
let mut out = Vec::with_capacity(64);
dispatch::dispatch_into_resp3(store, args, &mut out);
out
}
fn dispatch_into_resp3<A: ArgvView + ?Sized>(
&self,
store: &mut Store,
args: &A,
out: &mut Vec<u8>,
) {
dispatch::dispatch_into_resp3(store, args, out);
}
fn is_quit<A: ArgvView + ?Sized>(&self, args: &A) -> bool {
args.first()
.is_some_and(|c| c.eq_ignore_ascii_case(b"QUIT"))
}
fn on_shard_init(&self, store: &mut Store) {
let cfg = config_global::get();
store.set_max_memory(
cfg.memory.maxmemory,
map_eviction_policy(cfg.memory.maxmemory_policy),
);
}
fn on_shard_start(&self, shard: usize) {
ops::cluster::set_current_shard(shard);
ops::stats::register_shard(shard);
}
fn on_persist_stats(&self, in_flight: bool, aof_rewrites_total: u64) {
ops::set_persist_stats(in_flight, aof_rewrites_total);
}
fn on_replication_view(
&self,
master_repl_offset: u64,
replicas: Vec<(std::net::Ipv4Addr, u16, u64, Option<u64>)>,
) {
ops::replication::set_replication_view(master_repl_offset, replicas);
crate::elect_integration::set_view_offset(master_repl_offset);
}
fn on_command(&self) {
ops::stats::add_command();
}
fn on_connection(&self) {
ops::stats::add_connection();
}
fn shard_tick_interval_ms(&self) -> u64 {
let cfg = config_global::get();
let hz = cfg.expiry.hz;
if hz == 0 {
0
} else {
(1000 / u64::from(hz)).clamp(1, 10_000)
}
}
fn on_write(&self, store: &mut Store, key: &[u8]) {
crate::index_runtime::on_write(store, key);
crate::view_runtime::on_write(store, key);
}
fn extension_op(&self, store: &mut Store, argv: &[Vec<u8>]) -> Vec<u8> {
if argv.first().is_some_and(|v| v.eq_ignore_ascii_case(b"PREFIX.DIGEST")) {
return crate::cmd_digest::extension_op(store, argv);
}
if argv.first().is_some_and(|v| v.len() > 5 && v[..5].eq_ignore_ascii_case(b"VIEW.")) {
return crate::cmd_view::extension_op(store, argv);
}
crate::cmd_index_query::extension_op(store, argv)
}
fn extension_reduce(&self, argv: &[Vec<u8>], chunks: Vec<Vec<u8>>) -> Vec<u8> {
if argv.first().is_some_and(|v| v.eq_ignore_ascii_case(b"PREFIX.DIGEST")) {
return crate::cmd_digest::extension_reduce(chunks);
}
if argv.first().is_some_and(|v| v.len() > 5 && v[..5].eq_ignore_ascii_case(b"VIEW.")) {
return crate::cmd_view::extension_reduce(argv, chunks);
}
crate::cmd_index_reduce::extension_reduce(argv, chunks)
}
fn write_denied(&self) -> Option<Vec<u8>> {
crate::replica_state::write_denied_reply()
}
fn read_denied(&self) -> Option<Vec<u8>> {
crate::replica_state::read_denied_reply()
}
fn extension_reduce_v3(
&self,
argv: &[Vec<u8>],
chunks: Vec<Vec<u8>>,
proto: kevy_resp::RespVersion,
) -> Vec<u8> {
let reply = self.extension_reduce(argv, chunks);
if proto == kevy_resp::RespVersion::V3 {
return crate::cmd_index_reduce::resp3_upgrade(argv, reply);
}
reply
}
fn on_shard_tick(&self, store: &mut Store) {
crate::index_runtime::on_tick(store);
crate::view_runtime::on_tick(store);
let cfg = config_global::get();
if !crate::replica_state::is_replica() {
let samples = cfg.expiry.sample as usize;
store.tick_expire(samples, 16);
let _ = store.tick_hash_ttl(64);
}
store.set_max_memory(
cfg.memory.maxmemory,
map_eviction_policy(cfg.memory.maxmemory_policy),
);
ops::stats::publish_gauges(store);
ops::stats::sample_ops_if_lead();
}
fn live_runtime_config(&self) -> kevy_rt::LiveRuntimeConfig {
if !config_global::is_initialised() {
return kevy_rt::LiveRuntimeConfig {
promotion_epoch: crate::replica_state::promotion_epoch(),
..kevy_rt::LiveRuntimeConfig::default()
};
}
let cfg = config_global::get();
let hz = cfg.expiry.hz;
let tick_ms = if hz == 0 {
Some(0)
} else {
Some((1000u64 / u64::from(hz)).clamp(1, 10_000))
};
kevy_rt::LiveRuntimeConfig {
appendfsync: Some(map_appendfsync(cfg.persistence.appendfsync)),
auto_aof_rewrite_pct: Some(cfg.persistence.auto_aof_rewrite_percentage),
auto_aof_rewrite_min_size: Some(cfg.persistence.auto_aof_rewrite_min_size),
tick_interval_ms: tick_ms,
notify_flags: Some(kevy_config::parse_notification_flags(
&cfg.notification.notify_keyspace_events,
)),
slowlog_slower_than_micros: Some(cfg.slowlog.slower_than_micros),
slowlog_max_len: Some(cfg.slowlog.max_len),
promotion_epoch: crate::replica_state::promotion_epoch(),
}
}
fn hello_reply<A: ArgvView + ?Sized>(
&self,
args: &A,
current_proto: RespVersion,
) -> (RespVersion, Vec<u8>) {
cmd_hello::kevy_hello_reply(args, current_proto)
}
fn is_write<A: ArgvView + ?Sized>(&self, args: &A) -> bool {
let Some(name) = args.first() else {
return false;
};
let mut buf = [0u8; 32];
cmd::is_write_verb(upper_verb(name, &mut buf))
}
fn notify_class<A: ArgvView + ?Sized>(&self, args: &A) -> Option<NotifyClass> {
let name = args.first()?;
let mut buf = [0u8; 32];
cmd::notify_class_for_verb(upper_verb(name, &mut buf))
}
fn txn_kind<A: ArgvView + ?Sized>(&self, args: &A) -> TxnKind {
let Some(name) = args.first() else {
return TxnKind::Other;
};
let mut buf = [0u8; 32];
match upper_verb(name, &mut buf) {
b"MULTI" => TxnKind::Multi,
b"EXEC" => TxnKind::Exec,
b"DISCARD" => TxnKind::Discard,
b"WATCH" => TxnKind::Watch,
_ => TxnKind::Other,
}
}
fn resolve_block_argv<A: ArgvView + ?Sized>(
&self,
store: &mut Store,
args: &A,
kind: BlockKind,
) -> Argv {
match kind {
BlockKind::XReadBlock => cmd_block::xread_resolve_argv(store, args),
_ => args.to_argv(),
}
}
fn block_serve_argv<A: ArgvView + ?Sized>(
&self,
args: &A,
kind: BlockKind,
key: &[u8],
) -> Argv {
cmd_block_serve::block_serve_argv(args, kind, key)
}
fn block_ready<A: ArgvView + ?Sized>(
&self,
store: &mut Store,
serve_argv: &A,
kind: BlockKind,
) -> bool {
cmd_block_serve::block_ready(store, serve_argv, kind)
}
fn wake_idx<A: ArgvView + ?Sized>(&self, args: &A) -> Option<u8> {
let name = args.first()?;
let mut buf = [0u8; 32];
cmd_block::wake_idx_for_verb(upper_verb(name, &mut buf))
}
fn resolve<A: ArgvView + ?Sized>(&self, args: &A) -> ResolvedCmd {
cmd_resolve::kevy_resolve(args)
}
}