#![forbid(unsafe_code)]
use kevy_resp::{encode_error, parse_command};
use kevy_rt::Runtime;
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_lua;
mod cmd_resolve;
mod commands;
mod config_global;
mod replication;
mod dispatch;
mod dispatch_collections;
mod dispatch_collections_v127;
mod dispatch_resp3;
mod dispatch_geo;
mod dispatch_stream;
mod elect_integration;
mod ops;
mod replica_runner;
mod replica_state;
mod scope_integration;
pub use config_global::init as config_init;
pub use config_global::replace as config_replace;
pub use dispatch::dispatch;
pub use kevy_rt::Argv;
pub use kevy_store::Store as KeyspaceStore;
#[doc(hidden)]
pub fn install_replica_senders_for_test(senders: Vec<kevy_rt::ReplicaInboxSender>) {
replica_state::install_senders(senders);
}
#[doc(hidden)]
pub fn install_scope_integration_for_test(cfg: &kevy_config::Config) -> Result<(), String> {
scope_integration::install(cfg)?;
scope_integration::install_self_id(cfg);
Ok(())
}
pub enum AfterDrain {
KeepOpen,
Close,
}
#[derive(Clone, Copy, Default)]
pub struct KevyCommands;
pub(crate) 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 mut runtime = Runtime::new(ip, port, nshards, KevyCommands)
.with_data_dir(data_dir)
.with_accept_shards(cfg.server.accept_shards)
.with_max_clients(cfg.server.max_clients)
.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);
if cfg.cluster.enabled {
runtime = runtime.with_cluster(cluster_port_base(&cfg));
}
if let Ok(path) = std::env::var("KEVY_UNIX_SOCKET") {
if !path.is_empty() {
runtime = runtime.with_unix_socket(PathBuf::from(path));
}
}
let runtime = replication::apply(runtime, &cfg, nshards);
elect_integration::install_shard_offsets(nshards);
elect_integration::maybe_start(&cfg);
if let Err(msg) = scope_integration::install(&cfg) {
eprintln!("kevy: bad [cluster] scopes config: {msg}");
std::process::exit(1);
}
scope_integration::install_self_id(&cfg);
let stop = Arc::new(AtomicBool::new(false));
let run_result = runtime.run(stop);
elect_integration::shutdown();
if let Err(e) = run_result {
eprintln!("kevy: runtime error: {e}");
std::process::exit(1);
}
std::process::exit(0);
}
pub(crate) fn cluster_port_base(cfg: &kevy_config::Config) -> u16 {
match cfg.cluster.port_base {
0 => cfg.server.port.saturating_add(1),
base => base,
}
}
pub(crate) 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;