#![forbid(unsafe_code)]
use kevy_resp::{encode_error, parse_command};
use kevy_rt::{
ArgvView, BlockKind, Commands, NotifyClass, ResolvedCmd, RespVersion, Route, Runtime, TxnKind,
parse_slowlog_sub,
};
use kevy_store::Store;
use kevy_sys::Socket;
use std::io;
use std::path::PathBuf;
use std::sync::Arc;
use std::sync::atomic::AtomicBool;
mod cmd;
mod cmd_block;
mod cmd_block_serve;
mod cmd_data;
mod cmd_hello;
mod cmd_resolve;
mod config_global;
mod dispatch;
mod dispatch_collections;
mod dispatch_resp3;
mod dispatch_geo;
mod dispatch_stream;
mod ops;
pub use config_global::init as config_init;
pub use config_global::replace as config_replace;
use cmd::{scan_pattern, upper_verb};
pub use dispatch::dispatch;
pub use kevy_rt::Argv;
pub use kevy_store::Store as KeyspaceStore;
pub enum AfterDrain {
KeepOpen,
Close,
}
#[derive(Clone, Copy, Default)]
pub struct KevyCommands;
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"WAIT" | b"SHUTDOWN"
| b"CLIENT" | b"SELECT" => Route::Local,
b"DBSIZE" => Route::Dbsize,
b"FLUSHDB" | b"FLUSHALL" => Route::Flush,
b"SAVE" | b"BGSAVE" => Route::Save,
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"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"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" => {
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(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 shard_tick_interval_ms(&self) -> u64 {
let cfg = config_global::get();
let hz = cfg.expiry.hz;
if hz == 0 {
0
} else {
(1000 / hz as u64).clamp(1, 10_000)
}
}
fn on_shard_tick(&self, store: &mut Store) {
let cfg = config_global::get();
let samples = cfg.expiry.sample as usize;
store.tick_expire(samples, 16);
store.set_max_memory(
cfg.memory.maxmemory,
map_eviction_policy(cfg.memory.maxmemory_policy),
);
}
fn live_runtime_config(&self) -> kevy_rt::LiveRuntimeConfig {
if !config_global::is_initialised() {
return 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 / hz as u64).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),
}
}
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)
}
}
fn map_eviction_policy(p: kevy_config::EvictionPolicy) -> kevy_store::EvictionPolicy {
use kevy_config::EvictionPolicy as C;
use kevy_store::EvictionPolicy as S;
match p {
C::NoEviction => S::NoEviction,
C::AllKeysLru => S::AllKeysLru,
C::AllKeysLfu => S::AllKeysLfu,
C::AllKeysRandom => S::AllKeysRandom,
C::VolatileLru => S::VolatileLru,
C::VolatileLfu => S::VolatileLfu,
C::VolatileRandom => S::VolatileRandom,
C::VolatileTtl => S::VolatileTtl,
}
}
pub fn serve(ip: [u8; 4], port: u16, nshards: usize, data_dir: PathBuf, enable_aof: bool) -> ! {
let cfg = config_global::get();
let fsync = map_appendfsync(cfg.persistence.appendfsync);
let runtime = Runtime::new(ip, port, nshards, KevyCommands)
.with_data_dir(data_dir)
.with_aof(enable_aof)
.with_appendfsync(fsync)
.with_auto_aof_rewrite(
cfg.persistence.auto_aof_rewrite_percentage,
cfg.persistence.auto_aof_rewrite_min_size,
)
.with_advanced(
cfg.advanced.spin_limit,
cfg.advanced.park_timeout_ms,
cfg.advanced.tick_check_every,
cfg.advanced.ring_capacity,
)
.with_slowlog(cfg.slowlog.slower_than_micros, cfg.slowlog.max_len);
let stop = Arc::new(AtomicBool::new(false));
if let Err(e) = runtime.run(stop) {
eprintln!("kevy: runtime error: {e}");
std::process::exit(1);
}
std::process::exit(0);
}
fn map_appendfsync(p: kevy_config::AppendFsync) -> kevy_persist::Fsync {
use kevy_config::AppendFsync as C;
use kevy_persist::Fsync as P;
match p {
C::Always => P::Always,
C::EverySec => P::EverySec,
C::No => P::No,
}
}
pub fn drain_commands(store: &mut Store, input: &mut Vec<u8>, output: &mut Vec<u8>) -> AfterDrain {
loop {
match parse_command(input) {
Ok(Some((args, consumed))) => {
let reply = dispatch(store, &args);
output.extend_from_slice(&reply);
input.drain(..consumed);
if args
.first()
.is_some_and(|c| c.eq_ignore_ascii_case(b"QUIT"))
{
return AfterDrain::Close;
}
}
Ok(None) => return AfterDrain::KeepOpen,
Err(_) => {
encode_error(output, "ERR Protocol error");
return AfterDrain::Close;
}
}
}
}
pub fn handle_conn(conn: &Socket, store: &mut Store) -> io::Result<()> {
let mut input: Vec<u8> = Vec::with_capacity(4096);
let mut output: Vec<u8> = Vec::new();
let mut chunk = [0u8; 4096];
loop {
let after = drain_commands(store, &mut input, &mut output);
if !output.is_empty() {
conn.write_all(&output)?;
output.clear();
}
if matches!(after, AfterDrain::Close) {
return Ok(());
}
let n = conn.read(&mut chunk)?;
if n == 0 {
return Ok(());
}
input.extend_from_slice(&chunk[..n]);
}
}
#[cfg(test)]
mod tests;